@cryptotaxi247 / kubo / commits / 63e15abd8

refactor dagmodifier to work with trickledag format

Jeromy committed Mar 7, 2015 at 00:58 UTC 63e15abd8f22cfed3266bb08bfdf4bde5453c667
6 files changed +937 -453
fuse/ipns/ipns_unix.go
+9 -4
@@ -24,6 +24,7 @@ import (
24 path "github.com/jbenet/go-ipfs/path"
25 ft "github.com/jbenet/go-ipfs/unixfs"
26 uio "github.com/jbenet/go-ipfs/unixfs/io"
27 + mod "github.com/jbenet/go-ipfs/unixfs/mod"
28 ftpb "github.com/jbenet/go-ipfs/unixfs/pb"
29 u "github.com/jbenet/go-ipfs/util"
30 lgbl "github.com/jbenet/go-ipfs/util/eventlog/loggables"
@@ -211,7 +212,7 @@ type Node struct {
212
213 Ipfs *core.IpfsNode
214 Nd *mdag.Node
214 - dagMod *uio.DagModifier
215 + dagMod *mod.DagModifier
216 cached *ftpb.Data
217 }
218
@@ -238,7 +239,11 @@ func (s *Node) Attr() fuse.Attr {
239 size = 0
240 }
241 if size == 0 {
241 - size = s.dagMod.Size()
242 + dmsize, err := s.dagMod.Size()
243 + if err != nil {
244 + log.Error(err)
245 + }
246 + size = uint64(dmsize)
247 }
248
249 mode := os.FileMode(0666)
@@ -344,13 +349,13 @@ func (n *Node) Write(ctx context.Context, req *fuse.WriteRequest, resp *fuse.Wri
349
350 if n.dagMod == nil {
351 // Create a DagModifier to allow us to change the existing dag node
347 - dmod, err := uio.NewDagModifier(n.Nd, n.Ipfs.DAG, chunk.DefaultSplitter)
352 + dmod, err := mod.NewDagModifier(ctx, n.Nd, n.Ipfs.DAG, n.Ipfs.Pinning.GetManual(), chunk.DefaultSplitter)
353 if err != nil {
354 return err
355 }
356 n.dagMod = dmod
357 }
353 - wrote, err := n.dagMod.WriteAt(req.Data, uint64(req.Offset))
358 + wrote, err := n.dagMod.WriteAt(req.Data, int64(req.Offset))
359 if err != nil {
360 return err
361 }
importer/trickle/trickledag.go
+14
@@ -244,6 +244,20 @@ func verifyTDagRec(nd *dag.Node, depth, direct, layerRepeat int, ds dag.DAGServi
244 return nil
245 }
246
247 + // Verify this is a branch node
248 + pbn, err := ft.FromBytes(nd.Data)
249 + if err != nil {
250 + return err
251 + }
252 +
253 + if pbn.GetType() != ft.TFile {
254 + return errors.New("expected file as branch node")
255 + }
256 +
257 + if len(pbn.Data) > 0 {
258 + return errors.New("branch node should not have data")
259 + }
260 +
261 for i := 0; i < len(nd.Links); i++ {
262 child, err := nd.Links[i].GetNode(ds)
263 if err != nil {
unixfs/io/dagmodifier.go deleted
-202
@@ -1,202 +0,0 @@
1 -package io
2 -
3 -import (
4 - "bytes"
5 - "errors"
6 -
7 - proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
8 -
9 - chunk "github.com/jbenet/go-ipfs/importer/chunk"
10 - mdag "github.com/jbenet/go-ipfs/merkledag"
11 - ft "github.com/jbenet/go-ipfs/unixfs"
12 - ftpb "github.com/jbenet/go-ipfs/unixfs/pb"
13 - u "github.com/jbenet/go-ipfs/util"
14 -)
15 -
16 -var log = u.Logger("dagio")
17 -
18 -// DagModifier is the only struct licensed and able to correctly
19 -// perform surgery on a DAG 'file'
20 -// Dear god, please rename this to something more pleasant
21 -type DagModifier struct {
22 - dagserv mdag.DAGService
23 - curNode *mdag.Node
24 -
25 - pbdata *ftpb.Data
26 - splitter chunk.BlockSplitter
27 -}
28 -
29 -func NewDagModifier(from *mdag.Node, serv mdag.DAGService, spl chunk.BlockSplitter) (*DagModifier, error) {
30 - pbd, err := ft.FromBytes(from.Data)
31 - if err != nil {
32 - return nil, err
33 - }
34 -
35 - return &DagModifier{
36 - curNode: from.Copy(),
37 - dagserv: serv,
38 - pbdata: pbd,
39 - splitter: spl,
40 - }, nil
41 -}
42 -
43 -// WriteAt will modify a dag file in place
44 -// NOTE: it currently assumes only a single level of indirection
45 -func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
46 -
47 - // Check bounds
48 - if dm.pbdata.GetFilesize() < offset {
49 - return 0, errors.New("Attempted to perform write starting past end of file")
50 - }
51 -
52 - // First need to find where we are writing at
53 - end := uint64(len(b)) + offset
54 -
55 - // This shouldnt be necessary if we do subblocks sizes properly
56 - newsize := dm.pbdata.GetFilesize()
57 - if end > dm.pbdata.GetFilesize() {
58 - newsize = end
59 - }
60 - zeroblocklen := uint64(len(dm.pbdata.Data))
61 - origlen := len(b)
62 -
63 - if end <= zeroblocklen {
64 - log.Debug("Writing into zero block")
65 - // Replacing zeroeth data block (embedded in the root node)
66 - //TODO: check chunking here
67 - copy(dm.pbdata.Data[offset:], b)
68 - return len(b), nil
69 - }
70 -
71 - // Find where write should start
72 - var traversed uint64
73 - startsubblk := len(dm.pbdata.Blocksizes)
74 - if offset < zeroblocklen {
75 - dm.pbdata.Data = dm.pbdata.Data[:offset]
76 - startsubblk = 0
77 - } else {
78 - traversed = uint64(zeroblocklen)
79 - for i, size := range dm.pbdata.Blocksizes {
80 - if uint64(offset) < traversed+size {
81 - log.Debugf("Starting mod at block %d. [%d < %d + %d]", i, offset, traversed, size)
82 - // Here is where we start
83 - startsubblk = i
84 - lnk := dm.curNode.Links[i]
85 - node, err := dm.dagserv.Get(u.Key(lnk.Hash))
86 - if err != nil {
87 - return 0, err
88 - }
89 - data, err := ft.UnwrapData(node.Data)
90 - if err != nil {
91 - return 0, err
92 - }
93 -
94 - // We have to rewrite the data before our write in this block.
95 - b = append(data[:offset-traversed], b...)
96 - break
97 - }
98 - traversed += size
99 - }
100 - if startsubblk == len(dm.pbdata.Blocksizes) {
101 - // TODO: Im not sure if theres any case that isnt being handled here.
102 - // leaving this note here as a future reference in case something breaks
103 - }
104 - }
105 -
106 - // Find blocks that need to be overwritten
107 - var changed []int
108 - mid := -1
109 - var midoff uint64
110 - for i, size := range dm.pbdata.Blocksizes[startsubblk:] {
111 - if end > traversed {
112 - changed = append(changed, i+startsubblk)
113 - } else {
114 - break
115 - }
116 - traversed += size
117 - if end < traversed {
118 - mid = i + startsubblk
119 - midoff = end - (traversed - size)
120 - break
121 - }
122 - }
123 -
124 - // If our write starts in the middle of a block...
125 - var midlnk *mdag.Link
126 - if mid >= 0 {
127 - midlnk = dm.curNode.Links[mid]
128 - midnode, err := dm.dagserv.Get(u.Key(midlnk.Hash))
129 - if err != nil {
130 - return 0, err
131 - }
132 -
133 - // NOTE: this may have to be changed later when we have multiple
134 - // layers of indirection
135 - data, err := ft.UnwrapData(midnode.Data)
136 - if err != nil {
137 - return 0, err
138 - }
139 - b = append(b, data[midoff:]...)
140 - }
141 -
142 - // Generate new sub-blocks, and sizes
143 - subblocks := splitBytes(b, dm.splitter)
144 - var links []*mdag.Link
145 - var sizes []uint64
146 - for _, sb := range subblocks {
147 - n := &mdag.Node{Data: ft.WrapData(sb)}
148 - _, err := dm.dagserv.Add(n)
149 - if err != nil {
150 - log.Warningf("Failed adding node to DAG service: %s", err)
151 - return 0, err
152 - }
153 - lnk, err := mdag.MakeLink(n)
154 - if err != nil {
155 - return 0, err
156 - }
157 - links = append(links, lnk)
158 - sizes = append(sizes, uint64(len(sb)))
159 - }
160 -
161 - // This is disgusting (and can be rewritten if performance demands)
162 - if len(changed) > 0 {
163 - sechalflink := append(links, dm.curNode.Links[changed[len(changed)-1]+1:]...)
164 - dm.curNode.Links = append(dm.curNode.Links[:changed[0]], sechalflink...)
165 - sechalfblks := append(sizes, dm.pbdata.Blocksizes[changed[len(changed)-1]+1:]...)
166 - dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes[:changed[0]], sechalfblks...)
167 - } else {
168 - dm.curNode.Links = append(dm.curNode.Links, links...)
169 - dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes, sizes...)
170 - }
171 - dm.pbdata.Filesize = proto.Uint64(newsize)
172 -
173 - return origlen, nil
174 -}
175 -
176 -func (dm *DagModifier) Size() uint64 {
177 - if dm == nil {
178 - return 0
179 - }
180 - return dm.pbdata.GetFilesize()
181 -}
182 -
183 -// splitBytes uses a splitterFunc to turn a large array of bytes
184 -// into many smaller arrays of bytes
185 -func splitBytes(b []byte, spl chunk.BlockSplitter) [][]byte {
186 - out := spl.Split(bytes.NewReader(b))
187 - var arr [][]byte
188 - for blk := range out {
189 - arr = append(arr, blk)
190 - }
191 - return arr
192 -}
193 -
194 -// GetNode gets the modified DAG Node
195 -func (dm *DagModifier) GetNode() (*mdag.Node, error) {
196 - b, err := proto.Marshal(dm.pbdata)
197 - if err != nil {
198 - return nil, err
199 - }
200 - dm.curNode.Data = b
201 - return dm.curNode.Copy(), nil
202 -}
unixfs/io/dagmodifier_test.go deleted
-247
@@ -1,247 +0,0 @@
1 -package io
2 -
3 -import (
4 - "fmt"
5 - "io"
6 - "io/ioutil"
7 - "testing"
8 -
9 - "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
10 - "github.com/jbenet/go-ipfs/blocks/blockstore"
11 - bs "github.com/jbenet/go-ipfs/blockservice"
12 - "github.com/jbenet/go-ipfs/exchange/offline"
13 - imp "github.com/jbenet/go-ipfs/importer"
14 - "github.com/jbenet/go-ipfs/importer/chunk"
15 - mdag "github.com/jbenet/go-ipfs/merkledag"
16 - ft "github.com/jbenet/go-ipfs/unixfs"
17 - u "github.com/jbenet/go-ipfs/util"
18 -
19 - ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
20 - context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
21 -)
22 -
23 -func getMockDagServ(t *testing.T) mdag.DAGService {
24 - dstore := ds.NewMapDatastore()
25 - tsds := sync.MutexWrap(dstore)
26 - bstore := blockstore.NewBlockstore(tsds)
27 - bserv, err := bs.New(bstore, offline.Exchange(bstore))
28 - if err != nil {
29 - t.Fatal(err)
30 - }
31 - return mdag.NewDAGService(bserv)
32 -}
33 -
34 -func getNode(t *testing.T, dserv mdag.DAGService, size int64) ([]byte, *mdag.Node) {
35 - in := io.LimitReader(u.NewTimeSeededRand(), size)
36 - node, err := imp.BuildDagFromReader(in, dserv, nil, &chunk.SizeSplitter{500})
37 - if err != nil {
38 - t.Fatal(err)
39 - }
40 -
41 - dr, err := NewDagReader(context.Background(), node, dserv)
42 - if err != nil {
43 - t.Fatal(err)
44 - }
45 -
46 - b, err := ioutil.ReadAll(dr)
47 - if err != nil {
48 - t.Fatal(err)
49 - }
50 -
51 - return b, node
52 -}
53 -
54 -func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier) []byte {
55 - newdata := make([]byte, size)
56 - r := u.NewTimeSeededRand()
57 - r.Read(newdata)
58 -
59 - if size+beg > uint64(len(orig)) {
60 - orig = append(orig, make([]byte, (size+beg)-uint64(len(orig)))...)
61 - }
62 - copy(orig[beg:], newdata)
63 -
64 - nmod, err := dm.WriteAt(newdata, uint64(beg))
65 - if err != nil {
66 - t.Fatal(err)
67 - }
68 -
69 - if nmod != int(size) {
70 - t.Fatalf("Mod length not correct! %d != %d", nmod, size)
71 - }
72 -
73 - nd, err := dm.GetNode()
74 - if err != nil {
75 - t.Fatal(err)
76 - }
77 -
78 - rd, err := NewDagReader(context.Background(), nd, dm.dagserv)
79 - if err != nil {
80 - t.Fatal(err)
81 - }
82 -
83 - after, err := ioutil.ReadAll(rd)
84 - if err != nil {
85 - t.Fatal(err)
86 - }
87 -
88 - err = arrComp(after, orig)
89 - if err != nil {
90 - t.Fatal(err)
91 - }
92 - return orig
93 -}
94 -
95 -func TestDagModifierBasic(t *testing.T) {
96 - t.Skip("DAGModifier needs to be fixed to work with indirect blocks.")
97 - if err := u.SetLogLevel("blockservice", "critical"); err != nil {
98 - t.Fatalf("testlog prepare failed: %s", err)
99 - }
100 - if err := u.SetLogLevel("merkledag", "critical"); err != nil {
101 - t.Fatalf("testlog prepare failed: %s", err)
102 - }
103 - dserv := getMockDagServ(t)
104 - b, n := getNode(t, dserv, 50000)
105 -
106 - dagmod, err := NewDagModifier(n, dserv, &chunk.SizeSplitter{Size: 512})
107 - if err != nil {
108 - t.Fatal(err)
109 - }
110 -
111 - // Within zero block
112 - beg := uint64(15)
113 - length := uint64(60)
114 -
115 - t.Log("Testing mod within zero block")
116 - b = testModWrite(t, beg, length, b, dagmod)
117 -
118 - // Within bounds of existing file
119 - beg = 1000
120 - length = 4000
121 - t.Log("Testing mod within bounds of existing file.")
122 - b = testModWrite(t, beg, length, b, dagmod)
123 -
124 - // Extend bounds
125 - beg = 49500
126 - length = 4000
127 -
128 - t.Log("Testing mod that extends file.")
129 - b = testModWrite(t, beg, length, b, dagmod)
130 -
131 - // "Append"
132 - beg = uint64(len(b))
133 - length = 3000
134 - b = testModWrite(t, beg, length, b, dagmod)
135 -
136 - // Verify reported length
137 - node, err := dagmod.GetNode()
138 - if err != nil {
139 - t.Fatal(err)
140 - }
141 -
142 - size, err := ft.DataSize(node.Data)
143 - if err != nil {
144 - t.Fatal(err)
145 - }
146 -
147 - expected := uint64(50000 + 3500 + 3000)
148 - if size != expected {
149 - t.Fatalf("Final reported size is incorrect [%d != %d]", size, expected)
150 - }
151 -}
152 -
153 -func TestMultiWrite(t *testing.T) {
154 - t.Skip("DAGModifier needs to be fixed to work with indirect blocks.")
155 - dserv := getMockDagServ(t)
156 - _, n := getNode(t, dserv, 0)
157 -
158 - dagmod, err := NewDagModifier(n, dserv, &chunk.SizeSplitter{Size: 512})
159 - if err != nil {
160 - t.Fatal(err)
161 - }
162 -
163 - data := make([]byte, 4000)
164 - u.NewTimeSeededRand().Read(data)
165 -
166 - for i := 0; i < len(data); i++ {
167 - n, err := dagmod.WriteAt(data[i:i+1], uint64(i))
168 - if err != nil {
169 - t.Fatal(err)
170 - }
171 - if n != 1 {
172 - t.Fatal("Somehow wrote the wrong number of bytes! (n != 1)")
173 - }
174 - }
175 - nd, err := dagmod.GetNode()
176 - if err != nil {
177 - t.Fatal(err)
178 - }
179 -
180 - read, err := NewDagReader(context.Background(), nd, dserv)
181 - if err != nil {
182 - t.Fatal(err)
183 - }
184 - rbuf, err := ioutil.ReadAll(read)
185 - if err != nil {
186 - t.Fatal(err)
187 - }
188 -
189 - err = arrComp(rbuf, data)
190 - if err != nil {
191 - t.Fatal(err)
192 - }
193 -}
194 -
195 -func TestMultiWriteCoal(t *testing.T) {
196 - t.Skip("Skipping test until DagModifier is fixed")
197 - dserv := getMockDagServ(t)
198 - _, n := getNode(t, dserv, 0)
199 -
200 - dagmod, err := NewDagModifier(n, dserv, &chunk.SizeSplitter{Size: 512})
201 - if err != nil {
202 - t.Fatal(err)
203 - }
204 -
205 - data := make([]byte, 4000)
206 - u.NewTimeSeededRand().Read(data)
207 -
208 - for i := 0; i < len(data); i++ {
209 - n, err := dagmod.WriteAt(data[:i+1], 0)
210 - if err != nil {
211 - t.Fatal(err)
212 - }
213 - if n != i+1 {
214 - t.Fatal("Somehow wrote the wrong number of bytes! (n != 1)")
215 - }
216 - }
217 - nd, err := dagmod.GetNode()
218 - if err != nil {
219 - t.Fatal(err)
220 - }
221 -
222 - read, err := NewDagReader(context.Background(), nd, dserv)
223 - if err != nil {
224 - t.Fatal(err)
225 - }
226 - rbuf, err := ioutil.ReadAll(read)
227 - if err != nil {
228 - t.Fatal(err)
229 - }
230 -
231 - err = arrComp(rbuf, data)
232 - if err != nil {
233 - t.Fatal(err)
234 - }
235 -}
236 -
237 -func arrComp(a, b []byte) error {
238 - if len(a) != len(b) {
239 - return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
240 - }
241 - for i, v := range a {
242 - if v != b[i] {
243 - return fmt.Errorf("Arrays differ at index: %d", i)
244 - }
245 - }
246 - return nil
247 -}
unixfs/mod/dagmodifier.go new
+450
@@ -0,0 +1,450 @@
1 +package mod
2 +
3 +import (
4 + "bytes"
5 + "errors"
6 + "io"
7 + "os"
8 +
9 + proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
10 + mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
11 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
12 +
13 + chunk "github.com/jbenet/go-ipfs/importer/chunk"
14 + help "github.com/jbenet/go-ipfs/importer/helpers"
15 + trickle "github.com/jbenet/go-ipfs/importer/trickle"
16 + mdag "github.com/jbenet/go-ipfs/merkledag"
17 + pin "github.com/jbenet/go-ipfs/pin"
18 + ft "github.com/jbenet/go-ipfs/unixfs"
19 + uio "github.com/jbenet/go-ipfs/unixfs/io"
20 + ftpb "github.com/jbenet/go-ipfs/unixfs/pb"
21 + u "github.com/jbenet/go-ipfs/util"
22 +)
23 +
24 +// 2MB
25 +var writebufferSize = 1 << 21
26 +
27 +var log = u.Logger("dagio")
28 +
29 +// DagModifier is the only struct licensed and able to correctly
30 +// perform surgery on a DAG 'file'
31 +// Dear god, please rename this to something more pleasant
32 +type DagModifier struct {
33 + dagserv mdag.DAGService
34 + curNode *mdag.Node
35 + mp pin.ManualPinner
36 +
37 + splitter chunk.BlockSplitter
38 + ctx context.Context
39 + readCancel func()
40 +
41 + writeStart uint64
42 + curWrOff uint64
43 + wrBuf *bytes.Buffer
44 +
45 + read *uio.DagReader
46 +}
47 +
48 +func NewDagModifier(ctx context.Context, from *mdag.Node, serv mdag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*DagModifier, error) {
49 + return &DagModifier{
50 + curNode: from.Copy(),
51 + dagserv: serv,
52 + splitter: spl,
53 + ctx: ctx,
54 + mp: mp,
55 + }, nil
56 +}
57 +
58 +// WriteAt will modify a dag file in place
59 +// NOTE: it currently assumes only a single level of indirection
60 +func (dm *DagModifier) WriteAt(b []byte, offset int64) (int, error) {
61 + // TODO: this is currently VERY inneficient
62 + if uint64(offset) != dm.curWrOff {
63 + size, err := dm.Size()
64 + if err != nil {
65 + return 0, err
66 + }
67 + if offset > size {
68 + err := dm.expandSparse(offset - size)
69 + if err != nil {
70 + return 0, err
71 + }
72 + }
73 +
74 + err = dm.Flush()
75 + if err != nil {
76 + return 0, err
77 + }
78 + dm.writeStart = uint64(offset)
79 + }
80 +
81 + return dm.Write(b)
82 +}
83 +
84 +// A reader that just returns zeros
85 +type zeroReader struct{}
86 +
87 +func (zr zeroReader) Read(b []byte) (int, error) {
88 + for i, _ := range b {
89 + b[i] = 0
90 + }
91 + return len(b), nil
92 +}
93 +
94 +func (dm *DagModifier) expandSparse(size int64) error {
95 + spl := chunk.SizeSplitter{4096}
96 + r := io.LimitReader(zeroReader{}, size)
97 + blks := spl.Split(r)
98 + nnode, err := dm.appendData(dm.curNode, blks)
99 + if err != nil {
100 + return err
101 + }
102 + _, err = dm.dagserv.Add(nnode)
103 + if err != nil {
104 + return err
105 + }
106 + dm.curNode = nnode
107 + return nil
108 +}
109 +
110 +func (dm *DagModifier) Write(b []byte) (int, error) {
111 + if dm.read != nil {
112 + dm.read = nil
113 + }
114 + if dm.wrBuf == nil {
115 + dm.wrBuf = new(bytes.Buffer)
116 + }
117 + n, err := dm.wrBuf.Write(b)
118 + if err != nil {
119 + return n, err
120 + }
121 + dm.curWrOff += uint64(n)
122 + if dm.wrBuf.Len() > writebufferSize {
123 + err := dm.Flush()
124 + if err != nil {
125 + return n, err
126 + }
127 + }
128 + return n, nil
129 +}
130 +
131 +func (dm *DagModifier) Size() (int64, error) {
132 + // TODO: compute size without flushing, should be easy
133 + err := dm.Flush()
134 + if err != nil {
135 + return 0, err
136 + }
137 +
138 + pbn, err := ft.FromBytes(dm.curNode.Data)
139 + if err != nil {
140 + return 0, err
141 + }
142 +
143 + return int64(pbn.GetFilesize()), nil
144 +}
145 +
146 +func (dm *DagModifier) Flush() error {
147 + if dm.wrBuf == nil {
148 + return nil
149 + }
150 +
151 + // If we have an active reader, kill it
152 + if dm.read != nil {
153 + dm.read = nil
154 + dm.readCancel()
155 + }
156 +
157 + buflen := dm.wrBuf.Len()
158 +
159 + k, _, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
160 + if err != nil {
161 + return err
162 + }
163 +
164 + nd, err := dm.dagserv.Get(k)
165 + if err != nil {
166 + return err
167 + }
168 +
169 + dm.curNode = nd
170 +
171 + if !done {
172 + blks := dm.splitter.Split(dm.wrBuf)
173 + nd, err = dm.appendData(dm.curNode, blks)
174 + if err != nil {
175 + return err
176 + }
177 +
178 + _, err := dm.dagserv.Add(nd)
179 + if err != nil {
180 + return err
181 + }
182 +
183 + dm.curNode = nd
184 + }
185 +
186 + dm.writeStart += uint64(buflen)
187 +
188 + dm.wrBuf = nil
189 + return nil
190 +}
191 +
192 +func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader) (u.Key, int, bool, error) {
193 + f, err := ft.FromBytes(node.Data)
194 + if err != nil {
195 + return "", 0, false, err
196 + }
197 +
198 + if len(node.Links) == 0 && (f.GetType() == ftpb.Data_Raw || f.GetType() == ftpb.Data_File) {
199 + n, err := data.Read(f.Data[offset:])
200 + if err != nil && err != io.EOF {
201 + return "", 0, false, err
202 + }
203 +
204 + // Update newly written node..
205 + b, err := proto.Marshal(f)
206 + if err != nil {
207 + return "", 0, false, err
208 + }
209 +
210 + nd := &mdag.Node{Data: b}
211 + k, err := dm.dagserv.Add(nd)
212 + if err != nil {
213 + return "", 0, false, err
214 + }
215 +
216 + // Hey look! we're done!
217 + var done bool
218 + if n < len(f.Data) {
219 + done = true
220 + }
221 +
222 + return k, n, done, nil
223 + }
224 +
225 + var cur uint64
226 + var done bool
227 + var totread int
228 + for i, bs := range f.GetBlocksizes() {
229 + if cur+bs > offset {
230 + child, err := node.Links[i].GetNode(dm.dagserv)
231 + if err != nil {
232 + return "", 0, false, err
233 + }
234 + k, nread, sdone, err := dm.modifyDag(child, offset-cur, data)
235 + if err != nil {
236 + return "", 0, false, err
237 + }
238 + totread += nread
239 +
240 + offset += bs
241 + node.Links[i].Hash = mh.Multihash(k)
242 +
243 + if sdone {
244 + done = true
245 + break
246 + }
247 + }
248 + cur += bs
249 + }
250 +
251 + k, err := dm.dagserv.Add(node)
252 + return k, totread, done, err
253 +}
254 +
255 +func (dm *DagModifier) appendData(node *mdag.Node, blks <-chan []byte) (*mdag.Node, error) {
256 + dbp := &help.DagBuilderParams{
257 + Dagserv: dm.dagserv,
258 + Maxlinks: help.DefaultLinksPerBlock,
259 + Pinner: dm.mp,
260 + }
261 +
262 + return trickle.TrickleAppend(node, dbp.New(blks))
263 +}
264 +
265 +func (dm *DagModifier) Read(b []byte) (int, error) {
266 + err := dm.Flush()
267 + if err != nil {
268 + return 0, err
269 + }
270 +
271 + if dm.read == nil {
272 + dr, err := uio.NewDagReader(dm.ctx, dm.curNode, dm.dagserv)
273 + if err != nil {
274 + return 0, err
275 + }
276 +
277 + i, err := dr.Seek(int64(dm.curWrOff), os.SEEK_SET)
278 + if err != nil {
279 + return 0, err
280 + }
281 +
282 + if i != int64(dm.curWrOff) {
283 + return 0, errors.New("failed to seek properly")
284 + }
285 +
286 + dm.read = dr
287 + }
288 +
289 + n, err := dm.read.Read(b)
290 + dm.curWrOff += uint64(n)
291 + return n, err
292 +}
293 +
294 +// splitBytes uses a splitterFunc to turn a large array of bytes
295 +// into many smaller arrays of bytes
296 +func (dm *DagModifier) splitBytes(in io.Reader) ([]u.Key, error) {
297 + var out []u.Key
298 + blks := dm.splitter.Split(in)
299 + for blk := range blks {
300 + nd := help.NewUnixfsNode()
301 + nd.SetData(blk)
302 + dagnd, err := nd.GetDagNode()
303 + if err != nil {
304 + return nil, err
305 + }
306 +
307 + k, err := dm.dagserv.Add(dagnd)
308 + if err != nil {
309 + return nil, err
310 + }
311 + out = append(out, k)
312 + }
313 + return out, nil
314 +}
315 +
316 +// GetNode gets the modified DAG Node
317 +func (dm *DagModifier) GetNode() (*mdag.Node, error) {
318 + err := dm.Flush()
319 + if err != nil {
320 + return nil, err
321 + }
322 + return dm.curNode.Copy(), nil
323 +}
324 +
325 +func (dm *DagModifier) HasChanges() bool {
326 + return dm.wrBuf != nil
327 +}
328 +
329 +func (dm *DagModifier) Seek(offset int64, whence int) (int64, error) {
330 + err := dm.Flush()
331 + if err != nil {
332 + return 0, err
333 + }
334 +
335 + switch whence {
336 + case os.SEEK_CUR:
337 + dm.curWrOff += uint64(offset)
338 + dm.writeStart = dm.curWrOff
339 + case os.SEEK_SET:
340 + dm.curWrOff = uint64(offset)
341 + dm.writeStart = uint64(offset)
342 + case os.SEEK_END:
343 + return 0, errors.New("SEEK_END currently not implemented")
344 + default:
345 + return 0, errors.New("unrecognized whence")
346 + }
347 +
348 + if dm.read != nil {
349 + _, err = dm.read.Seek(offset, whence)
350 + if err != nil {
351 + return 0, err
352 + }
353 + }
354 +
355 + return int64(dm.curWrOff), nil
356 +}
357 +
358 +func (dm *DagModifier) Truncate(size int64) error {
359 + err := dm.Flush()
360 + if err != nil {
361 + return err
362 + }
363 +
364 + realSize, err := dm.Size()
365 + if err != nil {
366 + return err
367 + }
368 +
369 + if size > int64(realSize) {
370 + return errors.New("Cannot extend file through truncate")
371 + }
372 +
373 + nnode, err := dagTruncate(dm.curNode, uint64(size), dm.dagserv)
374 + if err != nil {
375 + return err
376 + }
377 +
378 + _, err = dm.dagserv.Add(nnode)
379 + if err != nil {
380 + return err
381 + }
382 +
383 + dm.curNode = nnode
384 + return nil
385 +}
386 +
387 +func dagTruncate(nd *mdag.Node, size uint64, ds mdag.DAGService) (*mdag.Node, error) {
388 + if len(nd.Links) == 0 {
389 + // TODO: this can likely be done without marshaling and remarshaling
390 + pbn, err := ft.FromBytes(nd.Data)
391 + if err != nil {
392 + return nil, err
393 + }
394 +
395 + nd.Data = ft.WrapData(pbn.Data[:size])
396 + return nd, nil
397 + }
398 +
399 + var cur uint64
400 + end := 0
401 + var modified *mdag.Node
402 + ndata := new(ft.FSNode)
403 + for i, lnk := range nd.Links {
404 + child, err := lnk.GetNode(ds)
405 + if err != nil {
406 + return nil, err
407 + }
408 +
409 + childsize, err := ft.DataSize(child.Data)
410 + if err != nil {
411 + return nil, err
412 + }
413 +
414 + if size < cur+childsize {
415 + nchild, err := dagTruncate(child, size-cur, ds)
416 + if err != nil {
417 + return nil, err
418 + }
419 +
420 + // TODO: sanity check size of truncated block
421 + ndata.AddBlockSize(size - cur)
422 +
423 + modified = nchild
424 + end = i
425 + break
426 + }
427 + cur += childsize
428 + ndata.AddBlockSize(childsize)
429 + }
430 +
431 + _, err := ds.Add(modified)
432 + if err != nil {
433 + return nil, err
434 + }
435 +
436 + nd.Links = nd.Links[:end]
437 + err = nd.AddNodeLinkClean("", modified)
438 + if err != nil {
439 + return nil, err
440 + }
441 +
442 + d, err := ndata.GetBytes()
443 + if err != nil {
444 + return nil, err
445 + }
446 +
447 + nd.Data = d
448 +
449 + return nd, nil
450 +}
unixfs/mod/dagmodifier_test.go new
+464
@@ -0,0 +1,464 @@
1 +package mod
2 +
3 +import (
4 + "fmt"
5 + "io"
6 + "io/ioutil"
7 + "testing"
8 +
9 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
10 + "github.com/jbenet/go-ipfs/blocks/blockstore"
11 + bs "github.com/jbenet/go-ipfs/blockservice"
12 + "github.com/jbenet/go-ipfs/exchange/offline"
13 + imp "github.com/jbenet/go-ipfs/importer"
14 + "github.com/jbenet/go-ipfs/importer/chunk"
15 + h "github.com/jbenet/go-ipfs/importer/helpers"
16 + trickle "github.com/jbenet/go-ipfs/importer/trickle"
17 + mdag "github.com/jbenet/go-ipfs/merkledag"
18 + ft "github.com/jbenet/go-ipfs/unixfs"
19 + uio "github.com/jbenet/go-ipfs/unixfs/io"
20 + u "github.com/jbenet/go-ipfs/util"
21 +
22 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
23 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
24 +)
25 +
26 +func getMockDagServ(t *testing.T) mdag.DAGService {
27 + dstore := ds.NewMapDatastore()
28 + tsds := sync.MutexWrap(dstore)
29 + bstore := blockstore.NewBlockstore(tsds)
30 + bserv, err := bs.New(bstore, offline.Exchange(bstore))
31 + if err != nil {
32 + t.Fatal(err)
33 + }
34 + return mdag.NewDAGService(bserv)
35 +}
36 +
37 +func getNode(t *testing.T, dserv mdag.DAGService, size int64) ([]byte, *mdag.Node) {
38 + in := io.LimitReader(u.NewTimeSeededRand(), size)
39 + node, err := imp.BuildTrickleDagFromReader(in, dserv, nil, &chunk.SizeSplitter{500})
40 + if err != nil {
41 + t.Fatal(err)
42 + }
43 +
44 + dr, err := uio.NewDagReader(context.Background(), node, dserv)
45 + if err != nil {
46 + t.Fatal(err)
47 + }
48 +
49 + b, err := ioutil.ReadAll(dr)
50 + if err != nil {
51 + t.Fatal(err)
52 + }
53 +
54 + return b, node
55 +}
56 +
57 +func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier) []byte {
58 + newdata := make([]byte, size)
59 + r := u.NewTimeSeededRand()
60 + r.Read(newdata)
61 +
62 + if size+beg > uint64(len(orig)) {
63 + orig = append(orig, make([]byte, (size+beg)-uint64(len(orig)))...)
64 + }
65 + copy(orig[beg:], newdata)
66 +
67 + nmod, err := dm.WriteAt(newdata, int64(beg))
68 + if err != nil {
69 + t.Fatal(err)
70 + }
71 +
72 + if nmod != int(size) {
73 + t.Fatalf("Mod length not correct! %d != %d", nmod, size)
74 + }
75 +
76 + nd, err := dm.GetNode()
77 + if err != nil {
78 + t.Fatal(err)
79 + }
80 +
81 + err = trickle.VerifyTrickleDagStructure(nd, dm.dagserv, h.DefaultLinksPerBlock, 4)
82 + if err != nil {
83 + t.Fatal(err)
84 + }
85 +
86 + rd, err := uio.NewDagReader(context.Background(), nd, dm.dagserv)
87 + if err != nil {
88 + t.Fatal(err)
89 + }
90 +
91 + after, err := ioutil.ReadAll(rd)
92 + if err != nil {
93 + t.Fatal(err)
94 + }
95 +
96 + err = arrComp(after, orig)
97 + if err != nil {
98 + t.Fatal(err)
99 + }
100 + return orig
101 +}
102 +
103 +func TestDagModifierBasic(t *testing.T) {
104 + dserv := getMockDagServ(t)
105 + b, n := getNode(t, dserv, 50000)
106 + ctx, cancel := context.WithCancel(context.Background())
107 + defer cancel()
108 +
109 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
110 + if err != nil {
111 + t.Fatal(err)
112 + }
113 +
114 + // Within zero block
115 + beg := uint64(15)
116 + length := uint64(60)
117 +
118 + t.Log("Testing mod within zero block")
119 + b = testModWrite(t, beg, length, b, dagmod)
120 +
121 + // Within bounds of existing file
122 + beg = 1000
123 + length = 4000
124 + t.Log("Testing mod within bounds of existing multiblock file.")
125 + b = testModWrite(t, beg, length, b, dagmod)
126 +
127 + // Extend bounds
128 + beg = 49500
129 + length = 4000
130 +
131 + t.Log("Testing mod that extends file.")
132 + b = testModWrite(t, beg, length, b, dagmod)
133 +
134 + // "Append"
135 + beg = uint64(len(b))
136 + length = 3000
137 + t.Log("Testing pure append")
138 + b = testModWrite(t, beg, length, b, dagmod)
139 +
140 + // Verify reported length
141 + node, err := dagmod.GetNode()
142 + if err != nil {
143 + t.Fatal(err)
144 + }
145 +
146 + size, err := ft.DataSize(node.Data)
147 + if err != nil {
148 + t.Fatal(err)
149 + }
150 +
151 + expected := uint64(50000 + 3500 + 3000)
152 + if size != expected {
153 + t.Fatalf("Final reported size is incorrect [%d != %d]", size, expected)
154 + }
155 +}
156 +
157 +func TestMultiWrite(t *testing.T) {
158 + dserv := getMockDagServ(t)
159 + _, n := getNode(t, dserv, 0)
160 +
161 + ctx, cancel := context.WithCancel(context.Background())
162 + defer cancel()
163 +
164 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
165 + if err != nil {
166 + t.Fatal(err)
167 + }
168 +
169 + data := make([]byte, 4000)
170 + u.NewTimeSeededRand().Read(data)
171 +
172 + for i := 0; i < len(data); i++ {
173 + n, err := dagmod.WriteAt(data[i:i+1], int64(i))
174 + if err != nil {
175 + t.Fatal(err)
176 + }
177 + if n != 1 {
178 + t.Fatal("Somehow wrote the wrong number of bytes! (n != 1)")
179 + }
180 + }
181 + nd, err := dagmod.GetNode()
182 + if err != nil {
183 + t.Fatal(err)
184 + }
185 +
186 + read, err := uio.NewDagReader(context.Background(), nd, dserv)
187 + if err != nil {
188 + t.Fatal(err)
189 + }
190 + rbuf, err := ioutil.ReadAll(read)
191 + if err != nil {
192 + t.Fatal(err)
193 + }
194 +
195 + err = arrComp(rbuf, data)
196 + if err != nil {
197 + t.Fatal(err)
198 + }
199 +}
200 +
201 +func TestMultiWriteAndFlush(t *testing.T) {
202 + dserv := getMockDagServ(t)
203 + _, n := getNode(t, dserv, 0)
204 +
205 + ctx, cancel := context.WithCancel(context.Background())
206 + defer cancel()
207 +
208 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
209 + if err != nil {
210 + t.Fatal(err)
211 + }
212 +
213 + data := make([]byte, 20)
214 + u.NewTimeSeededRand().Read(data)
215 +
216 + for i := 0; i < len(data); i++ {
217 + n, err := dagmod.WriteAt(data[i:i+1], int64(i))
218 + if err != nil {
219 + t.Fatal(err)
220 + }
221 + if n != 1 {
222 + t.Fatal("Somehow wrote the wrong number of bytes! (n != 1)")
223 + }
224 + err = dagmod.Flush()
225 + if err != nil {
226 + t.Fatal(err)
227 + }
228 + }
229 + nd, err := dagmod.GetNode()
230 + if err != nil {
231 + t.Fatal(err)
232 + }
233 +
234 + read, err := uio.NewDagReader(context.Background(), nd, dserv)
235 + if err != nil {
236 + t.Fatal(err)
237 + }
238 + rbuf, err := ioutil.ReadAll(read)
239 + if err != nil {
240 + t.Fatal(err)
241 + }
242 +
243 + err = arrComp(rbuf, data)
244 + if err != nil {
245 + t.Fatal(err)
246 + }
247 +}
248 +
249 +func TestWriteNewFile(t *testing.T) {
250 + dserv := getMockDagServ(t)
251 + _, n := getNode(t, dserv, 0)
252 +
253 + ctx, cancel := context.WithCancel(context.Background())
254 + defer cancel()
255 +
256 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
257 + if err != nil {
258 + t.Fatal(err)
259 + }
260 +
261 + towrite := make([]byte, 2000)
262 + u.NewTimeSeededRand().Read(towrite)
263 +
264 + nw, err := dagmod.Write(towrite)
265 + if err != nil {
266 + t.Fatal(err)
267 + }
268 + if nw != len(towrite) {
269 + t.Fatal("Wrote wrong amount")
270 + }
271 +
272 + nd, err := dagmod.GetNode()
273 + if err != nil {
274 + t.Fatal(err)
275 + }
276 +
277 + read, err := uio.NewDagReader(ctx, nd, dserv)
278 + if err != nil {
279 + t.Fatal(err)
280 + }
281 +
282 + data, err := ioutil.ReadAll(read)
283 + if err != nil {
284 + t.Fatal(err)
285 + }
286 +
287 + if err := arrComp(data, towrite); err != nil {
288 + t.Fatal(err)
289 + }
290 +}
291 +
292 +func TestMultiWriteCoal(t *testing.T) {
293 + dserv := getMockDagServ(t)
294 + _, n := getNode(t, dserv, 0)
295 +
296 + ctx, cancel := context.WithCancel(context.Background())
297 + defer cancel()
298 +
299 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
300 + if err != nil {
301 + t.Fatal(err)
302 + }
303 +
304 + data := make([]byte, 1000)
305 + u.NewTimeSeededRand().Read(data)
306 +
307 + for i := 0; i < len(data); i++ {
308 + n, err := dagmod.WriteAt(data[:i+1], 0)
309 + if err != nil {
310 + fmt.Println("FAIL AT ", i)
311 + t.Fatal(err)
312 + }
313 + if n != i+1 {
314 + t.Fatal("Somehow wrote the wrong number of bytes! (n != 1)")
315 + }
316 +
317 + // TEMP
318 + nn, err := dagmod.GetNode()
319 + if err != nil {
320 + t.Fatal(err)
321 + }
322 +
323 + r, err := uio.NewDagReader(ctx, nn, dserv)
324 + if err != nil {
325 + t.Fatal(err)
326 + }
327 +
328 + out, err := ioutil.ReadAll(r)
329 + if err != nil {
330 + t.Fatal(err)
331 + }
332 +
333 + if err := arrComp(out, data[:i+1]); err != nil {
334 + fmt.Println("A ", len(out))
335 + fmt.Println(out)
336 + fmt.Println(data[:i+1])
337 + t.Fatal(err)
338 + }
339 + //
340 + }
341 + nd, err := dagmod.GetNode()
342 + if err != nil {
343 + t.Fatal(err)
344 + }
345 +
346 + read, err := uio.NewDagReader(context.Background(), nd, dserv)
347 + if err != nil {
348 + t.Fatal(err)
349 + }
350 + rbuf, err := ioutil.ReadAll(read)
351 + if err != nil {
352 + t.Fatal(err)
353 + }
354 +
355 + err = arrComp(rbuf, data)
356 + if err != nil {
357 + t.Fatal(err)
358 + }
359 +}
360 +
361 +func TestLargeWriteChunks(t *testing.T) {
362 + dserv := getMockDagServ(t)
363 + _, n := getNode(t, dserv, 0)
364 +
365 + ctx, cancel := context.WithCancel(context.Background())
366 + defer cancel()
367 +
368 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
369 + if err != nil {
370 + t.Fatal(err)
371 + }
372 +
373 + wrsize := 1000
374 + datasize := 10000000
375 + data := make([]byte, datasize)
376 +
377 + u.NewTimeSeededRand().Read(data)
378 +
379 + for i := 0; i < datasize/wrsize; i++ {
380 + n, err := dagmod.WriteAt(data[i*wrsize:(i+1)*wrsize], int64(i*wrsize))
381 + if err != nil {
382 + t.Fatal(err)
383 + }
384 + if n != wrsize {
385 + t.Fatal("failed to write buffer")
386 + }
387 + }
388 +
389 + out, err := ioutil.ReadAll(dagmod)
390 + if err != nil {
391 + t.Fatal(err)
392 + }
393 +
394 + if err = arrComp(out, data); err != nil {
395 + t.Fatal(err)
396 + }
397 +
398 +}
399 +
400 +func TestDagTruncate(t *testing.T) {
401 + dserv := getMockDagServ(t)
402 + b, n := getNode(t, dserv, 50000)
403 + ctx, cancel := context.WithCancel(context.Background())
404 + defer cancel()
405 +
406 + dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
407 + if err != nil {
408 + t.Fatal(err)
409 + }
410 +
411 + err = dagmod.Truncate(12345)
412 + if err != nil {
413 + t.Fatal(err)
414 + }
415 +
416 + out, err := ioutil.ReadAll(dagmod)
417 + if err != nil {
418 + t.Fatal(err)
419 + }
420 +
421 + if err = arrComp(out, b[:12345]); err != nil {
422 + t.Fatal(err)
423 + }
424 +}
425 +
426 +func arrComp(a, b []byte) error {
427 + if len(a) != len(b) {
428 + return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
429 + }
430 + for i, v := range a {
431 + if v != b[i] {
432 + return fmt.Errorf("Arrays differ at index: %d", i)
433 + }
434 + }
435 + return nil
436 +}
437 +
438 +func printDag(nd *mdag.Node, ds mdag.DAGService, indent int) {
439 + pbd, err := ft.FromBytes(nd.Data)
440 + if err != nil {
441 + panic(err)
442 + }
443 +
444 + for i := 0; i < indent; i++ {
445 + fmt.Print(" ")
446 + }
447 + fmt.Printf("{size = %d, type = %s, children = %d", pbd.GetFilesize(), pbd.GetType().String(), len(pbd.GetBlocksizes()))
448 + if len(nd.Links) > 0 {
449 + fmt.Println()
450 + }
451 + for _, lnk := range nd.Links {
452 + child, err := lnk.GetNode(ds)
453 + if err != nil {
454 + panic(err)
455 + }
456 + printDag(child, ds, indent+1)
457 + }
458 + if len(nd.Links) > 0 {
459 + for i := 0; i < indent; i++ {
460 + fmt.Print(" ")
461 + }
462 + }
463 + fmt.Println("}")
464 +}