@cryptotaxi247 / kubo / commits / 472753516

fixing mutability issues in ipns

Jeromy committed Oct 3, 2014 at 19:22 UTC 47275351603c1beb9772781302240728e9a87280
4 files changed +149 -45
fuse/ipns/ipns_unix.go
+107 -28
@@ -1,6 +1,7 @@
1 package ipns
2
3 import (
4 + "fmt"
5 "io/ioutil"
6 "os"
7 "path/filepath"
@@ -72,7 +73,7 @@ func CreateRoot(n *core.IpfsNode, keys []ci.PrivKey, ipfsroot string) (*Root, er
73 nd := new(Node)
74 nd.Ipfs = n
75 nd.key = k
75 - nd.repub = NewRepublisher(nd, time.Millisecond*10)
76 + nd.repub = NewRepublisher(nd, time.Millisecond*10, time.Second)
77
78 go nd.repub.Run()
79
@@ -182,11 +183,15 @@ func (r *Root) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
183
184 // Node is the core object representing a filesystem tree node.
185 type Node struct {
186 + root *Root
187 nsRoot *Node
188 + parent *Node
189
190 repub *Republisher
191
189 - // Name really only for logging purposes
192 + // This nodes name in its parent dir.
193 + // NOTE: this strategy wont work well if we allow hard links
194 + // (im all for murdering the thought of hard links)
195 name string
196
197 // Private keys held by nodes at the root of a keyspace
@@ -244,9 +249,10 @@ func (s *Node) Lookup(name string, intr fs.Intr) (fs.Node, fuse.Error) {
249
250 func (n *Node) makeChild(name string, node *mdag.Node) *Node {
251 child := &Node{
247 - Ipfs: n.Ipfs,
248 - Nd: node,
249 - name: n.name + "/" + name,
252 + Ipfs: n.Ipfs,
253 + Nd: node,
254 + name: name,
255 + nsRoot: n.nsRoot,
256 }
257
258 if n.nsRoot == nil {
@@ -289,7 +295,7 @@ func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
295 }
296
297 func (n *Node) Write(req *fuse.WriteRequest, resp *fuse.WriteResponse, intr fs.Intr) fuse.Error {
292 - log.Debug("ipns: Node Write: flags = %s, offset = %d, size = %d", req.Flags.String(), req.Offset, len(req.Data))
298 + log.Debug("ipns: Node Write [%s]: flags = %s, offset = %d, size = %d", n.name, req.Flags.String(), req.Offset, len(req.Data))
299 if n.dataBuf == nil {
300 n.dataBuf = new(bytes.Buffer)
301 }
@@ -312,19 +318,45 @@ func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
318 // This operation holds everything in memory,
319 // should be changed to stream the block creation/storage
320 // but for now, since the buf is all in memory anyways...
315 - err := imp.NewDagInNode(n.dataBuf, n.Nd)
321 +
322 + //NOTE:
323 + // This should only occur on a file object, if this were to be a
324 + // folder, bad things would happen.
325 + newNode, err := imp.NewDagFromReader(n.dataBuf)
326 if err != nil {
317 - log.Error("ipns: Flush error: %s", err)
318 - // return fuse.EVERYBAD
327 + log.Critical("error creating dag from dataBuf: %s", err)
328 return fuse.ENODATA
329 }
330 + if n.parent != nil {
331 + err := n.parent.update(n.name, newNode)
332 + if err != nil {
333 + log.Critical("error in updating ipns dag tree: %s", err)
334 + // return fuse.ETHISISPRETTYBAD
335 + return fuse.ENOSYS
336 + }
337 + }
338 + n.Nd = newNode
339
322 - n.sendPublishSignal()
340 + //TEMP
341 + dr, err := mdag.NewDagReader(n.Nd, n.Ipfs.DAG)
342 + if err != nil {
343 + log.Critical("Verification read failed.")
344 + }
345 + b, err := ioutil.ReadAll(dr)
346 + if err != nil {
347 + log.Critical("Verification read failed.")
348 + }
349 + fmt.Println("VERIFICATION READ")
350 + fmt.Printf("READ %d BYTES\n", len(b))
351 + fmt.Println(string(b))
352 + //
353 +
354 + n.wasChanged()
355 }
356 return nil
357 }
358
327 -func (n *Node) sendPublishSignal() {
359 +func (n *Node) wasChanged() {
360 root := n.nsRoot
361 if root == nil {
362 root = n
@@ -333,7 +365,8 @@ func (n *Node) sendPublishSignal() {
365 root.repub.Publish <- struct{}{}
366 }
367
336 -func (n *Node) updateTree() error {
368 +func (n *Node) republishRoot() error {
369 + log.Debug("Republish root")
370 var root *Node
371 if n.nsRoot != nil {
372 root = n.nsRoot
@@ -341,19 +374,13 @@ func (n *Node) updateTree() error {
374 root = n
375 }
376
344 - err := root.Nd.Update()
345 - if err != nil {
346 - log.Error("ipns: dag tree update failed: %s", err)
347 - return err
348 - }
349 -
350 - err = n.Ipfs.DAG.AddRecursive(root.Nd)
377 + // Add any nodes that may be new to the DAG service
378 + err := n.Ipfs.DAG.AddRecursive(root.Nd)
379 if err != nil {
380 log.Critical("ipns: Dag Add Error: %s", err)
381 return err
382 }
383
356 - n.changed = false
384 n.dataBuf = nil
385
386 ndkey, err := root.Nd.Key()
@@ -380,11 +407,13 @@ func (n *Node) Fsync(req *fuse.FsyncRequest, intr fs.Intr) fuse.Error {
407 func (n *Node) Mkdir(req *fuse.MkdirRequest, intr fs.Intr) (fs.Node, fuse.Error) {
408 log.Debug("Got mkdir request!")
409 dagnd := &mdag.Node{Data: mdag.FolderPBData()}
383 - n.Nd.AddNodeLink(req.Name, dagnd)
410 + nnode := n.Nd.Copy()
411 + nnode.AddNodeLink(req.Name, dagnd)
412
413 child := &Node{
414 Ipfs: n.Ipfs,
415 Nd: dagnd,
416 + name: req.Name,
417 }
418
419 if n.nsRoot == nil {
@@ -393,8 +422,17 @@ func (n *Node) Mkdir(req *fuse.MkdirRequest, intr fs.Intr) (fs.Node, fuse.Error)
422 child.nsRoot = n.nsRoot
423 }
424
396 - n.changed = true
397 - n.sendPublishSignal()
425 + if n.parent != nil {
426 + err := n.parent.update(n.name, nnode)
427 + if err != nil {
428 + log.Critical("Error updating node: %s", err)
429 + // Can we panic, please?
430 + return nil, fuse.ENODATA
431 + }
432 + }
433 + n.Nd = nnode
434 +
435 + n.wasChanged()
436
437 return child, nil
438 }
@@ -405,7 +443,7 @@ func (n *Node) Mknod(req *fuse.MknodRequest, intr fs.Intr) (fs.Node, fuse.Error)
443 }
444
445 func (n *Node) Open(req *fuse.OpenRequest, resp *fuse.OpenResponse, intr fs.Intr) (fs.Handle, fuse.Error) {
408 - log.Debug("[%s] Received open request! flags = %s", n.name, req.Flags.String())
446 + //log.Debug("[%s] Received open request! flags = %s", n.name, req.Flags.String())
447 //TODO: check open flags and truncate if necessary
448 return n, nil
449 }
@@ -417,23 +455,46 @@ func (n *Node) Create(req *fuse.CreateRequest, resp *fuse.CreateResponse, intr f
455 nd := &mdag.Node{Data: mdag.FilePBData(nil)}
456 child := n.makeChild(req.Name, nd)
457
420 - err := n.Nd.AddNodeLink(req.Name, nd)
458 + nnode := n.Nd.Copy()
459 +
460 + err := nnode.AddNodeLink(req.Name, nd)
461 if err != nil {
462 log.Error("Error adding child to node: %s", err)
463 return nil, nil, fuse.ENOENT
464 }
465 + if n.parent != nil {
466 + err := n.parent.update(n.name, nnode)
467 + if err != nil {
468 + log.Critical("Error updating node: %s", err)
469 + // Can we panic, please?
470 + return nil, nil, fuse.ENODATA
471 + }
472 + }
473 + n.Nd = nnode
474 + n.wasChanged()
475 +
476 return child, child, nil
477 }
478
479 func (n *Node) Remove(req *fuse.RemoveRequest, intr fs.Intr) fuse.Error {
480 log.Debug("[%s] Got Remove request: %s", n.name, req.Name)
430 - err := n.Nd.RemoveNodeLink(req.Name)
481 + nnode := n.Nd.Copy()
482 + err := nnode.RemoveNodeLink(req.Name)
483 if err != nil {
484 log.Error("Remove: No such file.")
485 return fuse.ENOENT
486 }
435 - n.changed = true
436 - n.sendPublishSignal()
487 +
488 + if n.parent != nil {
489 + err := n.parent.update(n.name, nnode)
490 + if err != nil {
491 + log.Critical("Error updating node: %s", err)
492 + // Can we panic, please?
493 + return fuse.ENODATA
494 + }
495 + }
496 + n.Nd = nnode
497 + n.wasChanged()
498 return nil
499 }
500
@@ -464,3 +525,21 @@ func (n *Node) Rename(req *fuse.RenameRequest, newDir fs.Node, intr fs.Intr) fus
525 }
526 return nil
527 }
528 +
529 +func (n *Node) update(name string, newnode *mdag.Node) error {
530 + nnode := n.Nd.Copy()
531 + err := nnode.RemoveNodeLink(name)
532 + if err != nil {
533 + return err
534 + }
535 + nnode.AddNodeLink(name, newnode)
536 +
537 + if n.parent != nil {
538 + err := n.parent.update(n.name, newnode)
539 + if err != nil {
540 + return err
541 + }
542 + }
543 + n.Nd = nnode
544 + return nil
545 +}
fuse/ipns/repub_unix.go
+23 -12
@@ -3,34 +3,45 @@ package ipns
3 import "time"
4
5 type Republisher struct {
6 - Timeout time.Duration
7 - Publish chan struct{}
8 - node *Node
6 + TimeoutLong time.Duration
7 + TimeoutShort time.Duration
8 + Publish chan struct{}
9 + node *Node
10 }
11
11 -func NewRepublisher(n *Node, tout time.Duration) *Republisher {
12 +func NewRepublisher(n *Node, tshort, tlong time.Duration) *Republisher {
13 return &Republisher{
13 - Timeout: tout,
14 - Publish: make(chan struct{}),
15 - node: n,
14 + TimeoutShort: tshort,
15 + TimeoutLong: tlong,
16 + Publish: make(chan struct{}),
17 + node: n,
18 }
19 }
20
21 func (np *Republisher) Run() {
22 for _ = range np.Publish {
21 - timer := time.After(np.Timeout)
23 + quick := time.After(np.TimeoutShort)
24 + longer := time.After(np.TimeoutLong)
25 for {
26 select {
24 - case <-timer:
27 + case <-quick:
28 //Do the publish!
29 log.Info("Publishing Changes!")
27 - err := np.node.updateTree()
30 + err := np.node.republishRoot()
31 if err != nil {
29 - log.Critical("updateTree error: %s", err)
32 + log.Critical("republishRoot error: %s", err)
33 + }
34 + goto done
35 + case <-longer:
36 + //Do the publish!
37 + log.Info("Publishing Changes!")
38 + err := np.node.republishRoot()
39 + if err != nil {
40 + log.Critical("republishRoot error: %s", err)
41 }
42 goto done
43 case <-np.Publish:
33 - timer = time.After(np.Timeout)
44 + quick = time.After(np.TimeoutShort)
45 }
46 }
47 done:
merkledag/merkledag.go
+10
@@ -75,6 +75,16 @@ func (n *Node) RemoveNodeLink(name string) error {
75 return u.ErrNotFound
76 }
77
78 +func (n *Node) Copy() *Node {
79 + nnode := new(Node)
80 + nnode.Data = make([]byte, len(n.Data))
81 + copy(nnode.Data, n.Data)
82 +
83 + nnode.Links = make([]*Link, len(n.Links))
84 + copy(nnode.Links, n.Links)
85 + return nnode
86 +}
87 +
88 // Size returns the total size of the data addressed by node,
89 // including the total sizes of references.
90 func (n *Node) Size() (uint64, error) {
util/util.go
+9 -5
@@ -99,11 +99,15 @@ func DOut(format string, a ...interface{}) {
99 func SetupLogging() {
100 backend := logging.NewLogBackend(os.Stderr, "", 0)
101 logging.SetBackend(backend)
102 - if Debug {
103 - logging.SetLevel(logging.DEBUG, "")
104 - } else {
105 - logging.SetLevel(logging.ERROR, "")
106 - }
102 + /*
103 + if Debug {
104 + logging.SetLevel(logging.DEBUG, "")
105 + } else {
106 + logging.SetLevel(logging.ERROR, "")
107 + }
108 + */
109 + logging.SetLevel(logging.ERROR, "merkledag")
110 + logging.SetLevel(logging.ERROR, "blockservice")
111 logging.SetFormatter(logging.MustStringFormatter(LogFormat))
112 }
113