feat(bitswap:network) propagate errors up the stack
Rather than pushing errors back down to lower layers, propagate the errors upward. This commit adds a `ReceiveError` method to BitSwap's network receiver. Still TODO: rm the error return value from: net.service.handler.HandleMessage This is inspired by delegation patterns in found in the wild.
Brian Tiger Chow committed
Sep 22, 2014 at 12:34 UTC
0e494690b33283ea893d1e2892b1246824cb2bcf
5 files changed
+39
-46
exchange/bitswap/bitswap.go
+11
-6
@@ -1,8 +1,6 @@
1
package bitswap
2
3
import (
4
- "errors"
5
-
4
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
6
@@ -120,14 +118,16 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
118
119
// TODO(brian): handle errors
120
func (bs *bitswap) ReceiveMessage(ctx context.Context, p *peer.Peer, incoming bsmsg.BitSwapMessage) (
123
- *peer.Peer, bsmsg.BitSwapMessage, error) {
121
+ *peer.Peer, bsmsg.BitSwapMessage) {
122
u.DOut("ReceiveMessage from %v\n", p.Key().Pretty())
123
124
if p == nil {
127
- return nil, nil, errors.New("Received nil Peer")
125
+ // TODO propagate the error upward
126
+ return nil, nil
127
}
128
if incoming == nil {
130
- return nil, nil, errors.New("Received nil Message")
129
+ // TODO propagate the error upward
130
+ return nil, nil
131
}
132
133
bs.strategy.MessageReceived(p, incoming) // FIRST
@@ -157,7 +157,12 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p *peer.Peer, incoming bs
157
}
158
}
159
defer bs.strategy.MessageSent(p, message)
160
- return p, message, nil
160
+ return p, message
161
+}
162
+
163
+func (bs *bitswap) ReceiveError(err error) {
164
+ // TODO log the network error
165
+ // TODO bubble the network error up to the parent context/error logger
166
}
167
168
// send strives to ensure that accounting is always performed when a message is
exchange/bitswap/network/interface.go
+3
-1
@@ -33,7 +33,9 @@ type Adapter interface {
33
type Receiver interface {
34
ReceiveMessage(
35
ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
36
- destination *peer.Peer, outgoing bsmsg.BitSwapMessage, err error)
36
+ destination *peer.Peer, outgoing bsmsg.BitSwapMessage)
37
+
38
+ ReceiveError(error)
39
}
40
41
// TODO(brian): move this to go-ipfs/net package
exchange/bitswap/network/net_message_adapter.go
+6
-9
@@ -1,8 +1,6 @@
1
package network
2
3
import (
4
- "errors"
5
-
4
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
6
bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
@@ -34,18 +32,16 @@ func (adapter *impl) HandleMessage(
32
ctx context.Context, incoming netmsg.NetMessage) (netmsg.NetMessage, error) {
33
34
if adapter.receiver == nil {
37
- return nil, errors.New("No receiver. NetMessage dropped")
35
+ return nil, nil
36
}
37
38
received, err := bsmsg.FromNet(incoming)
39
if err != nil {
42
- return nil, err
40
+ adapter.receiver.ReceiveError(err)
41
+ return nil, nil
42
}
43
45
- p, bsmsg, err := adapter.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
46
- if err != nil {
47
- return nil, err
48
- }
44
+ p, bsmsg := adapter.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
45
46
// TODO(brian): put this in a helper function
47
if bsmsg == nil || p == nil {
@@ -54,7 +50,8 @@ func (adapter *impl) HandleMessage(
50
51
outgoing, err := bsmsg.ToNet(p)
52
if err != nil {
57
- return nil, err
53
+ adapter.receiver.ReceiveError(err)
54
+ return nil, nil
55
}
56
57
return outgoing, nil
exchange/bitswap/testnet/network.go
+5
-19
@@ -76,18 +76,7 @@ func (n *network) deliver(
76
return errors.New("Invalid input")
77
}
78
79
- nextPeer, nextMsg, err := r.ReceiveMessage(context.TODO(), from, message)
80
- if err != nil {
81
-
82
- // TODO should this error be returned across network boundary?
83
-
84
- // TODO this raises an interesting question about network contract. How
85
- // can the network be expected to behave under different failure
86
- // conditions? What if peer is unreachable? Will we know if messages
87
- // aren't delivered?
88
-
89
- return err
90
- }
79
+ nextPeer, nextMsg := r.ReceiveMessage(context.TODO(), from, message)
80
81
if (nextPeer == nil && nextMsg != nil) || (nextMsg == nil && nextPeer != nil) {
82
return errors.New("Malformed client request")
@@ -119,15 +108,12 @@ func (n *network) SendRequest(
108
if !ok {
109
return nil, errors.New("Cannot locate peer on network")
110
}
122
- nextPeer, nextMsg, err := r.ReceiveMessage(context.TODO(), from, message)
123
- if err != nil {
124
- return nil, err
125
- // TODO return nil, NoResponse
126
- }
111
+ nextPeer, nextMsg := r.ReceiveMessage(context.TODO(), from, message)
112
113
// TODO dedupe code
114
if (nextPeer == nil && nextMsg != nil) || (nextMsg == nil && nextPeer != nil) {
130
- return nil, errors.New("Malformed client request")
115
+ r.ReceiveError(errors.New("Malformed client request"))
116
+ return nil, nil
117
}
118
119
// TODO dedupe code
@@ -144,7 +130,7 @@ func (n *network) SendRequest(
130
}
131
n.deliver(nextReceiver, nextPeer, nextMsg)
132
}()
147
- return nil, NoResponse
133
+ return nil, nil
134
}
135
return nextMsg, nil
136
}
exchange/bitswap/testnet/network_test.go
+14
-11
@@ -26,7 +26,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
26
ctx context.Context,
27
from *peer.Peer,
28
incoming bsmsg.BitSwapMessage) (
29
- *peer.Peer, bsmsg.BitSwapMessage, error) {
29
+ *peer.Peer, bsmsg.BitSwapMessage) {
30
31
t.Log("Recipient received a message from the network")
32
@@ -35,7 +35,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
35
m := bsmsg.New()
36
m.AppendBlock(testutil.NewBlockOrFail(t, expectedStr))
37
38
- return from, m, nil
38
+ return from, m
39
}))
40
41
t.Log("Build a message and send a synchronous request to recipient")
@@ -74,19 +74,19 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
74
ctx context.Context,
75
fromWaiter *peer.Peer,
76
msgFromWaiter bsmsg.BitSwapMessage) (
77
- *peer.Peer, bsmsg.BitSwapMessage, error) {
77
+ *peer.Peer, bsmsg.BitSwapMessage) {
78
79
msgToWaiter := bsmsg.New()
80
msgToWaiter.AppendBlock(testutil.NewBlockOrFail(t, expectedStr))
81
82
- return fromWaiter, msgToWaiter, nil
82
+ return fromWaiter, msgToWaiter
83
}))
84
85
waiter.SetDelegate(lambda(func(
86
ctx context.Context,
87
fromResponder *peer.Peer,
88
msgFromResponder bsmsg.BitSwapMessage) (
89
- *peer.Peer, bsmsg.BitSwapMessage, error) {
89
+ *peer.Peer, bsmsg.BitSwapMessage) {
90
91
// TODO assert that this came from the correct peer and that the message contents are as expected
92
ok := false
@@ -101,7 +101,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
101
t.Fatal("Message not received from the responder")
102
103
}
104
- return nil, nil, nil
104
+ return nil, nil
105
}))
106
107
messageSentAsync := bsmsg.New()
@@ -116,7 +116,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
116
}
117
118
type receiverFunc func(ctx context.Context, p *peer.Peer,
119
- incoming bsmsg.BitSwapMessage) (*peer.Peer, bsmsg.BitSwapMessage, error)
119
+ incoming bsmsg.BitSwapMessage) (*peer.Peer, bsmsg.BitSwapMessage)
120
121
// lambda returns a Receiver instance given a receiver function
122
func lambda(f receiverFunc) bsnet.Receiver {
@@ -126,13 +126,16 @@ func lambda(f receiverFunc) bsnet.Receiver {
126
}
127
128
type lambdaImpl struct {
129
- f func(ctx context.Context, p *peer.Peer,
130
- incoming bsmsg.BitSwapMessage) (
131
- *peer.Peer, bsmsg.BitSwapMessage, error)
129
+ f func(ctx context.Context, p *peer.Peer, incoming bsmsg.BitSwapMessage) (
130
+ *peer.Peer, bsmsg.BitSwapMessage)
131
}
132
133
func (lam *lambdaImpl) ReceiveMessage(ctx context.Context,
134
p *peer.Peer, incoming bsmsg.BitSwapMessage) (
136
- *peer.Peer, bsmsg.BitSwapMessage, error) {
135
+ *peer.Peer, bsmsg.BitSwapMessage) {
136
return lam.f(ctx, p, incoming)
137
}
138
+
139
+func (lam *lambdaImpl) ReceiveError(err error) {
140
+ // TODO log error
141
+}