forked from livekit/psrpc
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathbus_test.go
114 lines (94 loc) · 2.58 KB
/
bus_test.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
package psrpc
import (
"context"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/redis/go-redis/v9"
"github.com/stretchr/testify/require"
"github.com/livekit/psrpc/internal"
)
func TestMessageBus(t *testing.T) {
t.Run("Local", func(t *testing.T) {
bus := NewLocalMessageBus()
testSubscribe(t, bus)
testSubscribeQueue(t, bus)
testSubscribeClose(t, bus)
})
t.Run("Redis", func(t *testing.T) {
rc := redis.NewUniversalClient(&redis.UniversalOptions{Addrs: []string{"localhost:6379"}})
bus := NewRedisMessageBus(rc)
testSubscribe(t, bus)
testSubscribeQueue(t, bus)
testSubscribeClose(t, bus)
})
t.Run("Nats", func(t *testing.T) {
nc, _ := nats.Connect(nats.DefaultURL)
bus := NewNatsMessageBus(nc)
testSubscribe(t, bus)
testSubscribeQueue(t, bus)
testSubscribeClose(t, bus)
})
}
func testSubscribe(t *testing.T, bus MessageBus) {
ctx := context.Background()
channel := newID()
subA, err := Subscribe[*internal.Request](ctx, bus, channel, DefaultChannelSize)
require.NoError(t, err)
subB, err := Subscribe[*internal.Request](ctx, bus, channel, DefaultChannelSize)
require.NoError(t, err)
time.Sleep(time.Millisecond * 100)
require.NoError(t, bus.Publish(ctx, channel, &internal.Request{
RequestId: "1",
}))
msgA := <-subA.Channel()
msgB := <-subB.Channel()
require.NotNil(t, msgA)
require.NotNil(t, msgB)
require.Equal(t, "1", msgA.RequestId)
require.Equal(t, "1", msgB.RequestId)
}
func testSubscribeQueue(t *testing.T, bus MessageBus) {
ctx := context.Background()
channel := newID()
subA, err := SubscribeQueue[*internal.Request](ctx, bus, channel, DefaultChannelSize)
require.NoError(t, err)
subB, err := SubscribeQueue[*internal.Request](ctx, bus, channel, DefaultChannelSize)
require.NoError(t, err)
time.Sleep(time.Millisecond * 100)
require.NoError(t, bus.Publish(ctx, channel, &internal.Request{
RequestId: "2",
}))
received := 0
select {
case m := <-subA.Channel():
if m != nil {
received++
}
case <-time.After(DefaultClientTimeout):
// continue
}
select {
case m := <-subB.Channel():
if m != nil {
received++
}
case <-time.After(DefaultClientTimeout):
// continue
}
require.Equal(t, 1, received)
}
func testSubscribeClose(t *testing.T, bus MessageBus) {
ctx := context.Background()
channel := newID()
sub, err := Subscribe[*internal.Request](ctx, bus, channel, DefaultChannelSize)
require.NoError(t, err)
sub.Close()
time.Sleep(time.Millisecond * 100)
select {
case _, ok := <-sub.Channel():
require.False(t, ok)
default:
require.FailNow(t, "closed subscription channel should not block")
}
}