feat(bitswap) implement GetBlocks
@whyrusleeping @jbenet License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Nov 19, 2014 at 23:27 UTC
11f2856d31ad8424d0024a92c3e62810800a7627
1 file changed
+23
-13
exchange/bitswap/bitswap.go
+23
-13
@@ -79,9 +79,7 @@ type bitswap struct {
79
}
80
81
// GetBlock attempts to retrieve a particular block from peers within the
82
-// deadline enforced by the context
83
-//
84
-// TODO ensure only one active request per key
82
+// deadline enforced by the context.
83
func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, error) {
84
85
// make sure to derive a new |ctx| and pass it to children. It's correct to
@@ -95,26 +93,36 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
93
log.Event(ctx, "GetBlockRequestBegin", &k)
94
defer log.Event(ctx, "GetBlockRequestEnd", &k)
95
98
- promise := bs.notifications.Subscribe(ctx, k)
99
-
100
- select {
101
- case bs.batchRequests <- []u.Key{k}:
102
- case <-parent.Done():
103
- return nil, parent.Err()
96
+ promise, err := bs.GetBlocks(parent, []u.Key{k})
97
+ if err != nil {
98
+ return nil, err
99
}
100
101
select {
102
case block := <-promise:
108
- bs.wantlist.Remove(k)
103
return &block, nil
104
case <-parent.Done():
105
return nil, parent.Err()
106
}
107
}
108
115
-func (bs *bitswap) GetBlocks(parent context.Context, ks []u.Key) (*blocks.Block, error) {
116
- // TODO: something smart
117
- return nil, nil
109
+// GetBlocks returns a channel where the caller may receive blocks that
110
+// correspond to the provided |keys|. Returns an error if BitSwap is unable to
111
+// begin this request within the deadline enforced by the context.
112
+//
113
+// NB: Your request remains open until the context expires. To conserve
114
+// resources, provide a context with a reasonably short deadline (ie. not one
115
+// that lasts throughout the lifetime of the server)
116
+func (bs *bitswap) GetBlocks(ctx context.Context, keys []u.Key) (<-chan blocks.Block, error) {
117
+ // TODO log the request
118
+
119
+ promise := bs.notifications.Subscribe(ctx, keys...)
120
+ select {
121
+ case bs.batchRequests <- keys:
122
+ return promise, nil
123
+ case <-ctx.Done():
124
+ return nil, ctx.Err()
125
+ }
126
}
127
128
func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) error {
@@ -155,6 +163,7 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
163
return nil
164
}
165
166
+// TODO ensure only one active request per key
167
func (bs *bitswap) run(ctx context.Context) {
168
169
// Every so often, we should resend out our current want list
@@ -238,6 +247,7 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
247
continue // FIXME(brian): err ignored
248
}
249
bs.notifications.Publish(block)
250
+ bs.wantlist.Remove(block.Key())
251
err := bs.HasBlock(ctx, block)
252
if err != nil {
253
log.Warningf("HasBlock errored: %s", err)