@cryptotaxi247 / kubo / commits / d53deebad

wire GetBlocks into blockservice

Jeromy committed Nov 21, 2014 at 06:40 UTC d53deebada5f4ff2012f4fe9403862186200e28e
15 files changed +88 -55
blockservice/blockservice.go
+11 -1
@@ -68,7 +68,7 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
68 // consider moving this to an sync process.
69 if s.Remote != nil {
70 ctx := context.TODO()
71 - err = s.Remote.HasBlock(ctx, *b)
71 + err = s.Remote.HasBlock(ctx, b)
72 }
73 return k, err
74 }
@@ -98,6 +98,7 @@ func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, er
98 func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) <-chan *blocks.Block {
99 out := make(chan *blocks.Block, 32)
100 go func() {
101 + defer close(out)
102 var toFetch []u.Key
103 for _, k := range ks {
104 block, err := s.Blockstore.Get(k)
@@ -108,6 +109,15 @@ func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) <-chan *blocks
109 log.Debug("Blockservice: Got data in datastore.")
110 out <- block
111 }
112 +
113 + nblocks, err := s.Remote.GetBlocks(ctx, toFetch)
114 + if err != nil {
115 + log.Errorf("Error with GetBlocks: %s", err)
116 + return
117 + }
118 + for blk := range nblocks {
119 + out <- blk
120 + }
121 }()
122 return out
123 }
exchange/bitswap/bitswap.go
+6 -6
@@ -108,7 +108,7 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
108
109 select {
110 case block := <-promise:
111 - return &block, nil
111 + return block, nil
112 case <-parent.Done():
113 return nil, parent.Err()
114 }
@@ -122,7 +122,7 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
122 // NB: Your request remains open until the context expires. To conserve
123 // resources, provide a context with a reasonably short deadline (ie. not one
124 // that lasts throughout the lifetime of the server)
125 -func (bs *bitswap) GetBlocks(ctx context.Context, keys []u.Key) (<-chan blocks.Block, error) {
125 +func (bs *bitswap) GetBlocks(ctx context.Context, keys []u.Key) (<-chan *blocks.Block, error) {
126 // TODO log the request
127
128 promise := bs.notifications.Subscribe(ctx, keys...)
@@ -213,7 +213,7 @@ func (bs *bitswap) loop(parent context.Context) {
213
214 // HasBlock announces the existance of a block to this bitswap service. The
215 // service will potentially notify its peers.
216 -func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
216 +func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
217 log.Debugf("Has Block %v", blk.Key())
218 bs.wantlist.Remove(blk.Key())
219 bs.sendToPeersThatWant(ctx, blk)
@@ -244,7 +244,7 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
244
245 for _, block := range incoming.Blocks() {
246 // TODO verify blocks?
247 - if err := bs.blockstore.Put(&block); err != nil {
247 + if err := bs.blockstore.Put(block); err != nil {
248 log.Criticalf("error putting block: %s", err)
249 continue // FIXME(brian): err ignored
250 }
@@ -267,7 +267,7 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
267 if block, errBlockNotFound := bs.blockstore.Get(key); errBlockNotFound != nil {
268 continue
269 } else {
270 - message.AddBlock(*block)
270 + message.AddBlock(block)
271 }
272 }
273 }
@@ -290,7 +290,7 @@ func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage
290 bs.strategy.MessageSent(p, m)
291 }
292
293 -func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block) {
293 +func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block) {
294 log.Debugf("Sending %v to peers that want it", block.Key())
295
296 for _, p := range bs.strategy.Peers() {
exchange/bitswap/bitswap_test.go
+7 -7
@@ -83,7 +83,7 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
83 if err := hasBlock.blockstore.Put(block); err != nil {
84 t.Fatal(err)
85 }
86 - if err := hasBlock.exchange.HasBlock(context.Background(), *block); err != nil {
86 + if err := hasBlock.exchange.HasBlock(context.Background(), block); err != nil {
87 t.Fatal(err)
88 }
89
@@ -140,7 +140,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
140 first := instances[0]
141 for _, b := range blocks {
142 first.blockstore.Put(b)
143 - first.exchange.HasBlock(context.Background(), *b)
143 + first.exchange.HasBlock(context.Background(), b)
144 rs.Announce(first.peer, b.Key())
145 }
146
@@ -212,7 +212,7 @@ func TestSendToWantingPeer(t *testing.T) {
212 beta := bg.Next()
213 t.Logf("Peer %v announes availability of %v\n", w.peer, beta.Key())
214 ctx, _ = context.WithTimeout(context.Background(), timeout)
215 - if err := w.blockstore.Put(&beta); err != nil {
215 + if err := w.blockstore.Put(beta); err != nil {
216 t.Fatal(err)
217 }
218 w.exchange.HasBlock(ctx, beta)
@@ -225,7 +225,7 @@ func TestSendToWantingPeer(t *testing.T) {
225
226 t.Logf("%v announces availability of %v\n", o.peer, alpha.Key())
227 ctx, _ = context.WithTimeout(context.Background(), timeout)
228 - if err := o.blockstore.Put(&alpha); err != nil {
228 + if err := o.blockstore.Put(alpha); err != nil {
229 t.Fatal(err)
230 }
231 o.exchange.HasBlock(ctx, alpha)
@@ -254,16 +254,16 @@ type BlockGenerator struct {
254 seq int
255 }
256
257 -func (bg *BlockGenerator) Next() blocks.Block {
257 +func (bg *BlockGenerator) Next() *blocks.Block {
258 bg.seq++
259 - return *blocks.NewBlock([]byte(string(bg.seq)))
259 + return blocks.NewBlock([]byte(string(bg.seq)))
260 }
261
262 func (bg *BlockGenerator) Blocks(n int) []*blocks.Block {
263 blocks := make([]*blocks.Block, 0)
264 for i := 0; i < n; i++ {
265 b := bg.Next()
266 - blocks = append(blocks, &b)
266 + blocks = append(blocks, b)
267 }
268 return blocks
269 }
exchange/bitswap/message/message.go
+10 -10
@@ -19,7 +19,7 @@ type BitSwapMessage interface {
19 Wantlist() []u.Key
20
21 // Blocks returns a slice of unique blocks
22 - Blocks() []blocks.Block
22 + Blocks() []*blocks.Block
23
24 // AddWanted adds the key to the Wantlist.
25 //
@@ -32,7 +32,7 @@ type BitSwapMessage interface {
32 // implies Priority(A) > Priority(B)
33 AddWanted(u.Key)
34
35 - AddBlock(blocks.Block)
35 + AddBlock(*blocks.Block)
36 Exportable
37 }
38
@@ -42,14 +42,14 @@ type Exportable interface {
42 }
43
44 type impl struct {
45 - existsInWantlist map[u.Key]struct{} // map to detect duplicates
46 - wantlist []u.Key // slice to preserve ordering
47 - blocks map[u.Key]blocks.Block // map to detect duplicates
45 + existsInWantlist map[u.Key]struct{} // map to detect duplicates
46 + wantlist []u.Key // slice to preserve ordering
47 + blocks map[u.Key]*blocks.Block // map to detect duplicates
48 }
49
50 func New() BitSwapMessage {
51 return &impl{
52 - blocks: make(map[u.Key]blocks.Block),
52 + blocks: make(map[u.Key]*blocks.Block),
53 existsInWantlist: make(map[u.Key]struct{}),
54 wantlist: make([]u.Key, 0),
55 }
@@ -62,7 +62,7 @@ func newMessageFromProto(pbm pb.Message) BitSwapMessage {
62 }
63 for _, d := range pbm.GetBlocks() {
64 b := blocks.NewBlock(d)
65 - m.AddBlock(*b)
65 + m.AddBlock(b)
66 }
67 return m
68 }
@@ -71,8 +71,8 @@ func (m *impl) Wantlist() []u.Key {
71 return m.wantlist
72 }
73
74 -func (m *impl) Blocks() []blocks.Block {
75 - bs := make([]blocks.Block, 0)
74 +func (m *impl) Blocks() []*blocks.Block {
75 + bs := make([]*blocks.Block, 0)
76 for _, block := range m.blocks {
77 bs = append(bs, block)
78 }
@@ -88,7 +88,7 @@ func (m *impl) AddWanted(k u.Key) {
88 m.wantlist = append(m.wantlist, k)
89 }
90
91 -func (m *impl) AddBlock(b blocks.Block) {
91 +func (m *impl) AddBlock(b *blocks.Block) {
92 m.blocks[b.Key()] = b
93 }
94
exchange/bitswap/message/message_test.go
+7 -7
@@ -42,7 +42,7 @@ func TestAppendBlock(t *testing.T) {
42 m := New()
43 for _, str := range strs {
44 block := blocks.NewBlock([]byte(str))
45 - m.AddBlock(*block)
45 + m.AddBlock(block)
46 }
47
48 // assert strings are in proto message
@@ -133,10 +133,10 @@ func TestToNetFromNetPreservesWantList(t *testing.T) {
133 func TestToAndFromNetMessage(t *testing.T) {
134
135 original := New()
136 - original.AddBlock(*blocks.NewBlock([]byte("W")))
137 - original.AddBlock(*blocks.NewBlock([]byte("E")))
138 - original.AddBlock(*blocks.NewBlock([]byte("F")))
139 - original.AddBlock(*blocks.NewBlock([]byte("M")))
136 + original.AddBlock(blocks.NewBlock([]byte("W")))
137 + original.AddBlock(blocks.NewBlock([]byte("E")))
138 + original.AddBlock(blocks.NewBlock([]byte("F")))
139 + original.AddBlock(blocks.NewBlock([]byte("M")))
140
141 p := peer.WithIDString("X")
142 netmsg, err := original.ToNet(p)
@@ -180,8 +180,8 @@ func TestDuplicates(t *testing.T) {
180 t.Fatal("Duplicate in BitSwapMessage")
181 }
182
183 - msg.AddBlock(*b)
184 - msg.AddBlock(*b)
183 + msg.AddBlock(b)
184 + msg.AddBlock(b)
185 if len(msg.Blocks()) != 1 {
186 t.Fatal("Duplicate in BitSwapMessage")
187 }
exchange/bitswap/notifications/notifications.go
+6 -6
@@ -11,8 +11,8 @@ import (
11 const bufferSize = 16
12
13 type PubSub interface {
14 - Publish(block blocks.Block)
15 - Subscribe(ctx context.Context, keys ...u.Key) <-chan blocks.Block
14 + Publish(block *blocks.Block)
15 + Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Block
16 Shutdown()
17 }
18
@@ -24,7 +24,7 @@ type impl struct {
24 wrapped pubsub.PubSub
25 }
26
27 -func (ps *impl) Publish(block blocks.Block) {
27 +func (ps *impl) Publish(block *blocks.Block) {
28 topic := string(block.Key())
29 ps.wrapped.Pub(block, topic)
30 }
@@ -32,18 +32,18 @@ func (ps *impl) Publish(block blocks.Block) {
32 // Subscribe returns a one-time use |blockChannel|. |blockChannel| returns nil
33 // if the |ctx| times out or is cancelled. Then channel is closed after the
34 // blocks given by |keys| are sent.
35 -func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan blocks.Block {
35 +func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Block {
36 topics := make([]string, 0)
37 for _, key := range keys {
38 topics = append(topics, string(key))
39 }
40 subChan := ps.wrapped.SubOnce(topics...)
41 - blockChannel := make(chan blocks.Block, 1) // buffered so the sender doesn't wait on receiver
41 + blockChannel := make(chan *blocks.Block, 1) // buffered so the sender doesn't wait on receiver
42 go func() {
43 defer close(blockChannel)
44 select {
45 case val := <-subChan:
46 - block, ok := val.(blocks.Block)
46 + block, ok := val.(*blocks.Block)
47 if ok {
48 blockChannel <- block
49 }
exchange/bitswap/notifications/notifications_test.go
+4 -4
@@ -16,13 +16,13 @@ func TestPublishSubscribe(t *testing.T) {
16 defer n.Shutdown()
17 ch := n.Subscribe(context.Background(), blockSent.Key())
18
19 - n.Publish(*blockSent)
19 + n.Publish(blockSent)
20 blockRecvd, ok := <-ch
21 if !ok {
22 t.Fail()
23 }
24
25 - assertBlocksEqual(t, blockRecvd, *blockSent)
25 + assertBlocksEqual(t, blockRecvd, blockSent)
26
27 }
28
@@ -39,14 +39,14 @@ func TestCarryOnWhenDeadlineExpires(t *testing.T) {
39 assertBlockChannelNil(t, blockChannel)
40 }
41
42 -func assertBlockChannelNil(t *testing.T, blockChannel <-chan blocks.Block) {
42 +func assertBlockChannelNil(t *testing.T, blockChannel <-chan *blocks.Block) {
43 _, ok := <-blockChannel
44 if ok {
45 t.Fail()
46 }
47 }
48
49 -func assertBlocksEqual(t *testing.T, a, b blocks.Block) {
49 +func assertBlocksEqual(t *testing.T, a, b *blocks.Block) {
50 if !bytes.Equal(a.Data, b.Data) {
51 t.Fail()
52 }
exchange/bitswap/strategy/strategy_test.go
+1 -1
@@ -30,7 +30,7 @@ func TestConsistentAccounting(t *testing.T) {
30
31 m := message.New()
32 content := []string{"this", "is", "message", "i"}
33 - m.AddBlock(*blocks.NewBlock([]byte(strings.Join(content, " "))))
33 + m.AddBlock(blocks.NewBlock([]byte(strings.Join(content, " "))))
34
35 sender.MessageSent(receiver.Peer, m)
36 receiver.MessageReceived(sender.Peer, m)
exchange/bitswap/testnet/network_test.go
+4 -4
@@ -33,7 +33,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
33 // TODO test contents of incoming message
34
35 m := bsmsg.New()
36 - m.AddBlock(*blocks.NewBlock([]byte(expectedStr)))
36 + m.AddBlock(blocks.NewBlock([]byte(expectedStr)))
37
38 return from, m
39 }))
@@ -41,7 +41,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
41 t.Log("Build a message and send a synchronous request to recipient")
42
43 message := bsmsg.New()
44 - message.AddBlock(*blocks.NewBlock([]byte("data")))
44 + message.AddBlock(blocks.NewBlock([]byte("data")))
45 response, err := initiator.SendRequest(
46 context.Background(), peer.WithID(idOfRecipient), message)
47 if err != nil {
@@ -77,7 +77,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
77 peer.Peer, bsmsg.BitSwapMessage) {
78
79 msgToWaiter := bsmsg.New()
80 - msgToWaiter.AddBlock(*blocks.NewBlock([]byte(expectedStr)))
80 + msgToWaiter.AddBlock(blocks.NewBlock([]byte(expectedStr)))
81
82 return fromWaiter, msgToWaiter
83 }))
@@ -105,7 +105,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
105 }))
106
107 messageSentAsync := bsmsg.New()
108 - messageSentAsync.AddBlock(*blocks.NewBlock([]byte("data")))
108 + messageSentAsync.AddBlock(blocks.NewBlock([]byte("data")))
109 errSending := waiter.SendMessage(
110 context.Background(), peer.WithID(idOfResponder), messageSentAsync)
111 if errSending != nil {
exchange/interface.go
+3 -1
@@ -16,9 +16,11 @@ type Interface interface {
16 // GetBlock returns the block associated with a given key.
17 GetBlock(context.Context, u.Key) (*blocks.Block, error)
18
19 + GetBlocks(context.Context, []u.Key) (<-chan *blocks.Block, error)
20 +
21 // TODO Should callers be concerned with whether the block was made
22 // available on the network?
21 - HasBlock(context.Context, blocks.Block) error
23 + HasBlock(context.Context, *blocks.Block) error
24
25 io.Closer
26 }
exchange/offline/offline.go
+5 -1
@@ -30,7 +30,7 @@ func (_ *offlineExchange) GetBlock(context.Context, u.Key) (*blocks.Block, error
30 }
31
32 // HasBlock always returns nil.
33 -func (_ *offlineExchange) HasBlock(context.Context, blocks.Block) error {
33 +func (_ *offlineExchange) HasBlock(context.Context, *blocks.Block) error {
34 return nil
35 }
36
@@ -38,3 +38,7 @@ func (_ *offlineExchange) HasBlock(context.Context, blocks.Block) error {
38 func (_ *offlineExchange) Close() error {
39 return nil
40 }
41 +
42 +func (_ *offlineExchange) GetBlocks(context.Context, []u.Key) (<-chan *blocks.Block, error) {
43 + return nil, OfflineMode
44 +}
exchange/offline/offline_test.go
+1 -1
@@ -21,7 +21,7 @@ func TestBlockReturnsErr(t *testing.T) {
21 func TestHasBlockReturnsNil(t *testing.T) {
22 off := Exchange()
23 block := blocks.NewBlock([]byte("data"))
24 - err := off.HasBlock(context.Background(), *block)
24 + err := off.HasBlock(context.Background(), block)
25 if err != nil {
26 t.Fatal("")
27 }
importer/importer_test.go
+1
@@ -69,6 +69,7 @@ func testFileConsistency(t *testing.T, bs chunk.BlockSplitter, nbytes int) {
69 if err != nil {
70 t.Fatal(err)
71 }
72 +
73 r, err := uio.NewDagReader(nd, nil)
74 if err != nil {
75 t.Fatal(err)
merkledag/merkledag.go
+1 -1
@@ -328,7 +328,7 @@ func (ds *dagService) BatchFetch(ctx context.Context, root *Node) <-chan *Node {
328 if next == i {
329 sig <- nd
330 next++
331 - for ; nodes[next] != nil; next++ {
331 + for ; next < len(nodes) && nodes[next] != nil; next++ {
332 sig <- nodes[next]
333 }
334 }
unixfs/io/dagreader.go
+21 -5
@@ -17,10 +17,11 @@ var ErrIsDir = errors.New("this dag node is a directory")
17
18 // DagReader provides a way to easily read the data contained in a dag.
19 type DagReader struct {
20 - serv mdag.DAGService
21 - node *mdag.Node
22 - buf io.Reader
23 - fetchChan <-chan *mdag.Node
20 + serv mdag.DAGService
21 + node *mdag.Node
22 + buf io.Reader
23 + fetchChan <-chan *mdag.Node
24 + linkPosition int
25 }
26
27 // NewDagReader creates a new reader object that reads the data represented by the given
@@ -37,11 +38,15 @@ func NewDagReader(n *mdag.Node, serv mdag.DAGService) (io.Reader, error) {
38 // Dont allow reading directories
39 return nil, ErrIsDir
40 case ftpb.Data_File:
41 + var fetchChan <-chan *mdag.Node
42 + if serv != nil {
43 + fetchChan = serv.BatchFetch(context.TODO(), n)
44 + }
45 return &DagReader{
46 node: n,
47 serv: serv,
48 buf: bytes.NewBuffer(pb.GetData()),
44 - fetchChan: serv.BatchFetch(context.TODO(), n),
49 + fetchChan: fetchChan,
50 }, nil
51 case ftpb.Data_Raw:
52 // Raw block will just be a single level, return a byte buffer
@@ -61,6 +66,17 @@ func (dr *DagReader) precalcNextBuf() error {
66 if !ok {
67 return io.EOF
68 }
69 + default:
70 + // Only used when fetchChan is nil,
71 + // which only happens when passed in a nil dagservice
72 + // TODO: this logic is hard to follow, do it better.
73 + // NOTE: the only time this code is used, is during the
74 + // importer tests, consider just changing those tests
75 + if dr.linkPosition >= len(dr.node.Links) {
76 + return io.EOF
77 + }
78 + nxt = dr.node.Links[dr.linkPosition].Node
79 + dr.linkPosition++
80 }
81
82 pb := new(ftpb.Data)