@cryptotaxi247 / kubo / commits / e41ac9620

switch to new raceless routing event interface

fixes #5616 License: MIT Signed-off-by: Steven Allen <steven@stebalien.com>

Steven Allen committed Oct 19, 2018 at 17:08 UTC e41ac96207abfee7eaf3fdb813834cf028d12275
1 file changed +22 -20
core/commands/dht.go
+22 -20
@@ -72,23 +72,24 @@ var queryDhtCmd = &cmds.Command{
72 return
73 }
74
75 - events := make(chan *notif.QueryEvent)
76 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
77 -
75 id, err := peer.IDB58Decode(req.Arguments()[0])
76 if err != nil {
77 res.SetError(cmds.ClientError("invalid peer ID"), cmdkit.ErrClient)
78 return
79 }
80
81 + ctx, cancel := context.WithCancel(req.Context())
82 + ctx, events := notif.RegisterForQueryEvents(ctx)
83 +
84 closestPeers, err := n.DHT.GetClosestPeers(ctx, string(id))
85 if err != nil {
86 + cancel()
87 res.SetError(err, cmdkit.ErrNormal)
88 return
89 }
90
91 go func() {
91 - defer close(events)
92 + defer cancel()
93 for p := range closestPeers {
94 notif.PublishQueryEvent(ctx, &notif.QueryEvent{
95 ID: p,
@@ -182,8 +183,6 @@ var findProvidersDhtCmd = &cmds.Command{
183 return
184 }
185
185 - events := make(chan *notif.QueryEvent)
186 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
186 c, err := cid.Parse(req.Arguments()[0])
187
188 if err != nil {
@@ -194,6 +193,9 @@ var findProvidersDhtCmd = &cmds.Command{
193 outChan := make(chan interface{})
194 res.SetOutput((<-chan interface{})(outChan))
195
196 + ctx, cancel := context.WithCancel(req.Context())
197 + ctx, events := notif.RegisterForQueryEvents(ctx)
198 +
199 pchan := n.Routing.FindProvidersAsync(ctx, c, numProviders)
200 go func() {
201 defer close(outChan)
@@ -207,7 +209,7 @@ var findProvidersDhtCmd = &cmds.Command{
209 }()
210
211 go func() {
210 - defer close(events)
212 + defer cancel()
213 for p := range pchan {
214 np := p
215 notif.PublishQueryEvent(ctx, &notif.QueryEvent{
@@ -320,8 +322,8 @@ var provideRefDhtCmd = &cmds.Command{
322 outChan := make(chan interface{})
323 res.SetOutput((<-chan interface{})(outChan))
324
323 - events := make(chan *notif.QueryEvent)
324 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
325 + ctx, cancel := context.WithCancel(req.Context())
326 + ctx, events := notif.RegisterForQueryEvents(ctx)
327
328 go func() {
329 defer close(outChan)
@@ -335,7 +337,7 @@ var provideRefDhtCmd = &cmds.Command{
337 }()
338
339 go func() {
338 - defer close(events)
340 + defer cancel()
341 var err error
342 if rec {
343 err = provideKeysRec(ctx, n.Routing, n.DAG, cids)
@@ -449,8 +451,8 @@ var findPeerDhtCmd = &cmds.Command{
451 outChan := make(chan interface{})
452 res.SetOutput((<-chan interface{})(outChan))
453
452 - events := make(chan *notif.QueryEvent)
453 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
454 + ctx, cancel := context.WithCancel(req.Context())
455 + ctx, events := notif.RegisterForQueryEvents(ctx)
456
457 go func() {
458 defer close(outChan)
@@ -464,7 +466,7 @@ var findPeerDhtCmd = &cmds.Command{
466 }()
467
468 go func() {
467 - defer close(events)
469 + defer cancel()
470 pi, err := n.Routing.FindPeer(ctx, pid)
471 if err != nil {
472 notif.PublishQueryEvent(ctx, &notif.QueryEvent{
@@ -554,8 +556,8 @@ Different key types can specify other 'best' rules.
556 outChan := make(chan interface{})
557 res.SetOutput((<-chan interface{})(outChan))
558
557 - events := make(chan *notif.QueryEvent)
558 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
559 + ctx, cancel := context.WithCancel(req.Context())
560 + ctx, events := notif.RegisterForQueryEvents(ctx)
561
562 go func() {
563 defer close(outChan)
@@ -568,7 +570,7 @@ Different key types can specify other 'best' rules.
570 }()
571
572 go func() {
571 - defer close(events)
573 + defer cancel()
574 val, err := n.Routing.GetValue(ctx, dhtkey)
575 if err != nil {
576 notif.PublishQueryEvent(ctx, &notif.QueryEvent{
@@ -659,9 +661,6 @@ NOTE: A value may not exceed 2048 bytes.
661 return
662 }
663
662 - events := make(chan *notif.QueryEvent)
663 - ctx := notif.RegisterForQueryEvents(req.Context(), events)
664 -
664 key, err := escapeDhtKey(req.Arguments()[0])
665 if err != nil {
666 res.SetError(err, cmdkit.ErrNormal)
@@ -673,6 +672,9 @@ NOTE: A value may not exceed 2048 bytes.
672
673 data := req.Arguments()[1]
674
675 + ctx, cancel := context.WithCancel(req.Context())
676 + ctx, events := notif.RegisterForQueryEvents(ctx)
677 +
678 go func() {
679 defer close(outChan)
680 for e := range events {
@@ -685,7 +687,7 @@ NOTE: A value may not exceed 2048 bytes.
687 }()
688
689 go func() {
688 - defer close(events)
690 + defer cancel()
691 err := n.Routing.PutValue(ctx, key, []byte(data))
692 if err != nil {
693 notif.PublishQueryEvent(ctx, &notif.QueryEvent{