refac(bs:msg) let msg.Blocks() return []blocks
discard erroneous values wherever blocks cannot be nil, use value rather than pointer. only use pointers when absolutely necessary.
Brian Tiger Chow committed
Sep 13, 2014 at 16:56 UTC
78f0f5b0b975101355807d81462df986b583a31f
4 files changed
+42
-27
bitswap/bitswap.go
+17
-11
@@ -93,7 +93,7 @@ func (bs *BitSwap) GetBlock(k u.Key, timeout time.Duration) (
93
tleft := timeout - time.Now().Sub(begin)
94
provs_ch := bs.routing.FindProvidersAsync(k, 20, timeout)
95
96
- blockChannel := make(chan *blocks.Block)
96
+ blockChannel := make(chan blocks.Block)
97
after := time.After(tleft)
98
99
// TODO: when the data is received, shut down this for loop ASAP
@@ -106,7 +106,7 @@ func (bs *BitSwap) GetBlock(k u.Key, timeout time.Duration) (
106
return
107
}
108
select {
109
- case blockChannel <- blk:
109
+ case blockChannel <- *blk:
110
default:
111
}
112
}(p)
@@ -116,7 +116,7 @@ func (bs *BitSwap) GetBlock(k u.Key, timeout time.Duration) (
116
select {
117
case block := <-blockChannel:
118
close(blockChannel)
119
- return block, nil
119
+ return &block, nil
120
case <-after:
121
return nil, u.ErrTimeout
122
}
@@ -137,7 +137,7 @@ func (bs *BitSwap) getBlock(k u.Key, p *peer.Peer, timeout time.Duration) (*bloc
137
u.PErr("getBlock for '%s' timed out.\n", k.Pretty())
138
return nil, u.ErrTimeout
139
}
140
- return block, nil
140
+ return &block, nil
141
}
142
143
// HaveBlock announces the existance of a block to BitSwap, potentially sending
@@ -173,12 +173,7 @@ func (bs *BitSwap) handleMessages() {
173
}
174
175
if bsmsg.Blocks() != nil {
176
- for _, blkData := range bsmsg.Blocks() {
177
- blk, err := blocks.NewBlock(blkData)
178
- if err != nil {
179
- u.PErr("%v\n", err)
180
- continue
181
- }
176
+ for _, blk := range bsmsg.Blocks() {
177
go bs.blockReceive(mes.Peer, blk)
178
}
179
}
@@ -231,7 +226,7 @@ func (bs *BitSwap) peerWantsBlock(p *peer.Peer, want string) {
226
}
227
}
228
234
-func (bs *BitSwap) blockReceive(p *peer.Peer, blk *blocks.Block) {
229
+func (bs *BitSwap) blockReceive(p *peer.Peer, blk blocks.Block) {
230
u.DOut("blockReceive: %s\n", blk.Key().Pretty())
231
err := bs.datastore.Put(ds.NewKey(string(blk.Key())), blk.Data)
232
if err != nil {
@@ -286,5 +281,16 @@ func (bs *BitSwap) SetStrategy(sf StrategyFunc) {
281
func (bs *BitSwap) ReceiveMessage(
282
ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
283
bsmsg.BitSwapMessage, *peer.Peer, error) {
284
+ if incoming.Blocks() != nil {
285
+ for _, block := range incoming.Blocks() {
286
+ go bs.blockReceive(sender, block)
287
+ }
288
+ }
289
+
290
+ if incoming.Wantlist() != nil {
291
+ for _, want := range incoming.Wantlist() {
292
+ go bs.peerWantsBlock(sender, want)
293
+ }
294
+ }
295
return nil, nil, errors.New("TODO implement")
296
}
bitswap/message/message.go
+11
-3
@@ -15,7 +15,7 @@ import (
15
16
type BitSwapMessage interface {
17
Wantlist() []string
18
- Blocks() [][]byte
18
+ Blocks() []blocks.Block
19
AppendWanted(k u.Key)
20
AppendBlock(b *blocks.Block)
21
Exportable
@@ -46,8 +46,16 @@ func (m *message) Wantlist() []string {
46
}
47
48
// TODO(brian): convert these into blocks
49
-func (m *message) Blocks() [][]byte {
50
- return m.pb.Blocks
49
+func (m *message) Blocks() []blocks.Block {
50
+ bs := make([]blocks.Block, len(m.pb.Blocks))
51
+ for _, data := range m.pb.Blocks {
52
+ b, err := blocks.NewBlock(data)
53
+ if err != nil {
54
+ continue
55
+ }
56
+ bs = append(bs, *b)
57
+ }
58
+ return bs
59
}
60
61
func (m *message) AppendWanted(k u.Key) {
bitswap/notifications/notifications.go
+6
-6
@@ -9,8 +9,8 @@ import (
9
)
10
11
type PubSub interface {
12
- Publish(block *blocks.Block)
13
- Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block
12
+ Publish(block blocks.Block)
13
+ Subscribe(ctx context.Context, k u.Key) <-chan blocks.Block
14
Shutdown()
15
}
16
@@ -23,7 +23,7 @@ type impl struct {
23
wrapped pubsub.PubSub
24
}
25
26
-func (ps *impl) Publish(block *blocks.Block) {
26
+func (ps *impl) Publish(block blocks.Block) {
27
topic := string(block.Key())
28
ps.wrapped.Pub(block, topic)
29
}
@@ -31,15 +31,15 @@ func (ps *impl) Publish(block *blocks.Block) {
31
// Subscribe returns a one-time use |blockChannel|. |blockChannel| returns nil
32
// if the |ctx| times out or is cancelled. Then channel is closed after the
33
// block given by |k| is sent.
34
-func (ps *impl) Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block {
34
+func (ps *impl) Subscribe(ctx context.Context, k u.Key) <-chan blocks.Block {
35
topic := string(k)
36
subChan := ps.wrapped.SubOnce(topic)
37
- blockChannel := make(chan *blocks.Block)
37
+ blockChannel := make(chan blocks.Block)
38
go func() {
39
defer close(blockChannel)
40
select {
41
case val := <-subChan:
42
- block, ok := val.(*blocks.Block)
42
+ block, ok := val.(blocks.Block)
43
if ok {
44
blockChannel <- block
45
}
bitswap/notifications/notifications_test.go
+8
-7
@@ -34,19 +34,20 @@ func TestCarryOnWhenDeadlineExpires(t *testing.T) {
34
35
n := New()
36
defer n.Shutdown()
37
- blockChannel := n.Subscribe(fastExpiringCtx, getBlockOrFail(t, "A Missed Connection").Key())
37
+ block := getBlockOrFail(t, "A Missed Connection")
38
+ blockChannel := n.Subscribe(fastExpiringCtx, block.Key())
39
40
assertBlockChannelNil(t, blockChannel)
41
}
42
42
-func assertBlockChannelNil(t *testing.T, blockChannel <-chan *blocks.Block) {
43
- blockReceived := <-blockChannel
44
- if blockReceived != nil {
43
+func assertBlockChannelNil(t *testing.T, blockChannel <-chan blocks.Block) {
44
+ _, ok := <-blockChannel
45
+ if ok {
46
t.Fail()
47
}
48
}
49
49
-func assertBlocksEqual(t *testing.T, a, b *blocks.Block) {
50
+func assertBlocksEqual(t *testing.T, a, b blocks.Block) {
51
if !bytes.Equal(a.Data, b.Data) {
52
t.Fail()
53
}
@@ -55,10 +56,10 @@ func assertBlocksEqual(t *testing.T, a, b *blocks.Block) {
56
}
57
}
58
58
-func getBlockOrFail(t *testing.T, msg string) *blocks.Block {
59
+func getBlockOrFail(t *testing.T, msg string) blocks.Block {
60
block, blockCreationErr := blocks.NewBlock([]byte(msg))
61
if blockCreationErr != nil {
62
t.Fail()
63
}
63
- return block
64
+ return *block
65
}