@cryptotaxi247 / kubo / commits / 136bece85

Code cleanup

Jeromy committed Mar 9, 2015 at 14:53 UTC 136bece85072d1c928af9b672f46bef60946a8c2
1 file changed +31 -17
unixfs/mod/dagmodifier.go
+31 -17
@@ -17,7 +17,6 @@ import (
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"
20 u "github.com/jbenet/go-ipfs/util"
21 )
22
@@ -56,9 +55,10 @@ func NewDagModifier(ctx context.Context, from *mdag.Node, serv mdag.DAGService,
55 }
56
57 // WriteAt will modify a dag file in place
59 -// NOTE: it currently assumes only a single level of indirection
58 func (dm *DagModifier) WriteAt(b []byte, offset int64) (int, error) {
59 // TODO: this is currently VERY inneficient
60 + // each write that happens at an offset other than the current one causes a
61 + // flush to disk, and dag rewrite
62 if uint64(offset) != dm.curWrOff {
63 size, err := dm.Size()
64 if err != nil {
@@ -91,6 +91,8 @@ func (zr zeroReader) Read(b []byte) (int, error) {
91 return len(b), nil
92 }
93
94 +// expandSparse grows the file with zero blocks of 4096
95 +// A small blocksize is chosen to aid in deduplication
96 func (dm *DagModifier) expandSparse(size int64) error {
97 spl := chunk.SizeSplitter{4096}
98 r := io.LimitReader(zeroReader{}, size)
@@ -107,6 +109,7 @@ func (dm *DagModifier) expandSparse(size int64) error {
109 return nil
110 }
111
112 +// Write continues writing to the dag at the current offset
113 func (dm *DagModifier) Write(b []byte) (int, error) {
114 if dm.read != nil {
115 dm.read = nil
@@ -114,6 +117,7 @@ func (dm *DagModifier) Write(b []byte) (int, error) {
117 if dm.wrBuf == nil {
118 dm.wrBuf = new(bytes.Buffer)
119 }
120 +
121 n, err := dm.wrBuf.Write(b)
122 if err != nil {
123 return n, err
@@ -143,7 +147,9 @@ func (dm *DagModifier) Size() (int64, error) {
147 return int64(pbn.GetFilesize()), nil
148 }
149
150 +// Flush writes changes to this dag to disk
151 func (dm *DagModifier) Flush() error {
152 + // No buffer? Nothing to do
153 if dm.wrBuf == nil {
154 return nil
155 }
@@ -154,9 +160,11 @@ func (dm *DagModifier) Flush() error {
160 dm.readCancel()
161 }
162
163 + // Number of bytes we're going to write
164 buflen := dm.wrBuf.Len()
165
159 - k, _, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
166 + // overwrite existing dag nodes
167 + k, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
168 if err != nil {
169 return err
170 }
@@ -168,6 +176,7 @@ func (dm *DagModifier) Flush() error {
176
177 dm.curNode = nd
178
179 + // need to write past end of current dag
180 if !done {
181 blks := dm.splitter.Split(dm.wrBuf)
182 nd, err = dm.appendData(dm.curNode, blks)
@@ -189,28 +198,30 @@ func (dm *DagModifier) Flush() error {
198 return nil
199 }
200
192 -func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader) (u.Key, int, bool, error) {
201 +// modifyDag writes the data in 'data' over the data in 'node' starting at 'offset'
202 +func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader) (u.Key, bool, error) {
203 f, err := ft.FromBytes(node.Data)
204 if err != nil {
195 - return "", 0, false, err
205 + return "", false, err
206 }
207
198 - if len(node.Links) == 0 && (f.GetType() == ftpb.Data_Raw || f.GetType() == ftpb.Data_File) {
208 + // If we've reached a leaf node.
209 + if len(node.Links) == 0 {
210 n, err := data.Read(f.Data[offset:])
211 if err != nil && err != io.EOF {
201 - return "", 0, false, err
212 + return "", false, err
213 }
214
215 // Update newly written node..
216 b, err := proto.Marshal(f)
217 if err != nil {
207 - return "", 0, false, err
218 + return "", false, err
219 }
220
221 nd := &mdag.Node{Data: b}
222 k, err := dm.dagserv.Add(nd)
223 if err != nil {
213 - return "", 0, false, err
224 + return "", false, err
225 }
226
227 // Hey look! we're done!
@@ -219,23 +230,21 @@ func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader)
230 done = true
231 }
232
222 - return k, n, done, nil
233 + return k, done, nil
234 }
235
236 var cur uint64
237 var done bool
227 - var totread int
238 for i, bs := range f.GetBlocksizes() {
239 if cur+bs > offset {
240 child, err := node.Links[i].GetNode(dm.dagserv)
241 if err != nil {
232 - return "", 0, false, err
242 + return "", false, err
243 }
234 - k, nread, sdone, err := dm.modifyDag(child, offset-cur, data)
244 + k, sdone, err := dm.modifyDag(child, offset-cur, data)
245 if err != nil {
236 - return "", 0, false, err
246 + return "", false, err
247 }
238 - totread += nread
248
249 offset += bs
250 node.Links[i].Hash = mh.Multihash(k)
@@ -249,9 +258,10 @@ func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader)
258 }
259
260 k, err := dm.dagserv.Add(node)
252 - return k, totread, done, err
261 + return k, done, err
262 }
263
264 +// appendData appends the blocks from the given chan to the end of this dag
265 func (dm *DagModifier) appendData(node *mdag.Node, blks <-chan []byte) (*mdag.Node, error) {
266 dbp := &help.DagBuilderParams{
267 Dagserv: dm.dagserv,
@@ -262,6 +272,7 @@ func (dm *DagModifier) appendData(node *mdag.Node, blks <-chan []byte) (*mdag.No
272 return trickle.TrickleAppend(node, dbp.New(blks))
273 }
274
275 +// Read data from this dag starting at the current offset
276 func (dm *DagModifier) Read(b []byte) (int, error) {
277 err := dm.Flush()
278 if err != nil {
@@ -322,6 +333,7 @@ func (dm *DagModifier) GetNode() (*mdag.Node, error) {
333 return dm.curNode.Copy(), nil
334 }
335
336 +// HasChanges returned whether or not there are unflushed changes to this dag
337 func (dm *DagModifier) HasChanges() bool {
338 return dm.wrBuf != nil
339 }
@@ -366,8 +378,9 @@ func (dm *DagModifier) Truncate(size int64) error {
378 return err
379 }
380
381 + // Truncate can also be used to expand the file
382 if size > int64(realSize) {
370 - return errors.New("Cannot extend file through truncate")
383 + return dm.expandSparse(int64(size) - realSize)
384 }
385
386 nnode, err := dagTruncate(dm.curNode, uint64(size), dm.dagserv)
@@ -384,6 +397,7 @@ func (dm *DagModifier) Truncate(size int64) error {
397 return nil
398 }
399
400 +// dagTruncate truncates the given node to 'size' and returns the modified Node
401 func dagTruncate(nd *mdag.Node, size uint64, ds mdag.DAGService) (*mdag.Node, error) {
402 if len(nd.Links) == 0 {
403 // TODO: this can likely be done without marshaling and remarshaling