-
Notifications
You must be signed in to change notification settings - Fork 14
/
sync_adapter_sink_test.go
90 lines (71 loc) · 2.15 KB
/
sync_adapter_sink_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
package substrate
import (
"context"
"errors"
"testing"
"time"
"github.com/stretchr/testify/assert"
)
func TestSyncProduceAdapterBasic(t *testing.T) {
assert := assert.New(t)
ap := &mockAsyncSink{5, make(chan struct{}, 1)}
sc := NewSynchronousMessageSink(ap)
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
m1, m2, m3, m4, m5 := message([]byte{'a'}), message([]byte{'b'}), message([]byte{'c'}), message([]byte{'d'}), message([]byte{'e'})
assert.NoError(sc.PublishMessage(ctx, &m1))
assert.NoError(sc.PublishMessage(ctx, &m2))
assert.NoError(sc.PublishMessage(ctx, &m3))
assert.NoError(sc.PublishMessage(ctx, &m4))
assert.NoError(sc.PublishMessage(ctx, &m5))
assert.NoError(sc.Close())
select {
case <-ap.closed:
default:
t.Error("underlying async sink didn't get closed")
}
assert.Equal(ErrSinkAlreadyClosed, sc.PublishMessage(ctx, &m5))
assert.Equal(ErrSinkAlreadyClosed, sc.Close())
assert.Equal(ErrSinkAlreadyClosed, sc.PublishMessage(ctx, &m1))
}
func TestSyncProduceAdapter_ErrorOnSend(t *testing.T) {
ap := &mockAsyncSink{0, make(chan struct{}, 1)}
sc := NewSynchronousMessageSink(ap)
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
msg := message([]byte{'t'})
assert.Equal(t, ErrSinkClosedOrFailedDuringSend, sc.PublishMessage(ctx, &msg))
assert.Equal(t, errSeenAllMessages, sc.Close())
assert.Equal(t, ErrSinkAlreadyClosed, sc.Close())
}
var errSeenAllMessages = errors.New("mock sink saw specified number of messages")
type mockAsyncSink struct {
toAckCount int
closed chan struct{}
}
func (mock *mockAsyncSink) PublishMessages(ctx context.Context, acks chan<- Message, messages <-chan Message) error {
for {
select {
case <-ctx.Done():
return ctx.Err()
case m := <-messages:
if mock.toAckCount > 0 {
select {
case <-ctx.Done():
return ctx.Err()
case acks <- m:
}
mock.toAckCount--
} else {
return errSeenAllMessages
}
}
}
}
func (mock *mockAsyncSink) Close() error {
close(mock.closed)
return nil
}
func (mock *mockAsyncSink) Status() (*Status, error) {
return &Status{Working: true}, nil
}