@cryptotaxi247 / kubo / commits / f81f6ca6a

implement some simple dht request read timeouts

License: MIT Signed-off-by: Jeromy <why@ipfs.io>

Jeromy committed Jun 17, 2016 at 11:01 UTC f81f6ca6a63ff50f23c1c37984d9533f26eb10f0
2 files changed +49 -3
routing/dht/dht.go
+40 -3
@@ -107,6 +107,16 @@ func (dht *IpfsDHT) putValueToPeer(ctx context.Context, p peer.ID,
107 pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(key), 0)
108 pmes.Record = rec
109 rpmes, err := dht.sendRequest(ctx, p, pmes)
110 + switch err {
111 + case ErrReadTimeout:
112 + log.Errorf("read timeout: %s %s", p, key)
113 + fallthrough
114 + default:
115 + return err
116 + case nil:
117 + break
118 + }
119 +
120 if err != nil {
121 return err
122 }
@@ -164,7 +174,16 @@ func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.ID,
174 defer log.EventBegin(ctx, "getValueSingle", p, &key).Done()
175
176 pmes := pb.NewMessage(pb.Message_GET_VALUE, string(key), 0)
167 - return dht.sendRequest(ctx, p, pmes)
177 + resp, err := dht.sendRequest(ctx, p, pmes)
178 + switch err {
179 + case nil:
180 + return resp, nil
181 + case ErrReadTimeout:
182 + log.Errorf("read timeout: %s %s", p, key)
183 + fallthrough
184 + default:
185 + return nil, err
186 + }
187 }
188
189 // getLocal attempts to retrieve the value from the datastore
@@ -238,14 +257,32 @@ func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.ID, id peer.ID) (
257 defer log.EventBegin(ctx, "findPeerSingle", p, id).Done()
258
259 pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), 0)
241 - return dht.sendRequest(ctx, p, pmes)
260 + resp, err := dht.sendRequest(ctx, p, pmes)
261 + switch err {
262 + case nil:
263 + return resp, nil
264 + case ErrReadTimeout:
265 + log.Errorf("read timeout: %s %s", p, id)
266 + fallthrough
267 + default:
268 + return nil, err
269 + }
270 }
271
272 func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.ID, key key.Key) (*pb.Message, error) {
273 defer log.EventBegin(ctx, "findProvidersSingle", p, &key).Done()
274
275 pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, string(key), 0)
248 - return dht.sendRequest(ctx, p, pmes)
276 + resp, err := dht.sendRequest(ctx, p, pmes)
277 + switch err {
278 + case nil:
279 + return resp, nil
280 + case ErrReadTimeout:
281 + log.Errorf("read timeout: %s %s", p, key)
282 + fallthrough
283 + default:
284 + return nil, err
285 + }
286 }
287
288 // nearestPeersToQuery returns the routing tables closest peers.
routing/dht/dht_net.go
+9
@@ -1,6 +1,7 @@
1 package dht
2
3 import (
4 + "fmt"
5 "sync"
6 "time"
7
@@ -12,6 +13,9 @@ import (
13 inet "gx/ipfs/QmdBpVuSYuTGDA8Kn66CbKvEThXqKUh2nTANZEhzSxqrmJ/go-libp2p/p2p/net"
14 )
15
16 +var dhtReadMessageTimeout = time.Minute
17 +var ErrReadTimeout = fmt.Errorf("timed out reading response")
18 +
19 // handleNewStream implements the inet.StreamHandler
20 func (dht *IpfsDHT) handleNewStream(s inet.Stream) {
21 go dht.handleNewMessage(s)
@@ -232,10 +236,15 @@ func (ms *messageSender) ctxReadMsg(ctx context.Context, mes *pb.Message) error
236 errc <- r.ReadMsg(mes)
237 }(ms.r)
238
239 + t := time.NewTimer(dhtReadMessageTimeout)
240 + defer t.Stop()
241 +
242 select {
243 case err := <-errc:
244 return err
245 case <-ctx.Done():
246 return ctx.Err()
247 + case <-t.C:
248 + return ErrReadTimeout
249 }
250 }