coreapi: dht: use shared set in provideKeysRec
License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>
Łukasz Magiera committed
Aug 10, 2018 at 13:40 UTC
4b252a200f2cce27c033e1480e44d7247792ee0b
1 file changed
+23
-14
core/coreapi/dht.go
+23
-14
@@ -7,6 +7,7 @@ import (
7
8
coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
9
caopts "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
10
+ "github.com/ipfs/go-ipfs/thirdparty/streaming-cid-set"
11
12
dag "gx/ipfs/QmNr4E8z9bGTztvHJktp7uQaMdx9p3r9Asrq6eYk7iCh4a/go-merkledag"
13
offline "gx/ipfs/QmPuLWvxK1vg6ckKUpT53Dow9VLCcQGdL5Trwxa8PTLp7r/go-ipfs-exchange-offline"
@@ -98,25 +99,33 @@ func provideKeys(ctx context.Context, r routing.IpfsRouting, cids []*cid.Cid) er
99
}
100
101
func provideKeysRec(ctx context.Context, r routing.IpfsRouting, bs blockstore.Blockstore, cids []*cid.Cid) error {
101
- provided := cid.NewSet()
102
- for _, c := range cids {
103
- dserv := dag.NewDAGService(blockservice.New(bs, offline.Exchange(bs)))
102
+ provided := streamingset.NewStreamingSet()
103
105
- err := dag.EnumerateChildrenAsync(ctx, dag.GetLinksDirect(dserv), c, provided.Visit)
106
- if err != nil {
107
- return err
108
- }
109
- }
104
+ errCh := make(chan error)
105
+ go func() {
106
+ for _, c := range cids {
107
+ dserv := dag.NewDAGService(blockservice.New(bs, offline.Exchange(bs)))
108
111
- for _, k := range provided.Keys() {
112
- err := r.Provide(ctx, k, true)
113
- if err != nil {
109
+ err := dag.EnumerateChildrenAsync(ctx, dag.GetLinksDirect(dserv), c, provided.Visitor(ctx))
110
+ if err != nil {
111
+ errCh <- err
112
+ }
113
+ }
114
+ }()
115
+
116
+ for {
117
+ select {
118
+ case k := <-provided.New:
119
+ err := r.Provide(ctx, k, true)
120
+ if err != nil {
121
+ return err
122
+ }
123
+ case err := <-errCh:
124
return err
125
+ case <-ctx.Done():
126
+ return ctx.Err()
127
}
116
- provided.Add(k)
128
}
118
-
119
- return nil
129
}
130
131
func (api *DhtAPI) core() coreiface.CoreAPI {