@cryptotaxi247 / kubo / commits / 5639042df

bitswap: send wantlist code reuse + debug logs

Juan Batiz-Benet committed Jan 3, 2015 at 03:01 UTC 5639042df52d1b6a9a4ffb5e493a89f90c4c8329
1 file changed +62 -23
exchange/bitswap/bitswap.go
+62 -23
@@ -3,6 +3,7 @@
3 package bitswap
4
5 import (
6 + "fmt"
7 "math"
8 "sync"
9 "time"
@@ -170,58 +171,96 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
171 return bs.network.Provide(ctx, blk.Key())
172 }
173
173 -func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.ID) error {
174 +func (bs *bitswap) sendWantlistMsgToPeer(ctx context.Context, m bsmsg.BitSwapMessage, p peer.ID) error {
175 + logd := fmt.Sprintf("%s bitswap.sendWantlistMsgToPeer(%d, %s)", bs.self, len(m.Wantlist()), p)
176 +
177 + log.Debugf("%s sending wantlist", logd)
178 + if err := bs.send(ctx, p, m); err != nil {
179 + log.Errorf("%s send wantlist error: %s", logd, err)
180 + return err
181 + }
182 + log.Debugf("%s send wantlist success", logd)
183 + return nil
184 +}
185 +
186 +func (bs *bitswap) sendWantlistMsgToPeers(ctx context.Context, m bsmsg.BitSwapMessage, peers <-chan peer.ID) error {
187 if peers == nil {
188 panic("Cant send wantlist to nil peerchan")
189 }
177 - message := bsmsg.New()
178 - for _, wanted := range bs.wantlist.Entries() {
179 - message.AddEntry(wanted.Key, wanted.Priority)
180 - }
190 +
191 + logd := fmt.Sprintf("%s bitswap.sendWantlistMsgTo(%d)", bs.self, len(m.Wantlist()))
192 + log.Debugf("%s begin", logd)
193 + defer log.Debugf("%s end", logd)
194 +
195 + set := pset.New()
196 wg := sync.WaitGroup{}
197 for peerToQuery := range peers {
198 log.Event(ctx, "PeerToQuery", peerToQuery)
199 + logd := fmt.Sprintf("%sto(%s)", logd, peerToQuery)
200 +
201 + if !set.TryAdd(peerToQuery) { //Do once per peer
202 + log.Debugf("%s skipped (already sent)", logd)
203 + continue
204 + }
205 +
206 wg.Add(1)
207 go func(p peer.ID) {
208 defer wg.Done()
187 - if err := bs.send(ctx, p, message); err != nil {
188 - log.Error(err)
189 - return
190 - }
209 + bs.sendWantlistMsgToPeer(ctx, m, p)
210 }(peerToQuery)
211 }
212 wg.Wait()
213 return nil
214 }
215
197 -func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wantlist.ThreadSafe) {
198 - ctx, cancel := context.WithCancel(ctx)
199 - defer cancel()
200 -
216 +func (bs *bitswap) sendWantlistToPeers(ctx context.Context, peers <-chan peer.ID) error {
217 message := bsmsg.New()
218 message.SetFull(true)
203 - for _, e := range bs.wantlist.Entries() {
204 - message.AddEntry(e.Key, e.Priority)
219 + for _, wanted := range bs.wantlist.Entries() {
220 + message.AddEntry(wanted.Key, wanted.Priority)
221 }
222 + return bs.sendWantlistMsgToPeers(ctx, message, peers)
223 +}
224
207 - set := pset.New()
225 +func (bs *bitswap) sendWantlistToProviders(ctx context.Context) {
226 + logd := fmt.Sprintf("%s bitswap.sendWantlistToProviders", bs.self)
227 + log.Debugf("%s begin", logd)
228 + defer log.Debugf("%s end", logd)
229 +
230 + ctx, cancel := context.WithCancel(ctx)
231 + defer cancel()
232 +
233 + // prepare a channel to hand off to sendWantlistToPeers
234 + sendToPeers := make(chan peer.ID)
235
236 // Get providers for all entries in wantlist (could take a while)
237 wg := sync.WaitGroup{}
211 - for _, e := range wantlist.Entries() {
238 + for _, e := range bs.wantlist.Entries() {
239 wg.Add(1)
240 go func(k u.Key) {
241 defer wg.Done()
242 +
243 + logd := fmt.Sprintf("%s(entry: %s)", logd, k)
244 + log.Debugf("%s asking dht for providers", logd)
245 +
246 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
247 providers := bs.network.FindProvidersAsync(child, k, maxProvidersPerRequest)
248 for prov := range providers {
218 - if set.TryAdd(prov) { //Do once per peer
219 - bs.send(ctx, prov, message)
220 - }
249 + log.Debugf("%s dht returned provider %s. send wantlist", logd, prov)
250 + sendToPeers <- prov
251 }
252 }(e.Key)
253 }
224 - wg.Wait()
254 +
255 + go func() {
256 + wg.Wait() // make sure all our children do finish.
257 + close(sendToPeers)
258 + }()
259 +
260 + err := bs.sendWantlistToPeers(ctx, sendToPeers)
261 + if err != nil {
262 + log.Errorf("%s sendWantlistToPeers error: %s", logd, err)
263 + }
264 }
265
266 func (bs *bitswap) taskWorker(ctx context.Context) {
@@ -247,7 +286,7 @@ func (bs *bitswap) clientWorker(parent context.Context) {
286 select {
287 case <-broadcastSignal:
288 // Resend unfulfilled wantlist keys
250 - bs.sendWantlistToProviders(ctx, bs.wantlist)
289 + bs.sendWantlistToProviders(ctx)
290 broadcastSignal = time.After(rebroadcastDelay.Get())
291 case ks := <-bs.batchRequests:
292 if len(ks) == 0 {
@@ -266,7 +305,7 @@ func (bs *bitswap) clientWorker(parent context.Context) {
305 // newer bitswap strategies.
306 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
307 providers := bs.network.FindProvidersAsync(child, ks[0], maxProvidersPerRequest)
269 - err := bs.sendWantListTo(ctx, providers)
308 + err := bs.sendWantlistToPeers(ctx, providers)
309 if err != nil {
310 log.Errorf("error sending wantlist: %s", err)
311 }