-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtransport_test.go
More file actions
188 lines (166 loc) · 5.75 KB
/
Copy pathtransport_test.go
File metadata and controls
188 lines (166 loc) · 5.75 KB
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
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
package events
import (
"context"
"errors"
"sync"
"testing"
"github.com/goforj/events/eventscore"
)
type fakeTransport struct {
mu sync.Mutex
published []eventscore.Message
subscribers map[string]eventscore.MessageHandler
readyErr error
publishErr error
subscribeErr error
lastCtx context.Context
}
// newFakeTransport initializes the in-memory subscription table used by transport delegation tests.
func newFakeTransport() *fakeTransport {
return &fakeTransport{subscribers: make(map[string]eventscore.MessageHandler)}
}
// Driver reports a stable concrete driver so bus construction can exercise transport configuration.
func (f *fakeTransport) Driver() eventscore.Driver { return eventscore.DriverNATS }
// Ready returns the configured readiness failure without contacting an external broker.
func (f *fakeTransport) Ready(context.Context) error {
return f.readyErr
}
// PublishContext records the message and synchronously invokes the topic subscriber.
func (f *fakeTransport) PublishContext(ctx context.Context, msg eventscore.Message) error {
f.lastCtx = ctx
if f.publishErr != nil {
return f.publishErr
}
f.mu.Lock()
f.published = append(f.published, msg)
handler := f.subscribers[msg.Topic]
f.mu.Unlock()
if handler != nil {
return handler(context.Background(), msg)
}
return nil
}
// SubscribeContext installs one topic handler and returns a subscription that removes it.
func (f *fakeTransport) SubscribeContext(_ context.Context, topic string, handler eventscore.MessageHandler) (eventscore.Subscription, error) {
if f.subscribeErr != nil {
return nil, f.subscribeErr
}
f.mu.Lock()
defer f.mu.Unlock()
f.subscribers[topic] = handler
return fakeTransportSubscription(func() {
f.mu.Lock()
defer f.mu.Unlock()
delete(f.subscribers, topic)
}), nil
}
type fakeTransportSubscription func()
// Close invokes the captured removal function exactly as a transport subscription would.
func (f fakeTransportSubscription) Close() error {
f()
return nil
}
// TestBusWithTransportDelegatesReady verifies readiness reaches the configured transport driver.
func TestBusWithTransportDelegatesReady(t *testing.T) {
want := errors.New("not ready")
transport := newFakeTransport()
transport.readyErr = want
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
if err := bus.Ready(); !errors.Is(err, want) {
t.Fatalf("Ready error = %v, want %v", err, want)
}
}
// TestBusWithTransportPublishesAndReceives verifies transport payloads round-trip through registered handlers.
func TestBusWithTransportPublishesAndReceives(t *testing.T) {
transport := newFakeTransport()
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
called := false
_, err = bus.Subscribe(func(userCreated) { called = true })
if err != nil {
t.Fatalf("Subscribe returned error: %v", err)
}
if err := bus.Publish(userCreated{}); err != nil {
t.Fatalf("Publish returned error: %v", err)
}
if !called {
t.Fatal("expected transport-backed delivery")
}
if len(transport.published) != 1 {
t.Fatalf("published count = %d, want 1", len(transport.published))
}
if transport.published[0].Topic != "user.created" {
t.Fatalf("published topic = %q, want %q", transport.published[0].Topic, "user.created")
}
}
// TestTransportSubscriptionClosesWhenLastHandlerRemoved verifies shared transport subscriptions outlive individual handlers only as needed.
func TestTransportSubscriptionClosesWhenLastHandlerRemoved(t *testing.T) {
transport := newFakeTransport()
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
sub, err := bus.Subscribe(func(userCreated) {})
if err != nil {
t.Fatalf("Subscribe returned error: %v", err)
}
if err := sub.Close(); err != nil {
t.Fatalf("Close returned error: %v", err)
}
transport.mu.Lock()
_, ok := transport.subscribers["user.created"]
transport.mu.Unlock()
if ok {
t.Fatal("expected transport subscriber removal")
}
}
// TestBusWithTransportPropagatesPublishError verifies driver publish failures remain observable.
func TestBusWithTransportPropagatesPublishError(t *testing.T) {
want := errors.New("publish failed")
transport := newFakeTransport()
transport.publishErr = want
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
if err := bus.Publish(userCreated{}); !errors.Is(err, want) {
t.Fatalf("Publish error = %v, want %v", err, want)
}
}
// TestBusWithTransportUsesBackgroundContextForNilPublishContext verifies nil publish contexts are normalized before driver calls.
func TestBusWithTransportUsesBackgroundContextForNilPublishContext(t *testing.T) {
transport := newFakeTransport()
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
if err := bus.WithContext(nil).Publish(userCreated{}); err != nil {
t.Fatalf("WithContext(nil).Publish returned error: %v", err)
}
if transport.lastCtx == nil {
t.Fatal("expected non-nil context to be forwarded")
}
}
// TestBusWithTransportRollsBackHandlerOnSubscribeError verifies failed transport setup leaves no registered handler behind.
func TestBusWithTransportRollsBackHandlerOnSubscribeError(t *testing.T) {
want := errors.New("subscribe failed")
transport := newFakeTransport()
transport.subscribeErr = want
bus, err := New(Config{Transport: transport})
if err != nil {
t.Fatalf("New returned error: %v", err)
}
if _, err := bus.Subscribe(func(userCreated) {}); !errors.Is(err, want) {
t.Fatalf("Subscribe error = %v, want %v", err, want)
}
bus.mu.RLock()
defer bus.mu.RUnlock()
if len(bus.handlers["user.created"]) != 0 {
t.Fatal("expected failed subscribe to roll back registered handler")
}
}