@cryptotaxi247 / kubo / commits / e8536db35

make bitswap sub-RPC's timeout (slowly for now)

Jeromy committed Dec 2, 2014 at 07:34 UTC e8536db3510a1554f5ed7c9704e1a937bc923cb9
3 files changed +25 -5
cmd/ipfs/main.go
+2
@@ -6,6 +6,7 @@ import (
6 "io"
7 "os"
8 "os/signal"
9 + "runtime"
10 "runtime/pprof"
11 "syscall"
12
@@ -54,6 +55,7 @@ type cmdInvocation struct {
55 // - output the response
56 // - if anything fails, print error, maybe with help
57 func main() {
58 + runtime.GOMAXPROCS(3)
59 ctx := context.Background()
60 var err error
61 var invoc cmdInvocation
exchange/bitswap/bitswap.go
+22 -4
@@ -26,6 +26,9 @@ var log = eventlog.Logger("bitswap")
26 // TODO: if a 'non-nice' strategy is implemented, consider increasing this value
27 const maxProvidersPerRequest = 3
28
29 +const providerRequestTimeout = time.Second * 10
30 +const hasBlockTimeout = time.Second * 15
31 +
32 // New initializes a BitSwap instance that communicates over the
33 // provided BitSwapNetwork. This function registers the returned instance as
34 // the network delegate.
@@ -181,7 +184,8 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, ks []u.Key) {
184 for _, k := range ks {
185 wg.Add(1)
186 go func(k u.Key) {
184 - providers := bs.routing.FindProvidersAsync(ctx, k, maxProvidersPerRequest)
187 + child, _ := context.WithTimeout(ctx, providerRequestTimeout)
188 + providers := bs.routing.FindProvidersAsync(child, k, maxProvidersPerRequest)
189
190 err := bs.sendWantListTo(ctx, providers)
191 if err != nil {
@@ -228,7 +232,8 @@ func (bs *bitswap) loop(parent context.Context) {
232 // pinning a file, you store and provide all blocks associated with
233 // it. Later, this assumption may not hold as true if we implement
234 // newer bitswap strategies.
231 - providers := bs.routing.FindProvidersAsync(ctx, ks[0], maxProvidersPerRequest)
235 + child, _ := context.WithTimeout(ctx, providerRequestTimeout)
236 + providers := bs.routing.FindProvidersAsync(child, ks[0], maxProvidersPerRequest)
237
238 err := bs.sendWantListTo(ctx, providers)
239 if err != nil {
@@ -247,8 +252,21 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
252 log.Debugf("Has Block %s", blk.Key())
253 bs.wantlist.Remove(blk.Key())
254 bs.notifications.Publish(blk)
250 - bs.sendToPeersThatWant(ctx, blk)
251 - return bs.routing.Provide(ctx, blk.Key())
255 +
256 + var err error
257 + wg := &sync.WaitGroup{}
258 + wg.Add(2)
259 + child, _ := context.WithTimeout(ctx, hasBlockTimeout)
260 + go func() {
261 + bs.sendToPeersThatWant(child, blk)
262 + wg.Done()
263 + }()
264 + go func() {
265 + err = bs.routing.Provide(child, blk.Key())
266 + wg.Done()
267 + }()
268 + wg.Wait()
269 + return err
270 }
271
272 // receiveBlock handles storing the block in the blockstore and calling HasBlock
routing/dht/routing.go
+1 -1
@@ -126,7 +126,7 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
126 }
127
128 func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int) <-chan peer.Peer {
129 - log.Debug("Find Providers: %s", key)
129 + log.Debugf("Find Providers: %s", key)
130 peerOut := make(chan peer.Peer, count)
131 go func() {
132 ps := newPeerSet()