@cryptotaxi247 / kubo / commits / 5806ac00c

rework FetchGraph to be less of a memory hog

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Feb 20, 2016 at 11:04 UTC 5806ac00c84bcb3301826aaad1f2dacafa54abf6
1 file changed +22 -17
merkledag/merkledag.go
+22 -17
@@ -149,6 +149,8 @@ func (ds *dagService) GetMany(ctx context.Context, keys []key.Key) (<-chan *Node
149 out := make(chan *Node)
150 errs := make(chan error, 1)
151 blocks := ds.Blocks.GetBlocks(ctx, keys)
152 + var count int
153 +
154 go func() {
155 defer close(out)
156 defer close(errs)
@@ -156,6 +158,9 @@ func (ds *dagService) GetMany(ctx context.Context, keys []key.Key) (<-chan *Node
158 select {
159 case b, ok := <-blocks:
160 if !ok {
161 + if count != len(keys) {
162 + errs <- fmt.Errorf("failed to fetch all nodes")
163 + }
164 return
165 }
166 nd, err := Decoded(b.Data)
@@ -165,6 +170,7 @@ func (ds *dagService) GetMany(ctx context.Context, keys []key.Key) (<-chan *Node
170 }
171 select {
172 case out <- nd:
173 + count++
174 case <-ctx.Done():
175 return
176 }
@@ -404,28 +410,27 @@ func EnumerateChildrenAsync(ctx context.Context, ds DAGService, root *Node, set
410 func fetchNodes(ctx context.Context, ds DAGService, in <-chan []key.Key, out chan<- *Node, errs chan<- error) {
411 defer close(out)
412
407 - get := func(g NodeGetter) {
408 - nd, err := g.Get(ctx)
409 - if err != nil {
413 + get := func(ks []key.Key) {
414 + nodes, errch := ds.GetMany(ctx, ks)
415 + for {
416 select {
411 - case errs <- err:
412 - case <-ctx.Done():
417 + case nd, ok := <-nodes:
418 + if !ok {
419 + return
420 + }
421 + select {
422 + case out <- nd:
423 + case <-ctx.Done():
424 + return
425 + }
426 + case err := <-errch:
427 + errs <- err
428 + return
429 }
414 - return
415 - }
416 -
417 - select {
418 - case out <- nd:
419 - case <-ctx.Done():
420 - return
430 }
431 }
432
433 for ks := range in {
425 - ng := GetNodes(ctx, ds, ks)
426 - for _, g := range ng {
427 - go get(g)
428 - }
434 + go get(ks)
435 }
430 -
436 }