@cryptotaxi247 / kubo / commits / 99ae43202

start working getBlocks up the call chain

Jeromy committed Nov 21, 2014 at 01:15 UTC 99ae432021fd2987d9d7e018676605962eca0975
2 files changed +73 -15
blockservice/blockservice.go
+23 -3
@@ -96,9 +96,29 @@ func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, er
96 }
97 }
98
99 -func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) (<-chan blocks.Block, error) {
100 - // TODO:
101 - return nil, nil
99 +func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) <-chan *blocks.Block {
100 + out := make(chan *blocks.Block, 32)
101 + go func() {
102 + var toFetch []u.Key
103 + for _, k := range ks {
104 + datai, err := s.Datastore.Get(k.DsKey())
105 + if err == nil {
106 + log.Debug("Blockservice: Got data in datastore.")
107 + bdata, ok := datai.([]byte)
108 + if !ok {
109 + log.Criticalf("data associated with %s is not a []byte", k)
110 + continue
111 + }
112 + out <- &blocks.Block{
113 + Multihash: mh.Multihash(k),
114 + Data: bdata,
115 + }
116 + } else {
117 + toFetch = append(toFetch, k)
118 + }
119 + }
120 + }()
121 + return out
122 }
123
124 // DeleteBlock deletes a block in the blockservice from the datastore
merkledag/merkledag.go
+50 -12
@@ -289,18 +289,56 @@ func FetchGraph(ctx context.Context, root *Node, serv DAGService) chan struct{}
289 // Take advantage of blockservice/bitswap batched requests to fetch all
290 // child nodes of a given node
291 // TODO: finish this
292 -func (ds *dagService) BatchFetch(ctx context.Context, root *Node) error {
293 - var keys []u.Key
294 - for _, lnk := range root.Links {
295 - keys = append(keys, u.Key(lnk.Hash))
296 - }
292 +func (ds *dagService) BatchFetch(ctx context.Context, root *Node) chan struct{} {
293 + sig := make(chan struct{})
294 + go func() {
295 + var keys []u.Key
296 + for _, lnk := range root.Links {
297 + keys = append(keys, u.Key(lnk.Hash))
298 + }
299
298 - blocks, err := ds.Blocks.GetBlocks(ctx, keys)
299 - if err != nil {
300 - return err
301 - }
300 + blkchan := ds.Blocks.GetBlocks(ctx, keys)
301 +
302 + //
303 + next := 0
304 + seen := make(map[int]struct{})
305 + //
306 +
307 + for blk := range blkchan {
308 + for i, lnk := range root.Links {
309 +
310 + //
311 + seen[i] = struct{}{}
312 + //
313 +
314 + if u.Key(lnk.Hash) != blk.Key() {
315 + continue
316 + }
317 + nd, err := Decoded(blk.Data)
318 + if err != nil {
319 + log.Error("Got back bad block!")
320 + break
321 + }
322 + lnk.Node = nd
323 +
324 + //
325 + if next == i {
326 + sig <- struct{}{}
327 + next++
328 + for {
329 + if _, ok := seen[next]; ok {
330 + sig <- struct{}{}
331 + next++
332 + } else {
333 + break
334 + }
335 + }
336 + }
337 + //
338 + }
339 + }
340 + }()
341
303 - _ = blocks
304 - //what do i do with blocks?
305 - return nil
342 + // TODO: return a channel, and signal when the 'Next' readable block is available
343 + return sig
344 }