-
Notifications
You must be signed in to change notification settings - Fork 14
/
Copy pathqueue_test.go
143 lines (107 loc) · 3.56 KB
/
queue_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
package goq_test
import (
. "github.com/masslessparticle/goq"
"github.com/masslessparticle/goq/testhelpers"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
)
var _ = Describe("Queue", func() {
Context("Enqueue", func() {
var publisher *testhelpers.TestPublisher
BeforeEach(func() {
publisher = testhelpers.NewTestPublisher()
publisher.Responses <- true
})
It("can enqueue a message", func() {
queue := NewGoQ(25, publisher)
err := queue.Enqueue(Message{
Id: "1",
Payload: "The Message",
})
Expect(err).ToNot(HaveOccurred())
})
It("throws an error when the queuedepth is exceeded", func() {
queue := NewGoQ(1, publisher)
err := queue.Enqueue(Message{
Id: "1",
Payload: "The Message",
})
Expect(err).ToNot(HaveOccurred())
err = queue.Enqueue(Message{
Id: "2",
Payload: "Another Message",
})
Expect(err).To(HaveOccurred())
})
It("returns an error when a message is sent to a done queue", func() {
publisher := testhelpers.NewTestPublisher()
publisher.Responses <- true
queue := NewGoQ(25, publisher)
queue.StartPublishing()
queue.Enqueue(Message{Id: "MessageId - 1"})
recievedMessage := Message{}
Eventually(publisher.Messages).Should(Receive(&recievedMessage))
Expect(recievedMessage.Id).To(Equal("MessageId - 1"))
queue.StopPublishing()
err := queue.Enqueue(Message{Id: "MessageId - 2"})
Expect(err).To(HaveOccurred())
})
})
Context("Notifications", func() {
It("passes to the pubsub whether or not the channel is done", func() {
publisher := testhelpers.NewTestPublisher()
publisher.Responses <- true
queue := NewGoQ(25, publisher)
queue.StartPublishing()
queue.Enqueue(Message{Id: "MessageId - 1"})
recievedMessage := Message{}
Eventually(publisher.Messages).Should(Receive(&recievedMessage))
Expect(recievedMessage.Id).To(Equal("MessageId - 1"))
queue.StopPublishing()
Eventually(publisher.DoneCalls).Should(Receive(Equal(true)))
})
It("sends the message to the publisher", func() {
publisher := testhelpers.NewTestPublisher()
publisher.Responses <- true
queue := NewGoQ(25, publisher)
queue.StartPublishing()
queue.Enqueue(Message{
Id: "MessageId - 1",
Payload: "This is the message",
})
recievedMessage := Message{}
Eventually(publisher.Messages).Should(Receive(&recievedMessage))
Expect(recievedMessage.Payload).To(Equal("This is the message"))
})
It("doesn't send notifications after stopping publishing", func() {
publisher := testhelpers.NewTestPublisher()
publisher.Responses <- true
queue := NewGoQ(25, publisher)
queue.StartPublishing()
queue.Enqueue(Message{Id: "MessageId - 1"})
recievedMessage := Message{}
Eventually(publisher.Messages).Should(Receive(&recievedMessage))
Expect(recievedMessage.Id).To(Equal("MessageId - 1"))
queue.PausePublishing()
queue.Enqueue(Message{Id: "MessageId - 2"})
Consistently(publisher.Messages).ShouldNot(Receive())
})
It("retries message if delivery fails", func() {
publisher := testhelpers.NewTestPublisher()
publisher.Responses <- false
publisher.Responses <- true
queue := NewGoQ(25, publisher)
queue.StartPublishing()
queue.Enqueue(Message{Id: "MessageId - 1"})
Eventually(func() int {
return len(publisher.Messages)
}).Should(Equal(2))
})
})
It("does not panic if StopPublishing multiple times", func() {
publisher := testhelpers.NewTestPublisher()
queue := NewGoQ(25, publisher)
queue.StopPublishing()
Expect(queue.StopPublishing).ToNot(Panic())
})
})