feat(util) ForwardNBlocks
License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Nov 21, 2014 at 13:09 UTC
438ffa1dd7a211c9b327f6759bc364406ba36a27
2 files changed
+85
util/async/forward.go
new
+30
@@ -0,0 +1,30 @@
1
+package async
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+ "github.com/jbenet/go-ipfs/blocks"
6
+)
7
+
8
+// ForwardN forwards up to |num| blocks to the returned channel.
9
+func ForwardN(ctx context.Context, in <-chan *blocks.Block, num int) <-chan *blocks.Block {
10
+ out := make(chan *blocks.Block)
11
+ go func() {
12
+ defer close(out)
13
+ for i := 0; i < num; i++ {
14
+ select {
15
+ case block, ok := <-in:
16
+ if !ok {
17
+ return
18
+ }
19
+ select {
20
+ case out <- block:
21
+ case <-ctx.Done():
22
+ return
23
+ }
24
+ case <-ctx.Done():
25
+ return
26
+ }
27
+ }
28
+ }()
29
+ return out
30
+}
util/async/forward_test.go
new
+55
@@ -0,0 +1,55 @@
1
+package async
2
+
3
+import (
4
+ "testing"
5
+
6
+ "code.google.com/p/go.net/context"
7
+ "github.com/jbenet/go-ipfs/blocks"
8
+)
9
+
10
+func TestForwardTwo(t *testing.T) {
11
+ const n = 2
12
+ in := make(chan *blocks.Block, n)
13
+ ctx := context.Background()
14
+ out := ForwardN(ctx, in, n)
15
+
16
+ in <- blocks.NewBlock([]byte("one"))
17
+ in <- blocks.NewBlock([]byte("two"))
18
+
19
+ _ = <-out // 1
20
+ _ = <-out // 2
21
+
22
+ _, ok := <-out // closed
23
+ if !ok {
24
+ return
25
+ }
26
+ t.Fail()
27
+}
28
+
29
+func TestCloseInput(t *testing.T) {
30
+ const n = 2
31
+ in := make(chan *blocks.Block, 0)
32
+ ctx := context.Background()
33
+ out := ForwardN(ctx, in, n)
34
+
35
+ close(in)
36
+ _, ok := <-out // closed
37
+ if !ok {
38
+ return
39
+ }
40
+ t.Fatal("input channel closed, but output channel not")
41
+
42
+}
43
+
44
+func TestContextClosedWhenBlockingOnInput(t *testing.T) {
45
+ const n = 1 // but we won't ever send a block
46
+ ctx, cancel := context.WithCancel(context.Background())
47
+ out := ForwardN(ctx, make(chan *blocks.Block), n)
48
+
49
+ cancel() // before sending anything
50
+ _, ok := <-out
51
+ if !ok {
52
+ return
53
+ }
54
+ t.Fail()
55
+}