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

Function TestBus

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

Source from the content-addressed store, hash-verified

10)
11
12func TestBus(t *testing.T) {
13 bus := pubsub.NewBus()
14 defer bus.Close()
15
16 did := ed25519.GenPrivKey().PubKey().Address()
17
18 ev := newEvent(did)
19
20 assert.NoError(t, bus.Publish(ev))
21
22 sub1, err := bus.Subscribe()
23 require.NoError(t, err)
24
25 sub2, err := bus.Subscribe()
26 require.NoError(t, err)
27
28 assert.NoError(t, bus.Publish(ev))
29
30 select {
31 case newEv := <-sub1.Events():
32 assert.Equal(t, ev, newEv)
33 case <-pubsub.AfterThreadStart(t):
34 require.Fail(t, "time out")
35 }
36
37 select {
38 case newEv := <-sub2.Events():
39 assert.Equal(t, ev, newEv)
40 case <-pubsub.AfterThreadStart(t):
41 require.Fail(t, "time out")
42 }
43
44 sub2.Close()
45
46 select {
47 case <-sub2.Done():
48 case <-pubsub.AfterThreadStart(t):
49 require.Fail(t, "time out")
50 }
51
52 assert.NoError(t, bus.Publish(ev))
53
54 select {
55 case newEv := <-sub1.Events():
56 assert.Equal(t, ev, newEv)
57 case <-pubsub.AfterThreadStart(t):
58 require.Fail(t, "time out")
59 }
60
61 select {
62 case <-sub2.Events():
63 require.Fail(t, "spurious event")
64 case <-pubsub.AfterThreadStart(t):
65 }
66
67 bus.Close()
68
69 select {

Callers

nothing calls this directly

Calls 9

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

Tested by

no test coverage detected