address concerns from PR
License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>
Jeromy committed
Jul 8, 2015 at 08:48 UTC
dc3b9ed1407b05ee536db6399fce9f04b1963d7d
2 files changed
+94
-94
merkledag/merkledag.go
+81
-78
@@ -122,84 +122,7 @@ func (n *dagService) Remove(nd *Node) error {
122
123
// FetchGraph fetches all nodes that are children of the given node
124
func FetchGraph(ctx context.Context, root *Node, serv DAGService) error {
125
- toprocess := make(chan []key.Key, 8)
126
- nodes := make(chan *Node, 8)
127
- errs := make(chan error, 1)
128
-
129
- ctx, cancel := context.WithCancel(ctx)
130
- defer cancel()
131
- defer close(toprocess)
132
-
133
- go fetchNodes(ctx, serv, toprocess, nodes, errs)
134
-
135
- nodes <- root
136
- live := 1
137
-
138
- for {
139
- select {
140
- case nd, ok := <-nodes:
141
- if !ok {
142
- return nil
143
- }
144
-
145
- var keys []key.Key
146
- for _, lnk := range nd.Links {
147
- keys = append(keys, key.Key(lnk.Hash))
148
- }
149
- keys = dedupeKeys(keys)
150
-
151
- // keep track of open request, when zero, we're done
152
- live += len(keys) - 1
153
-
154
- if live == 0 {
155
- return nil
156
- }
157
-
158
- if len(keys) > 0 {
159
- select {
160
- case toprocess <- keys:
161
- case <-ctx.Done():
162
- return ctx.Err()
163
- }
164
- }
165
- case err := <-errs:
166
- return err
167
- case <-ctx.Done():
168
- return ctx.Err()
169
- }
170
- }
171
-}
172
-
173
-func fetchNodes(ctx context.Context, ds DAGService, in <-chan []key.Key, out chan<- *Node, errs chan<- error) {
174
- defer close(out)
175
- for {
176
- select {
177
- case ks, ok := <-in:
178
- if !ok {
179
- return
180
- }
181
-
182
- ng := ds.GetNodes(ctx, ks)
183
- for _, g := range ng {
184
- go func(g NodeGetter) {
185
- nd, err := g.Get(ctx)
186
- if err != nil {
187
- select {
188
- case errs <- err:
189
- case <-ctx.Done():
190
- }
191
- return
192
- }
193
-
194
- select {
195
- case out <- nd:
196
- case <-ctx.Done():
197
- return
198
- }
199
- }(g)
200
- }
201
- }
202
- }
125
+ return EnumerateChildrenAsync(ctx, serv, root, key.NewKeySet())
126
}
127
128
// FindLinks searches this nodes links for the given key,
@@ -383,3 +306,83 @@ func EnumerateChildren(ctx context.Context, ds DAGService, root *Node, set key.K
306
}
307
return nil
308
}
309
+
310
+func EnumerateChildrenAsync(ctx context.Context, ds DAGService, root *Node, set key.KeySet) error {
311
+ toprocess := make(chan []key.Key, 8)
312
+ nodes := make(chan *Node, 8)
313
+ errs := make(chan error, 1)
314
+
315
+ ctx, cancel := context.WithCancel(ctx)
316
+ defer cancel()
317
+ defer close(toprocess)
318
+
319
+ go fetchNodes(ctx, ds, toprocess, nodes, errs)
320
+
321
+ nodes <- root
322
+ live := 1
323
+
324
+ for {
325
+ select {
326
+ case nd, ok := <-nodes:
327
+ if !ok {
328
+ return nil
329
+ }
330
+ // a node has been fetched
331
+ live--
332
+
333
+ var keys []key.Key
334
+ for _, lnk := range nd.Links {
335
+ k := key.Key(lnk.Hash)
336
+ if !set.Has(k) {
337
+ set.Add(k)
338
+ live++
339
+ keys = append(keys, k)
340
+ }
341
+ }
342
+
343
+ if live == 0 {
344
+ return nil
345
+ }
346
+
347
+ if len(keys) > 0 {
348
+ select {
349
+ case toprocess <- keys:
350
+ case <-ctx.Done():
351
+ return ctx.Err()
352
+ }
353
+ }
354
+ case err := <-errs:
355
+ return err
356
+ case <-ctx.Done():
357
+ return ctx.Err()
358
+ }
359
+ }
360
+}
361
+
362
+func fetchNodes(ctx context.Context, ds DAGService, in <-chan []key.Key, out chan<- *Node, errs chan<- error) {
363
+ defer close(out)
364
+
365
+ get := func(g NodeGetter) {
366
+ nd, err := g.Get(ctx)
367
+ if err != nil {
368
+ select {
369
+ case errs <- err:
370
+ case <-ctx.Done():
371
+ }
372
+ return
373
+ }
374
+
375
+ select {
376
+ case out <- nd:
377
+ case <-ctx.Done():
378
+ return
379
+ }
380
+ }
381
+
382
+ for ks := range in {
383
+ ng := ds.GetNodes(ctx, ks)
384
+ for _, g := range ng {
385
+ go get(g)
386
+ }
387
+ }
388
+}
merkledag/merkledag_test.go
+13
-16
@@ -299,38 +299,35 @@ func TestCantGet(t *testing.T) {
299
}
300
301
func TestFetchGraph(t *testing.T) {
302
- bsi := bstest.Mocks(t, 1)[0]
303
- ds := NewDAGService(bsi)
302
+ var dservs []DAGService
303
+ bsis := bstest.Mocks(t, 2)
304
+ for _, bsi := range bsis {
305
+ dservs = append(dservs, NewDAGService(bsi))
306
+ }
307
308
read := io.LimitReader(u.NewTimeSeededRand(), 1024*32)
309
spl := &chunk.SizeSplitter{512}
310
308
- root, err := imp.BuildDagFromReader(read, ds, spl, nil)
311
+ root, err := imp.BuildDagFromReader(read, dservs[0], spl, nil)
312
if err != nil {
313
t.Fatal(err)
314
}
315
313
- err = FetchGraph(context.TODO(), root, ds)
316
+ err = FetchGraph(context.TODO(), root, dservs[1])
317
if err != nil {
318
t.Fatal(err)
319
}
317
-}
318
-
319
-func TestFetchGraphOther(t *testing.T) {
320
- var dservs []DAGService
321
- for _, bsi := range bstest.Mocks(t, 2) {
322
- dservs = append(dservs, NewDAGService(bsi))
323
- }
324
-
325
- read := io.LimitReader(u.NewTimeSeededRand(), 1024*32)
326
- spl := &chunk.SizeSplitter{512}
320
328
- root, err := imp.BuildDagFromReader(read, dservs[0], spl, nil)
321
+ // create an offline dagstore and ensure all blocks were fetched
322
+ bs, err := bserv.New(bsis[1].Blockstore, offline.Exchange(bsis[1].Blockstore))
323
if err != nil {
324
t.Fatal(err)
325
}
326
333
- err = FetchGraph(context.TODO(), root, dservs[1])
327
+ offline_ds := NewDAGService(bs)
328
+ ks := key.NewKeySet()
329
+
330
+ err = EnumerateChildren(context.Background(), offline_ds, root, ks)
331
if err != nil {
332
t.Fatal(err)
333
}