This repository was archived by the owner on Mar 9, 2026. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 215
Expand file tree
/
Copy pathtest_ordered_sequencer.py
More file actions
329 lines (224 loc) · 9.59 KB
/
Copy pathtest_ordered_sequencer.py
File metadata and controls
329 lines (224 loc) · 9.59 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
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
# Copyright 2019, Google LLC All rights reserved.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import concurrent.futures as futures
from unittest import mock
import pytest
from google.auth import credentials
from google.cloud.pubsub_v1 import publisher
from google.cloud.pubsub_v1.publisher._sequencer import ordered_sequencer
from google.pubsub_v1 import types as gapic_types
_ORDERING_KEY = "ordering_key_1"
def create_message():
return gapic_types.PubsubMessage(data=b"foo", attributes={"bar": "baz"})
def create_client():
creds = mock.Mock(spec=credentials.Credentials)
return publisher.Client(credentials=creds)
def create_ordered_sequencer(client):
return ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
def test_stop_makes_sequencer_invalid():
client = create_client()
message = create_message()
sequencer = create_ordered_sequencer(client)
sequencer.stop()
# Publish after stop() throws
with pytest.raises(RuntimeError):
sequencer.publish(message)
# Commit after stop() throws
with pytest.raises(RuntimeError):
sequencer.commit()
# Stop after stop() throws
with pytest.raises(RuntimeError):
sequencer.stop()
def test_stop_no_batches():
client = create_client()
sequencer = create_ordered_sequencer(client)
# No exceptions thrown if there are no batches.
sequencer.stop()
def test_stop_one_batch():
client = create_client()
sequencer = create_ordered_sequencer(client)
batch1 = mock.Mock(spec=client._batch_class)
sequencer._set_batches([batch1])
sequencer.stop()
# Assert that the first batch is committed.
assert batch1.commit.call_count == 1
assert batch1.cancel.call_count == 0
def test_stop_many_batches():
client = create_client()
sequencer = create_ordered_sequencer(client)
batch1 = mock.Mock(spec=client._batch_class)
batch2 = mock.Mock(spec=client._batch_class)
sequencer._set_batches([batch1, batch2])
sequencer.stop()
# Assert that the first batch is committed and the rest cancelled.
assert batch1.commit.call_count == 1
assert batch1.cancel.call_count == 0
assert batch2.commit.call_count == 0
assert batch2.cancel.call_count == 1
def test_commit():
client = create_client()
sequencer = create_ordered_sequencer(client)
batch1 = mock.Mock(spec=client._batch_class)
batch2 = mock.Mock(spec=client._batch_class)
sequencer._set_batches([batch1, batch2])
sequencer.commit()
# Only commit the first batch.
assert batch1.commit.call_count == 1
assert batch2.commit.call_count == 0
def test_commit_empty_batch_list():
client = create_client()
sequencer = create_ordered_sequencer(client)
# Test nothing bad happens.
sequencer.commit()
def test_no_commit_when_paused():
client = create_client()
batch = mock.Mock(spec=client._batch_class)
sequencer = create_ordered_sequencer(client)
sequencer._set_batch(batch)
sequencer._pause()
sequencer.commit()
assert batch.commit.call_count == 0
def test_pause_and_unpause():
client = create_client()
message = create_message()
sequencer = create_ordered_sequencer(client)
# Unpausing without pausing throws.
with pytest.raises(RuntimeError):
sequencer.unpause()
sequencer._pause()
# Publishing while paused returns a future with an exception.
future = sequencer.publish(message)
assert future.exception().ordering_key == _ORDERING_KEY
sequencer.unpause()
# Assert publish does not set exception after unpause().
future = sequencer.publish(message)
with pytest.raises(futures._base.TimeoutError):
future.exception(timeout=0)
def test_basic_publish():
client = create_client()
message = create_message()
batch = mock.Mock(spec=client._batch_class)
sequencer = create_ordered_sequencer(client)
sequencer._set_batch(batch)
sequencer.publish(message)
batch.publish.assert_called_once_with(message)
def test_publish_custom_retry():
client = create_client()
message = create_message()
sequencer = create_ordered_sequencer(client)
sequencer.publish(message, retry=mock.sentinel.custom_retry)
assert sequencer._ordered_batches # batch exists
batch = sequencer._ordered_batches[0]
assert batch._commit_retry is mock.sentinel.custom_retry
def test_publish_custom_timeout():
client = create_client()
message = create_message()
sequencer = create_ordered_sequencer(client)
sequencer.publish(message, timeout=mock.sentinel.custom_timeout)
assert sequencer._ordered_batches # batch exists
batch = sequencer._ordered_batches[0]
assert batch._commit_timeout is mock.sentinel.custom_timeout
def test_publish_batch_full():
client = create_client()
message = create_message()
batch = mock.Mock(spec=client._batch_class)
# Make batch full.
batch.publish.return_value = None
sequencer = create_ordered_sequencer(client)
sequencer._set_batch(batch)
# Will create a new batch since the old one is full, and return a future.
future = sequencer.publish(message)
batch.publish.assert_called_once_with(message)
assert future is not None
# There's now the old and the new batches.
assert len(sequencer._get_batches()) == 2
def test_batch_done_successfully():
client = create_client()
batch = mock.Mock(spec=client._batch_class)
sequencer = ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
sequencer._set_batch(batch)
sequencer._batch_done_callback(success=True)
# One batch is done, so the OrderedSequencer has no more work, and should
# return true for is_finished().
assert sequencer.is_finished()
# No batches remain in the batches list.
assert len(sequencer._get_batches()) == 0
def test_batch_done_successfully_one_batch_remains():
client = create_client()
batch1 = mock.Mock(spec=client._batch_class)
batch2 = mock.Mock(spec=client._batch_class)
sequencer = ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
sequencer._set_batches([batch1, batch2])
sequencer._batch_done_callback(success=True)
# One batch is done, but the OrderedSequencer has more work, so is_finished()
# should return false.
assert not sequencer.is_finished()
# Second batch should be not be committed since the it may still be able to
# accept messages.
assert batch2.commit.call_count == 0
# Only the second batch remains in the batches list.
assert len(sequencer._get_batches()) == 1
def test_batch_done_successfully_many_batches_remain():
client = create_client()
batch1 = mock.Mock(spec=client._batch_class)
batch2 = mock.Mock(spec=client._batch_class)
batch3 = mock.Mock(spec=client._batch_class)
sequencer = ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
sequencer._set_batches([batch1, batch2, batch3])
sequencer._batch_done_callback(success=True)
# One batch is done, but the OrderedSequencer has more work, so DO NOT
# return true for is_finished().
assert not sequencer.is_finished()
# Second batch should be committed since it is full. We know it's full
# because there exists a third batch. Batches are created only if the
# previous one can't accept messages any more / is full.
assert batch2.commit.call_count == 1
# Both the second and third batches remain in the batches list.
assert len(sequencer._get_batches()) == 2
def test_batch_done_unsuccessfully():
client = create_client()
message = create_message()
batch1 = mock.Mock(spec=client._batch_class)
batch2 = mock.Mock(spec=client._batch_class)
batch3 = mock.Mock(spec=client._batch_class)
sequencer = ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
sequencer._set_batches([batch1, batch2, batch3])
# Make the batch fail.
sequencer._batch_done_callback(success=False)
# Sequencer should remain as a sentinel to indicate this ordering key is
# paused. Therefore, don't call the cleanup callback.
assert not sequencer.is_finished()
# Cancel the remaining batches.
assert batch2.cancel.call_count == 1
assert batch3.cancel.call_count == 1
# Remove all the batches.
assert len(sequencer._get_batches()) == 0
# Verify that the sequencer is paused. Publishing while paused returns a
# future with an exception.
future = sequencer.publish(message)
assert future.exception().ordering_key == _ORDERING_KEY
def test_publish_after_finish():
client = create_client()
batch = mock.Mock(spec=client._batch_class)
sequencer = ordered_sequencer.OrderedSequencer(client, "topic_name", _ORDERING_KEY)
sequencer._set_batch(batch)
sequencer._batch_done_callback(success=True)
# One batch is done, so the OrderedSequencer has no more work, and should
# return true for is_finished().
assert sequencer.is_finished()
message = create_message()
# It's legal to publish after being finished.
sequencer.publish(message)
# Go back to accepting-messages mode.
assert not sequencer.is_finished()