fix(bs/notifications) prevent duplicates
@whyrusleeping now notifications _guarantees_ there won't be any duplicates License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Nov 21, 2014 at 19:11 UTC
9bf1ba6ab5db0ef92194a3b9e1d6b529865fa41a
2 files changed
+20
-3
exchange/bitswap/notifications/notifications.go
+18
-1
@@ -47,7 +47,12 @@ func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Blo
47
ps.wrapped.Unsub(valuesCh, toStrings(keys)...)
48
close(blocksCh)
49
}()
50
- for _, _ = range keys {
50
+ seen := make(map[u.Key]struct{})
51
+ i := 0 // req'd because it only counts unique block sends
52
+ for {
53
+ if i >= len(keys) {
54
+ return
55
+ }
56
select {
57
case <-ctx.Done():
58
return
@@ -59,10 +64,22 @@ func (ps *impl) Subscribe(ctx context.Context, keys ...u.Key) <-chan *blocks.Blo
64
if !ok {
65
return
66
}
67
+ if _, ok := seen[block.Key()]; ok {
68
+ continue
69
+ }
70
select {
71
case <-ctx.Done():
72
return
73
case blocksCh <- block: // continue
74
+ // Unsub alone is insufficient for keeping out duplicates.
75
+ // It's a race to unsubscribe before pubsub handles the
76
+ // next Publish call. Therefore, must also check for
77
+ // duplicates manually. Unsub is a performance
78
+ // consideration to avoid lots of unnecessary channel
79
+ // chatter.
80
+ ps.wrapped.Unsub(valuesCh, string(block.Key()))
81
+ i++
82
+ seen[block.Key()] = struct{}{}
83
}
84
}
85
}
exchange/bitswap/notifications/notifications_test.go
+2
-2
@@ -52,8 +52,8 @@ func TestPublishSubscribe(t *testing.T) {
52
}
53
54
func TestSubscribeMany(t *testing.T) {
55
- e1 := blocks.NewBlock([]byte("Greetings from The Interval"))
56
- e2 := blocks.NewBlock([]byte("Greetings from The Interval"))
55
+ e1 := blocks.NewBlock([]byte("1"))
56
+ e2 := blocks.NewBlock([]byte("2"))
57
58
n := New()
59
defer n.Shutdown()