@cryptotaxi247 / kubo / commits / 9652ada0d

implement publisher for ipns to wait until moments of rapid churn die down

Jeromy committed Sep 29, 2014 at 17:47 UTC 9652ada0d2029f8d97ca1219dbbf2e8d33e13699
3 files changed +25 -14
fuse/ipns/ipns_unix.go
+21 -9
@@ -76,6 +76,9 @@ func CreateRoot(n *core.IpfsNode, keys []ci.PrivKey, ipfsroot string) (*Root, er
76 nd := new(Node)
77 nd.Ipfs = n
78 nd.key = k
79 + nd.repub = NewRepublisher(nd, time.Millisecond*10)
80 +
81 + go nd.repub.Run()
82
83 pointsTo, err := n.Namesys.Resolve(name)
84 if err != nil {
@@ -185,6 +188,8 @@ func (r *Root) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
188 type Node struct {
189 nsRoot *Node
190
191 + repub *Republisher
192 +
193 // Name really only for logging purposes
194 name string
195
@@ -316,16 +321,20 @@ func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
321 return fuse.ENODATA
322 }
323
319 - err = n.updateTree()
320 - if err != nil {
321 - log.Error("updateTree failed: %s", err)
322 - return fuse.ENODATA
323 - }
324 -
324 + n.sendPublishSignal()
325 }
326 return nil
327 }
328
329 +func (n *Node) sendPublishSignal() {
330 + root := n.nsRoot
331 + if root == nil {
332 + root = n
333 + }
334 +
335 + root.repub.Publish <- struct{}{}
336 +}
337 +
338 func (n *Node) updateTree() error {
339 var root *Node
340 if n.nsRoot != nil {
@@ -387,7 +396,7 @@ func (n *Node) Mkdir(req *fuse.MkdirRequest, intr fs.Intr) (fs.Node, fuse.Error)
396 }
397
398 n.changed = true
390 - n.updateTree()
399 + n.sendPublishSignal()
400
401 return child, nil
402 }
@@ -425,6 +434,8 @@ func (n *Node) Remove(req *fuse.RemoveRequest, intr fs.Intr) fuse.Error {
434 log.Error("Remove: No such file.")
435 return fuse.ENOENT
436 }
437 + n.changed = true
438 + n.sendPublishSignal()
439 return nil
440 }
441
@@ -471,7 +482,7 @@ func Mount(ipfs *core.IpfsNode, fpath string, ipfspath string) error {
482 if err == nil {
483 return
484 }
474 - time.Sleep(time.Millisecond * 10)
485 + time.Sleep(time.Millisecond * 100)
486 }
487 ipfs.Network.Close()
488 }()
@@ -563,12 +574,13 @@ func NewRepublisher(n *Node, tout time.Duration) *Republisher {
574 }
575
576 func (np *Republisher) Run() {
566 - for _ := range np.Publish {
577 + for _ = range np.Publish {
578 timer := time.After(np.Timeout)
579 for {
580 select {
581 case <-timer:
582 //Do the publish!
583 + log.Info("Publishing Changes!")
584 err := np.node.updateTree()
585 if err != nil {
586 log.Critical("updateTree error: %s", err)
routing/dht/dht.go
+3 -4
@@ -177,8 +177,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
177 start := time.Now()
178
179 // Print out diagnostic
180 - log.Debug("[peer: %s] Sent message type: '%s' [to = %s]\n",
181 - dht.self.ID.Pretty(),
180 + log.Debug("Sent message type: '%s' [to = %s]",
181 Message_MessageType_name[int32(pmes.GetType())], p.ID.Pretty())
182
183 rmes, err := dht.sender.SendRequest(ctx, mes)
@@ -281,7 +280,7 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
280 }
281
282 if len(peers) > 0 {
284 - u.DOut("getValueOrPeers: peers")
283 + log.Debug("getValueOrPeers: peers")
284 return nil, peers, nil
285 }
286
@@ -400,7 +399,7 @@ func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer
399 for _, prov := range peers {
400 p, err := dht.peerFromInfo(prov)
401 if err != nil {
403 - u.PErr("error getting peer from info: %v\n", err)
402 + log.Error("error getting peer from info: %v", err)
403 continue
404 }
405
routing/dht/routing.go
+1 -1
@@ -18,7 +18,7 @@ import (
18 // PutValue adds value corresponding to given Key.
19 // This is the top level "Store" operation of the DHT
20 func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error {
21 - log.Debug("PutValue %s %v", key.Pretty(), value)
21 + log.Debug("PutValue %s", key.Pretty())
22 err := dht.putLocal(key, value)
23 if err != nil {
24 return err