MCPcopy Create free account
hub / github.com/akash-network/node / TestClone

Function TestClone

pubsub/bus_test.go:79–143  ·  view source on GitHub ↗
(t *testing.T)

Source from the content-addressed store, hash-verified

77}
78
79func 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

Callers

nothing calls this directly

Calls 11

CloseMethod · 0.95
SubscribeMethod · 0.95
NewBusFunction · 0.92
AfterThreadStartFunction · 0.92
SleepForThreadStartFunction · 0.92
newEventFunction · 0.85
AddressMethod · 0.80
PublishMethod · 0.65
EventsMethod · 0.65
CloneMethod · 0.65
DoneMethod · 0.65

Tested by

no test coverage detected