@cryptotaxi247 / kubo / commits / 5618fed0d

coreapi pubsub: better ctx for connectToPubSubPeers

License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>

Łukasz Magiera committed Oct 2, 2018 at 12:51 UTC 5618fed0d2adcb994123a588711ee4aab94e0aec
1 file changed +8 -4
core/coreapi/pubsub.go
+8 -4
@@ -20,6 +20,7 @@ import (
20 type PubSubAPI CoreAPI
21
22 type pubSubSubscription struct {
23 + cancel context.CancelFunc
24 subscription *floodsub.Subscription
25 }
26
@@ -75,19 +76,21 @@ func (api *PubSubAPI) Subscribe(ctx context.Context, topic string, opts ...caopt
76 return nil, err
77 }
78
79 + pubctx, cancel := context.WithCancel(api.node.Context())
80 +
81 if options.Discover {
82 go func() {
80 - blk, err := api.core().Block().Put(ctx, strings.NewReader("floodsub:"+topic))
83 + blk, err := api.core().Block().Put(pubctx, strings.NewReader("floodsub:"+topic))
84 if err != nil {
85 log.Error("pubsub discovery: ", err)
86 return
87 }
88
86 - connectToPubSubPeers(ctx, api.node, blk.Path().Cid())
89 + connectToPubSubPeers(pubctx, api.node, blk.Path().Cid())
90 }()
91 }
92
90 - return &pubSubSubscription{sub}, nil
93 + return &pubSubSubscription{cancel, sub}, nil
94 }
95
96 func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid cid.Cid) {
@@ -95,7 +98,7 @@ func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid cid.Cid) {
98 defer cancel()
99
100 provs := n.Routing.FindProvidersAsync(ctx, cid, 10)
98 - wg := &sync.WaitGroup{}
101 + var wg sync.WaitGroup
102 for p := range provs {
103 wg.Add(1)
104 go func(pi pstore.PeerInfo) {
@@ -127,6 +130,7 @@ func (api *PubSubAPI) checkNode() error {
130 }
131
132 func (sub *pubSubSubscription) Close() error {
133 + sub.cancel()
134 sub.subscription.Cancel()
135 return nil
136 }