bitswap: add a deadline to sendmsg calls
License: MIT Signed-off-by: Jeromy <why@ipfs.io>
Jeromy committed
Nov 29, 2016 at 15:22 UTC
02975bde9dc2d975c9ffb74604dfe9df6f0e6aea
4 files changed
+24
-8
exchange/bitswap/network/interface.go
+1
-1
@@ -38,7 +38,7 @@ type BitSwapNetwork interface {
38
}
39
40
type MessageSender interface {
41
- SendMsg(bsmsg.BitSwapMessage) error
41
+ SendMsg(context.Context, bsmsg.BitSwapMessage) error
42
Close() error
43
}
44
exchange/bitswap/network/ipfs_impl.go
+20
-4
@@ -4,6 +4,7 @@ import (
4
"context"
5
"fmt"
6
"io"
7
+ "time"
8
9
bsmsg "github.com/ipfs/go-ipfs/exchange/bitswap/message"
10
@@ -20,6 +21,8 @@ import (
21
22
var log = logging.Logger("bitswap_network")
23
24
+var sendMessageTimeout = time.Minute * 10
25
+
26
// NewFromIpfsHost returns a BitSwapNetwork supported by underlying IPFS host
27
func NewFromIpfsHost(host host.Host, r routing.ContentRouting) BitSwapNetwork {
28
bitswapNetwork := impl{
@@ -53,11 +56,20 @@ func (s *streamMessageSender) Close() error {
56
return s.s.Close()
57
}
58
56
-func (s *streamMessageSender) SendMsg(msg bsmsg.BitSwapMessage) error {
57
- return msgToStream(s.s, msg)
59
+func (s *streamMessageSender) SendMsg(ctx context.Context, msg bsmsg.BitSwapMessage) error {
60
+ return msgToStream(ctx, s.s, msg)
61
}
62
60
-func msgToStream(s inet.Stream, msg bsmsg.BitSwapMessage) error {
63
+func msgToStream(ctx context.Context, s inet.Stream, msg bsmsg.BitSwapMessage) error {
64
+ deadline := time.Now().Add(sendMessageTimeout)
65
+ if dl, ok := ctx.Deadline(); ok {
66
+ deadline = dl
67
+ }
68
+
69
+ if err := s.SetWriteDeadline(deadline); err != nil {
70
+ log.Warningf("error setting deadline: %s", err)
71
+ }
72
+
73
switch s.Protocol() {
74
case ProtocolBitswap:
75
if err := msg.ToNetV1(s); err != nil {
@@ -72,6 +84,10 @@ func msgToStream(s inet.Stream, msg bsmsg.BitSwapMessage) error {
84
default:
85
return fmt.Errorf("unrecognized protocol on remote: %s", s.Protocol())
86
}
87
+
88
+ if err := s.SetWriteDeadline(time.Time{}); err != nil {
89
+ log.Warningf("error resetting deadline: %s", err)
90
+ }
91
return nil
92
}
93
@@ -107,7 +123,7 @@ func (bsnet *impl) SendMessage(
123
}
124
defer s.Close()
125
110
- return msgToStream(s, outgoing)
126
+ return msgToStream(ctx, s, outgoing)
127
}
128
129
func (bsnet *impl) SetDelegate(r Receiver) {
exchange/bitswap/testnet/virtual.go
+2
-2
@@ -119,8 +119,8 @@ type messagePasser struct {
119
ctx context.Context
120
}
121
122
-func (mp *messagePasser) SendMsg(m bsmsg.BitSwapMessage) error {
123
- return mp.net.SendMessage(mp.ctx, mp.local, mp.target, m)
122
+func (mp *messagePasser) SendMsg(ctx context.Context, m bsmsg.BitSwapMessage) error {
123
+ return mp.net.SendMessage(ctx, mp.local, mp.target, m)
124
}
125
126
func (mp *messagePasser) Close() error {
exchange/bitswap/wantmanager.go
+1
-1
@@ -196,7 +196,7 @@ func (mq *msgQueue) doWork(ctx context.Context) {
196
197
// send wantlist updates
198
for { // try to send this message until we fail.
199
- err := mq.sender.SendMsg(wlm)
199
+ err := mq.sender.SendMsg(ctx, wlm)
200
if err == nil {
201
return
202
}