fix(bitswap/notifications) subscribe to many
License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Nov 21, 2014 at 17:18 UTC
03324f776517a681383eea2eae6dcf5b9dcd9a9b
1 file changed
+47
-4
exchange/bitswap/notifications/notifications.go
+47
-4
@@ -29,10 +29,7 @@ func (ps *impl) Publish(block *blocks.Block) {
29
ps.wrapped.Pub(block, topic)
30
}
31
32
-// Subscribe returns a channel of blocks for the given |keys|. |blockChannel|
33
-// is closed if the |ctx| times out or is cancelled, or after sending len(keys)
34
-// blocks.
35
-func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Block {
32
+func (ps *impl) SubscribeDeprec(ctx context.Context, keys ...u.Key) <-chan *blocks.Block {
33
topics := make([]string, 0)
34
for _, key := range keys {
35
topics = append(topics, string(key))
@@ -57,3 +54,49 @@ func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Blo
54
func (ps *impl) Shutdown() {
55
ps.wrapped.Shutdown()
56
}
57
+
58
+// Subscribe returns a channel of blocks for the given |keys|. |blockChannel|
59
+// is closed if the |ctx| times out or is cancelled, or after sending len(keys)
60
+// blocks.
61
+func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Block {
62
+ topics := toStrings(keys)
63
+ blocksCh := make(chan *blocks.Block, len(keys))
64
+ valuesCh := make(chan interface{}, len(keys))
65
+ ps.wrapped.AddSub(valuesCh, topics...)
66
+
67
+ go func() {
68
+ defer func() {
69
+ ps.wrapped.Unsub(valuesCh, topics...)
70
+ close(blocksCh)
71
+ }()
72
+ for _, _ = range keys {
73
+ select {
74
+ case <-ctx.Done():
75
+ return
76
+ case val, ok := <-valuesCh:
77
+ if !ok {
78
+ return
79
+ }
80
+ block, ok := val.(*blocks.Block)
81
+ if !ok {
82
+ return
83
+ }
84
+ select {
85
+ case <-ctx.Done():
86
+ return
87
+ case blocksCh <- block: // continue
88
+ }
89
+ }
90
+ }
91
+ }()
92
+
93
+ return blocksCh
94
+}
95
+
96
+func toStrings(keys []u.Key) []string {
97
+ strs := make([]string, 0)
98
+ for _, key := range keys {
99
+ strs = append(strs, string(key))
100
+ }
101
+ return strs
102
+}