feat(bitswap) ACTIVATE FULL CONCURRENCY cap'n
fix(bitswap) Put synchronously. Then notify async
Brian Tiger Chow committed
Sep 19, 2014 at 16:12 UTC
d0a5339547f876ed99bbe259439534dd83df980d
1 file changed
+15
-14
exchange/bitswap/bitswap.go
+15
-14
@@ -62,6 +62,7 @@ type bitswap struct {
62
func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error) {
63
64
ctx, cancelFunc := context.WithCancel(parent)
65
+ // TODO add to wantlist
66
promise := bs.notifications.Subscribe(ctx, k)
67
68
go func() {
@@ -69,8 +70,8 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
70
peersToQuery := bs.routing.FindProvidersAsync(ctx, k, maxProviders)
71
message := bsmsg.New()
72
message.AppendWanted(k)
72
- for i := range peersToQuery {
73
- func(p *peer.Peer) {
73
+ for iiiii := range peersToQuery {
74
+ go func(p *peer.Peer) {
75
response, err := bs.sender.SendRequest(ctx, p, message)
76
if err != nil {
77
return
@@ -84,13 +85,14 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
85
return
86
}
87
bs.ReceiveMessage(ctx, p, response)
87
- }(i)
88
+ }(iiiii)
89
}
90
}()
91
92
select {
93
case block := <-promise:
94
cancelFunc()
95
+ // TODO remove from wantlist
96
return &block, nil
97
case <-parent.Done():
98
return nil, parent.Err()
@@ -115,18 +117,17 @@ func (bs *bitswap) ReceiveMessage(
117
return nil, nil, errors.New("Received nil Message")
118
}
119
118
- bs.strategy.MessageReceived(p, incoming)
120
+ bs.strategy.MessageReceived(p, incoming) // FIRST
121
122
for _, block := range incoming.Blocks() {
121
- err := bs.blockstore.Put(block) // FIXME(brian): err ignored
122
- if err != nil {
123
- return nil, nil, err
124
- }
125
- bs.notifications.Publish(block)
126
- err = bs.HasBlock(ctx, block) // FIXME err ignored
127
- if err != nil {
128
- return nil, nil, err
123
+ // TODO verify blocks?
124
+ if err := bs.blockstore.Put(block); err != nil {
125
+ continue // FIXME(brian): err ignored
126
}
127
+ go bs.notifications.Publish(block)
128
+ go func() {
129
+ _ = bs.HasBlock(ctx, block) // FIXME err ignored
130
+ }()
131
}
132
133
for _, key := range incoming.Wantlist() {
@@ -148,7 +149,7 @@ func (bs *bitswap) ReceiveMessage(
149
// sent
150
func (bs *bitswap) send(ctx context.Context, p *peer.Peer, m bsmsg.BitSwapMessage) {
151
bs.sender.SendMessage(ctx, p, m)
151
- bs.strategy.MessageSent(p, m)
152
+ go bs.strategy.MessageSent(p, m)
153
}
154
155
func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block) {
@@ -157,7 +158,7 @@ func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block)
158
if bs.strategy.ShouldSendBlockToPeer(block.Key(), p) {
159
message := bsmsg.New()
160
message.AppendBlock(block)
160
- bs.send(ctx, p, message)
161
+ go bs.send(ctx, p, message)
162
}
163
}
164
}