(t *testing.T)
| 77 | } |
| 78 | |
| 79 | func TestClone(t *testing.T) { |
| 80 | bus := pubsub.NewBus() |
| 81 | defer bus.Close() |
| 82 | |
| 83 | did1 := ed25519.GenPrivKey().PubKey().Address() |
| 84 | ev1 := newEvent(did1) |
| 85 | |
| 86 | did2 := ed25519.GenPrivKey().PubKey().Address() |
| 87 | ev2 := newEvent(did2) |
| 88 | |
| 89 | assert.NoError(t, bus.Publish(ev1)) |
| 90 | |
| 91 | sub1, err := bus.Subscribe() |
| 92 | require.NoError(t, err) |
| 93 | |
| 94 | select { |
| 95 | case <-sub1.Events(): |
| 96 | require.Fail(t, "spurious event") |
| 97 | case <-pubsub.AfterThreadStart(t): |
| 98 | } |
| 99 | |
| 100 | assert.NoError(t, bus.Publish(ev1)) |
| 101 | assert.NoError(t, bus.Publish(ev2)) |
| 102 | |
| 103 | // allow event propagation |
| 104 | pubsub.SleepForThreadStart(t) |
| 105 | |
| 106 | // clone subscription |
| 107 | sub2, err := sub1.Clone() |
| 108 | require.NoError(t, err) |
| 109 | |
| 110 | // both subscriptions should receive both events |
| 111 | |
| 112 | for i, pev := range []pubsub.Event{ev1, ev2} { |
| 113 | select { |
| 114 | case ev := <-sub1.Events(): |
| 115 | assert.Equal(t, pev, ev, "sub1 event %v", i+1) |
| 116 | case <-pubsub.AfterThreadStart(t): |
| 117 | require.Fail(t, "timeout sub1 event %v", i+1) |
| 118 | } |
| 119 | |
| 120 | select { |
| 121 | case ev := <-sub2.Events(): |
| 122 | assert.Equal(t, pev, ev, "sub2 event %v", i+1) |
| 123 | case <-pubsub.AfterThreadStart(t): |
| 124 | require.Fail(t, "timeout sub2 event %v", i+1) |
| 125 | } |
| 126 | } |
| 127 | |
| 128 | // sub1 should close sub2 |
| 129 | sub1.Close() |
| 130 | |
| 131 | select { |
| 132 | case <-sub2.Done(): |
| 133 | case <-pubsub.AfterThreadStart(t): |
| 134 | require.Fail(t, "time out closing sub2") |
| 135 | } |
| 136 |
nothing calls this directly
no test coverage detected