@cryptotaxi247 / kubo / commits / 19da05701

remove buffer timing in bitswap in favor of manual batching

Jeromy committed Nov 20, 2014 at 04:58 UTC 19da05701d88aeef7f5a4ad8584786668d4231c3
3 files changed +39 -35
exchange/bitswap/bitswap.go
+18 -34
@@ -43,7 +43,7 @@ func New(ctx context.Context, p peer.Peer,
43 routing: routing,
44 sender: network,
45 wantlist: u.NewKeySet(),
46 - blockRequests: make(chan u.Key, 32),
46 + batchRequests: make(chan []u.Key, 32),
47 }
48 network.SetDelegate(bs)
49 go bs.run(ctx)
@@ -66,7 +66,10 @@ type bitswap struct {
66
67 notifications notifications.PubSub
68
69 - blockRequests chan u.Key
69 + // Requests for a set of related blocks
70 + // the assumption is made that the same peer is likely to
71 + // have more than a single block in the set
72 + batchRequests chan []u.Key
73
74 // strategy listens to network traffic and makes decisions about how to
75 // interact with partners.
@@ -97,7 +100,7 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
100 promise := bs.notifications.Subscribe(ctx, k)
101
102 select {
100 - case bs.blockRequests <- k:
103 + case bs.batchRequests <- []u.Key{k}:
104 case <-parent.Done():
105 return nil, parent.Err()
106 }
@@ -159,50 +162,31 @@ func (bs *bitswap) run(ctx context.Context) {
162 // Every so often, we should resend out our current want list
163 rebroadcastTime := time.Second * 5
164
162 - var providers <-chan peer.Peer // NB: must be initialized to zero value
163 - broadcastSignal := time.After(bs.strategy.GetRebroadcastDelay())
165 + broadcastSignal := time.NewTicker(bs.strategy.GetRebroadcastDelay())
166
165 - // Number of unsent keys for the current batch
166 - unsentKeys := 0
167 for {
168 select {
169 - case <-broadcastSignal:
170 - unsentKeys = 0
169 + case <-broadcastSignal.C:
170 wantlist := bs.wantlist.Keys()
171 if len(wantlist) == 0 {
172 continue
173 }
175 - if providers == nil {
176 - // rely on semi randomness of maps
177 - firstKey := wantlist[0]
178 - providers = bs.routing.FindProvidersAsync(ctx, firstKey, maxProvidersPerRequest)
179 - }
174 + providers := bs.routing.FindProvidersAsync(ctx, wantlist[0], maxProvidersPerRequest)
175 +
176 err := bs.sendWantListTo(ctx, providers)
177 if err != nil {
178 log.Errorf("error sending wantlist: %s", err)
179 }
184 - providers = nil
185 - broadcastSignal = time.After(bs.strategy.GetRebroadcastDelay())
186 -
187 - case k := <-bs.blockRequests:
188 - if unsentKeys == 0 {
189 - providers = bs.routing.FindProvidersAsync(ctx, k, maxProvidersPerRequest)
180 + case ks := <-bs.batchRequests:
181 + if len(ks) == 0 {
182 + log.Warning("Received batch request for zero blocks")
183 + continue
184 }
191 - unsentKeys++
192 -
193 - if unsentKeys >= bs.strategy.GetBatchSize() {
194 - // send wantlist to providers
195 - err := bs.sendWantListTo(ctx, providers)
196 - if err != nil {
197 - log.Errorf("error sending wantlist: %s", err)
198 - }
199 - unsentKeys = 0
200 - broadcastSignal = time.After(bs.strategy.GetRebroadcastDelay())
201 - providers = nil
202 - } else {
203 - // set a timeout to wait for more blocks or send current wantlist
185 + providers := bs.routing.FindProvidersAsync(ctx, ks[0], maxProvidersPerRequest)
186
205 - broadcastSignal = time.After(bs.strategy.GetBatchDelay())
187 + err := bs.sendWantListTo(ctx, providers)
188 + if err != nil {
189 + log.Errorf("error sending wantlist: %s", err)
190 }
191 case <-ctx.Done():
192 return
exchange/bitswap/bitswap_test.go
+1 -1
@@ -345,7 +345,7 @@ func session(net tn.Network, rs mock.RoutingServer, id peer.ID) instance {
345 routing: htc,
346 sender: adapter,
347 wantlist: util.NewKeySet(),
348 - blockRequests: make(chan util.Key, 32),
348 + batchRequests: make(chan []util.Key, 32),
349 }
350 adapter.SetDelegate(bs)
351 go bs.run(context.TODO())
merkledag/merkledag.go
+20
@@ -252,6 +252,7 @@ func (n *dagService) Remove(nd *Node) error {
252 // FetchGraph asynchronously fetches all nodes that are children of the given
253 // node, and returns a channel that may be waited upon for the fetch to complete
254 func FetchGraph(ctx context.Context, root *Node, serv DAGService) chan struct{} {
255 + log.Warning("Untested.")
256 var wg sync.WaitGroup
257 done := make(chan struct{})
258
@@ -284,3 +285,22 @@ func FetchGraph(ctx context.Context, root *Node, serv DAGService) chan struct{}
285
286 return done
287 }
288 +
289 +// Take advantage of blockservice/bitswap batched requests to fetch all
290 +// child nodes of a given node
291 +// TODO: finish this
292 +func (ds *dagService) BatchFetch(ctx context.Context, root *Node) error {
293 + var keys []u.Key
294 + for _, lnk := range root.Links {
295 + keys = append(keys, u.Key(lnk.Hash))
296 + }
297 +
298 + blocks, err := ds.Blocks.GetBlocks(keys)
299 + if err != nil {
300 + return err
301 + }
302 +
303 + _ = blocks
304 + //what do i do with blocks?
305 + return nil
306 +}