@cryptotaxi247 / kubo / commits / 98cde1578

integrate dagmodifier into ipns

Jeromy committed Oct 7, 2014 at 07:23 UTC 98cde1578d975fd655747643299e1db55e9ba54b
2 files changed +25 -30
fuse/ipns/ipns_unix.go
+19 -26
@@ -1,7 +1,6 @@
1 package ipns
2
3 import (
4 - "bytes"
4 "io/ioutil"
5 "os"
6 "path/filepath"
@@ -14,6 +13,7 @@ import (
13 "github.com/jbenet/go-ipfs/core"
14 ci "github.com/jbenet/go-ipfs/crypto"
15 imp "github.com/jbenet/go-ipfs/importer"
16 + dt "github.com/jbenet/go-ipfs/importer/dagwriter"
17 ft "github.com/jbenet/go-ipfs/importer/format"
18 mdag "github.com/jbenet/go-ipfs/merkledag"
19 u "github.com/jbenet/go-ipfs/util"
@@ -199,11 +199,8 @@ type Node struct {
199
200 Ipfs *core.IpfsNode
201 Nd *mdag.Node
202 - fd *mdag.DagReader
202 + dagMod *dt.DagModifier
203 cached *ft.PBData
204 -
205 - // For writing
206 - writerBuf WriteAtBuf
204 }
205
206 func (s *Node) loadData() error {
@@ -303,35 +300,32 @@ func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
300
301 func (n *Node) Write(req *fuse.WriteRequest, resp *fuse.WriteResponse, intr fs.Intr) fuse.Error {
302 log.Debug("ipns: Node Write [%s]: flags = %s, offset = %d, size = %d", n.name, req.Flags.String(), req.Offset, len(req.Data))
306 - if n.writerBuf == nil {
307 - n.writerBuf = NewWriterAtFromBytes(nil)
303 + if n.dagMod == nil {
304 + dmod, err := dt.NewDagModifier(n.Nd, n.Ipfs.DAG, imp.DefaultSplitter)
305 + if err != nil {
306 + log.Error("Error creating dag modifier: %s", err)
307 + return err
308 + }
309 + n.dagMod = dmod
310 }
309 - _, err := n.writerBuf.WriteAt(req.Data, req.Offset)
311 + wrote, err := n.dagMod.WriteAt(req.Data, uint64(req.Offset))
312 if err != nil {
313 return err
314 }
313 - resp.Size = len(req.Data)
315 + resp.Size = wrote
316 return nil
317 }
318
319 func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
320 log.Debug("Got flush request [%s]!", n.name)
321
320 - if n.writerBuf != nil {
321 - //TODO:
322 - // This operation holds everything in memory,
323 - // should be changed to stream the block creation/storage
324 - // but for now, since the buf is all in memory anyways...
325 -
326 - //NOTE:
327 - // This should only occur on a file object, if this were to be a
328 - // folder, bad things would happen.
329 - buf := bytes.NewReader(n.writerBuf.Bytes())
330 - newNode, err := imp.NewDagFromReader(buf)
322 + if n.dagMod != nil {
323 + newNode, err := n.dagMod.GetNode()
324 if err != nil {
332 - log.Critical("error creating dag from writerBuf: %s", err)
325 + log.Error("Error getting dag node from dagMod: %s", err)
326 return err
327 }
328 +
329 if n.parent != nil {
330 log.Debug("updating self in parent!")
331 err := n.parent.update(n.name, newNode)
@@ -358,7 +352,7 @@ func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
352 fmt.Println(b)
353 //*/
354
361 - n.writerBuf = nil
355 + n.dagMod = nil
356
357 n.wasChanged()
358 }
@@ -390,8 +384,6 @@ func (n *Node) republishRoot() error {
384 return err
385 }
386
393 - n.writerBuf = nil
394 -
387 ndkey, err := root.Nd.Key()
388 if err != nil {
389 log.Error("getKey error: %s", err)
@@ -451,8 +443,9 @@ func (n *Node) Open(req *fuse.OpenRequest, resp *fuse.OpenResponse, intr fs.Intr
443 //TODO: check open flags and truncate if necessary
444 if req.Flags&fuse.OpenTruncate != 0 {
445 log.Warning("Need to truncate file!")
454 - }
455 - if req.Flags&fuse.OpenAppend != 0 {
446 + n.cached = nil
447 + n.Nd = &mdag.Node{Data: ft.FilePBData(nil, 0)}
448 + } else if req.Flags&fuse.OpenAppend != 0 {
449 log.Warning("Need to append to file!")
450 }
451 return n, nil
importer/dagwriter/dagmodifier.go
+6 -4
@@ -109,8 +109,6 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
109 for i, size := range dm.pbdata.Blocksizes[startsubblk:] {
110 if end > traversed {
111 changed = append(changed, i+startsubblk)
112 - } else if end == traversed {
113 - break
112 } else {
113 break
114 }
@@ -122,6 +120,7 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
120 }
121 }
122
123 + // If our write starts in the middle of a block...
124 var midlnk *mdag.Link
125 if mid >= 0 {
126 midlnk = dm.curNode.Links[mid]
@@ -139,6 +138,7 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
138 b = append(b, data[midoff:]...)
139 }
140
141 + // Generate new sub-blocks, and sizes
142 subblocks := splitBytes(b, dm.splitter)
143 var links []*mdag.Link
144 var sizes []uint64
@@ -159,8 +159,10 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
159
160 // This is disgusting
161 if len(changed) > 0 {
162 - dm.curNode.Links = append(dm.curNode.Links[:changed[0]], append(links, dm.curNode.Links[changed[len(changed)-1]+1:]...)...)
163 - dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes[:changed[0]], append(sizes, dm.pbdata.Blocksizes[changed[len(changed)-1]+1:]...)...)
162 + sechalflink := append(links, dm.curNode.Links[changed[len(changed)-1]+1:]...)
163 + dm.curNode.Links = append(dm.curNode.Links[:changed[0]], sechalflink...)
164 + sechalfblks := append(sizes, dm.pbdata.Blocksizes[changed[len(changed)-1]+1:]...)
165 + dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes[:changed[0]], sechalfblks...)
166 } else {
167 dm.curNode.Links = append(dm.curNode.Links, links...)
168 dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes, sizes...)