@cryptotaxi247 / kubo / commits / 8f8230823

feat(bitswap/notifications) Subscribe to multiple keys

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Nov 19, 2014 at 14:19 UTC 8f8230823f898184a55e488edb0e976ec1d33f5f
1 file changed +9 -6
exchange/bitswap/notifications/notifications.go
+9 -6
@@ -12,7 +12,7 @@ const bufferSize = 16
12
13 type PubSub interface {
14 Publish(block blocks.Block)
15 - Subscribe(ctx context.Context, k u.Key) <-chan blocks.Block
15 + Subscribe(ctx context.Context, keys ...u.Key) <-chan blocks.Block
16 Shutdown()
17 }
18
@@ -31,10 +31,13 @@ func (ps *impl) Publish(block blocks.Block) {
31
32 // Subscribe returns a one-time use |blockChannel|. |blockChannel| returns nil
33 // if the |ctx| times out or is cancelled. Then channel is closed after the
34 -// block given by |k| is sent.
35 -func (ps *impl) Subscribe(ctx context.Context, k u.Key) <-chan blocks.Block {
36 - topic := string(k)
37 - subChan := ps.wrapped.SubOnce(topic)
34 +// blocks given by |keys| are sent.
35 +func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan blocks.Block {
36 + topics := make([]string, 0)
37 + for _, key := range keys {
38 + topics = append(topics, string(key))
39 + }
40 + subChan := ps.wrapped.SubOnce(topics...)
41 blockChannel := make(chan blocks.Block, 1) // buffered so the sender doesn't wait on receiver
42 go func() {
43 defer close(blockChannel)
@@ -45,7 +48,7 @@ func (ps *impl) Subscribe(ctx context.Context, k u.Key) <-chan blocks.Block {
48 blockChannel <- block
49 }
50 case <-ctx.Done():
48 - ps.wrapped.Unsub(subChan, topic)
51 + ps.wrapped.Unsub(subChan, topics...)
52 }
53 }()
54 return blockChannel