forked from acjzz/gokaf
-
Notifications
You must be signed in to change notification settings - Fork 0
/
engine_test.go
153 lines (141 loc) · 3.55 KB
/
engine_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
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
package gokaf
import (
"fmt"
"github.com/sirupsen/logrus"
"strings"
"sync"
"testing"
)
func TestEngine_Publish_Error(t *testing.T) {
topicName := "test"
type fields struct {
name string
logLevel logrus.Level
}
type args struct {
name string
obj interface{}
}
tests := []struct {
name string
fields fields
args args
wantErr bool
wantErrMsg string
startTopic bool
}{
{
"NoTopic",
fields{"testEngine", logrus.ErrorLevel},
args{topicName, "message"},
true,
fmt.Sprintf("topic '%s' does not exist", topicName),
false,
}, {
"TopicStopped",
fields{"testEngine", logrus.ErrorLevel},
args{topicName, "message"},
true,
"topic closed",
true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
ge := NewEngine(
tt.fields.name,
tt.fields.logLevel,
)
if tt.startTopic {
ge.AddTopic(topicName, func(topic string, obj interface{}) {})
ge.Stop()
}
err := ge.Publish(tt.args.name, tt.args.obj)
if !tt.wantErr && err != nil {
t.Errorf("publish() not error, wantErr %v", tt.wantErr)
} else if err != nil {
ErrMsg := fmt.Sprintf("%v", err)
if strings.Compare(ErrMsg, tt.wantErrMsg) != 0 {
t.Errorf("publish() error = '%v', wantErr '%v'", ErrMsg, tt.wantErrMsg)
}
}
ge.Stop()
})
}
}
func TestEngine_Publish(t *testing.T) {
testName := "publish & Consume"
t.Run(testName, func(t *testing.T) {
ge := NewEngine(testName, logrus.ErrorLevel)
topicName := "topic"
msg := "Test Message"
ge.AddTopic(topicName, func(topic string, obj interface{}) {
if strings.Compare(fmt.Sprintf("%v", obj), msg) != 0 {
t.Errorf("publish() received = '%v', expected '%v'", obj, msg)
} else if strings.Compare(topicName, topic) != 0 {
t.Errorf("publish() received from topic '%s', expected '%s'", topic, topicName)
}
})
err := ge.Publish(topicName, msg)
if err != nil {
t.Errorf("publish() error, wanted not Error\nError: %v", err)
}
ge.Stop()
})
}
func TestEngine_Publish_Multiple_Topics(t *testing.T) {
baseTopicName := "topic"
baseMsg := "Test Message"
type Message struct {
topic string
message string
}
tests := []struct {
name string
numTopics int
numConsumers int
}{
{"Two topics", 2, 1},
{"Three topics twenty Consumers", 3, 20},
{"Five topics", 5, 1},
{"Five topics two Consumers", 5, 2},
{"Twenty topics five Consumers", 20, 5},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
ge := NewEngine(tt.name, logrus.ErrorLevel)
for i := 0; i < tt.numTopics; i++ {
topicName := fmt.Sprintf("%s-%d", baseTopicName, i)
ge.AddTopic(topicName, func(topic string, obj interface{}) {
received := obj.(Message)
if strings.Compare(received.topic, topic) != 0 {
t.Errorf("publish() received from topic '%s', expected '%s'", received.topic, topicName)
}
}, tt.numConsumers)
}
var wg sync.WaitGroup
for i := 0; i < tt.numTopics-1; i++ {
wg.Add(1)
topicName := fmt.Sprintf("%s-%d", baseTopicName, i)
go func(topic string) {
msg := Message{topicName, baseMsg}
err := ge.Publish(topicName, msg)
if err != nil {
t.Errorf("publish() error, wanted not Error\nError: %v", err)
}
wg.Done()
}(topicName)
}
wg.Add(1)
topicName := fmt.Sprintf("%s-%d", baseTopicName, tt.numTopics-1)
msg := Message{topicName, baseMsg}
err := ge.Publish(topicName, msg)
if err != nil {
t.Errorf("publish() error, wanted not Error\nError: %v", err)
}
wg.Done()
wg.Wait()
ge.Stop()
})
}
}