-
Notifications
You must be signed in to change notification settings - Fork 35
Expand file tree
/
Copy pathredis_test.rb
More file actions
566 lines (468 loc) · 15.4 KB
/
Copy pathredis_test.rb
File metadata and controls
566 lines (468 loc) · 15.4 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
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
# frozen_string_literal: true
require 'test_helper'
class CI::Queue::RedisTest < Minitest::Test
include SharedQueueAssertions
EntryTest = Struct.new(:id, :queue_entry)
def setup
@redis_url = ENV.fetch('REDIS_URL', 'redis://localhost:6379/0')
@redis = ::Redis.new(url: @redis_url)
@redis.flushdb
super
@config = @queue.send(:config) # hack
end
def test_from_uri
second_queue = populate(
CI::Queue.from_uri(@redis_url, config)
)
assert_instance_of CI::Queue::Redis::Worker, second_queue
assert_equal @queue.to_a, second_queue.to_a
end
def test_requeue # redefine the shared one
previous_offset = CI::Queue::Redis.requeue_offset
CI::Queue::Redis.requeue_offset = 2
failed_once = false
test_order = poll(@queue, ->(test) {
if test == shuffled_test_list.last && !failed_once
failed_once = true
false
else
true
end
})
expected_order = shuffled_test_list.dup
expected_order.insert(-CI::Queue::Redis.requeue_offset, shuffled_test_list.last)
assert_equal expected_order, test_order
ensure
CI::Queue::Redis.requeue_offset = previous_offset
end
def test_retry_queue_with_all_tests_passing
poll(@queue)
retry_queue = @queue.retry_queue
populate(retry_queue)
retry_test_order = poll(retry_queue)
assert_equal [], retry_test_order
end
def test_retry_queue_with_all_tests_passing_2
poll(@queue)
retry_queue = @queue.retry_queue
populate(retry_queue)
retry_test_order = poll(retry_queue) do |test|
@queue.build.record_error(test.id, 'Failed')
end
assert_equal retry_test_order, retry_test_order
end
def test_shutdown
poll(@queue) do
@queue.shutdown!
end
assert_equal TEST_LIST.size - 1, @queue.size
end
def test_master_election
assert_predicate @queue, :master?
refute_predicate worker(2), :master?
@redis.flushdb
assert_predicate worker(2), :master?
refute_predicate worker(1), :master?
end
def test_exhausted_while_not_populated
assert_predicate @queue, :populated?
second_worker = worker(2, populate: false)
refute_predicate second_worker, :populated?
refute_predicate second_worker, :exhausted?
poll(@queue)
refute_predicate second_worker, :populated?
assert_predicate second_worker, :exhausted?
end
def test_monitor_boot_and_shutdown
@queue.config.max_missed_heartbeat_seconds = 1
@queue.boot_heartbeat_process!
status = @queue.stop_heartbeat!
assert_predicate status, :success?
ensure
@queue.config.max_missed_heartbeat_seconds = nil
end
def test_timed_out_test_are_picked_up_by_other_workers
second_queue = worker(2)
acquired = false
done = false
monitor = Monitor.new
condition = monitor.new_cond
Thread.start do
monitor.synchronize do
condition.wait_until { acquired }
poll(second_queue)
done = true
condition.signal
end
end
poll(@queue) do
acquired = true
monitor.synchronize do
condition.signal
condition.wait_until { done }
end
end
assert_predicate @queue, :exhausted?
assert_equal [], populate(@queue.retry_queue).to_a
assert_equal [], populate(second_queue.retry_queue).to_a.sort
end
def test_release_immediately_timeout_the_lease
second_queue = worker(2)
reserved_test = nil
poll(@queue) do |test|
reserved_test = test
break
end
refute_nil reserved_test
worker(1).release! # Use a new instance to ensure we don't depend on in-memory state
poll(second_queue) do |test|
assert_equal reserved_test, test
break
end
end
def test_test_isnt_requeued_if_it_was_picked_up_by_another_worker
second_queue = worker(2)
acquired = false
done = false
monitor = Monitor.new
condition = monitor.new_cond
Thread.start do
monitor.synchronize do
condition.wait_until { acquired }
poll(second_queue)
done = true
condition.signal
end
end
poll(@queue, false) do
break if acquired
acquired = true
monitor.synchronize do
condition.signal
condition.wait_until { done }
end
end
assert_predicate @queue, :exhausted?
end
def test_acknowledge_returns_false_if_the_test_was_picked_up_by_another_worker
second_queue = worker(2)
acquired = false
done = false
monitor = Monitor.new
condition = monitor.new_cond
Thread.start do
monitor.synchronize do
condition.wait_until { acquired }
second_queue.poll do |test|
assert_equal true, second_queue.acknowledge(test.id)
end
done = true
condition.signal
end
end
@queue.poll do |test|
break if acquired
acquired = true
monitor.synchronize do
condition.signal
condition.wait_until { done }
assert_equal false, @queue.acknowledge(test.id)
end
end
assert_predicate @queue, :exhausted?
end
def test_workers_register
assert_equal 1, @redis.scard(('build:42:workers'))
worker(2)
assert_equal 2, @redis.scard(('build:42:workers'))
end
def test_timeout_warning
begin
threads = 2.times.map do |i|
Thread.new do
queue = worker(i, tests: [TEST_LIST.first], build_id: '24')
queue.poll do |test|
sleep 1 # timeout
queue.acknowledge(test.id)
end
end
end
threads.each { |t| t.join(3) }
threads.each { |t| refute_predicate t, :alive? }
queue = worker(12, build_id: '24')
assert_equal [[:RESERVED_LOST_TEST, {test: 'ATest#test_foo', timeout: 0.2}]], queue.build.pop_warnings
ensure
threads.each(&:kill)
end
end
def test_streaming_waits_for_batches
leader = worker(1, populate: false, lazy_load_streaming_timeout: 2, queue_init_timeout: 2, build_id: 'streaming')
consumer = worker(2, populate: false, lazy_load_streaming_timeout: 2, queue_init_timeout: 2, build_id: 'streaming')
consumer.entry_resolver = ->(entry) { entry }
tests = [
EntryTest.new('ATest#test_foo', CI::Queue::QueueEntry.format('ATest#test_foo', '/tmp/a_test.rb')),
EntryTest.new('ATest#test_bar', CI::Queue::QueueEntry.format('ATest#test_bar', '/tmp/a_test.rb')),
]
streamed = Enumerator.new do |yielder|
sleep 0.2
tests.each { |test| yielder << test }
end
leader_thread = Thread.new do
leader.stream_populate(streamed, random: Random.new(0), batch_size: 1)
end
timeout_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 1
loop do
status = @redis.get(leader.send(:key, 'master-status'))
break if status == 'streaming' || status == 'ready'
raise "streaming status not set" if Process.clock_gettime(Process::CLOCK_MONOTONIC) > timeout_at
sleep 0.01
end
consumed = []
consumer_thread = Thread.new do
consumer.poll do |entry|
consumed << entry
consumer.acknowledge(entry)
end
end
sleep 0.05
leader_thread.join
consumer_thread.join(2)
assert_equal tests.map(&:queue_entry).sort, consumed.sort
assert_predicate consumer, :exhausted?
end
def test_reserve_lost_ignores_processed_entry_with_path
queue = worker(1, populate: false)
entry = CI::Queue::QueueEntry.format('ATest#test_foo', '/tmp/a_test.rb')
test_id = 'ATest#test_foo'
@redis.zadd(queue.send(:key, 'running'), 0, entry)
@redis.sadd(queue.send(:key, 'completed'), test_id)
@redis.hset(queue.send(:key, 'owners'), entry, queue.send(:key, 'worker', queue.config.worker_id, 'queue'))
lost = queue.send(:try_to_reserve_lost_test)
assert_nil lost
end
def test_streaming_timeout_raises_lost_master
queue = worker(1, populate: false, lazy_load_streaming_timeout: 1, queue_init_timeout: 1)
@redis.set(queue.send(:key, 'master-status'), 'streaming')
@redis.set(queue.send(:key, 'streaming-updated-at'), CI::Queue.time_now.to_f - 5)
assert_raises(CI::Queue::Redis::LostMaster) do
queue.poll { |_entry| }
end
end
def test_reserve_defers_own_requeued_test_once
queue = worker(1, populate: false, build_id: 'self-requeue-script')
queue.send(:register)
entry = CI::Queue::QueueEntry.format('ATest#test_foo', '/tmp/a_test.rb')
queue_key = queue.send(:key, 'queue')
requeued_by_key = queue.send(:key, 'requeued-by')
worker_queue_key = queue.send(:key, 'worker', queue.config.worker_id, 'queue')
workers_key = queue.send(:key, 'workers')
@redis.lpush(queue_key, entry)
@redis.hset(requeued_by_key, entry, worker_queue_key)
@redis.sadd(workers_key, '2')
first_try = queue.send(:try_to_reserve_test)
assert_nil first_try
assert_equal [entry], @redis.lrange(queue_key, 0, -1)
assert_nil @redis.hget(requeued_by_key, entry)
second_try = queue.send(:try_to_reserve_test)
assert_equal entry, second_try
end
def test_heartbeat_uses_test_id_for_processed_check
queue = worker(1, populate: false)
entry = CI::Queue::QueueEntry.format('ATest#test_foo', '/tmp/a_test.rb')
test_id = 'ATest#test_foo'
@redis.sadd(queue.send(:key, 'processed'), test_id)
result = queue.send(
:eval_script,
:heartbeat,
keys: [
queue.send(:key, 'running'),
queue.send(:key, 'processed'),
queue.send(:key, 'owners'),
queue.send(:key, 'worker', queue.config.worker_id, 'queue'),
],
argv: [CI::Queue.time_now.to_f, entry],
)
assert_nil result
end
def test_resolve_entry_falls_back_to_resolver
queue = worker(1, populate: false)
queue.instance_variable_set(:@index, { 'ATest#test_foo' => :ok })
queue.entry_resolver = ->(entry) { "resolved:#{entry}" }
missing_entry = CI::Queue::QueueEntry.format('MissingTest#test_bar', '/tmp/missing.rb')
resolved = queue.send(:resolve_entry, missing_entry)
assert_equal "resolved:#{missing_entry}", resolved
end
def test_continuously_timing_out_tests
3.times do
@redis.flushdb
begin
threads = 2.times.map do |i|
Thread.new do
queue = worker(i, tests: [TEST_LIST.first], build_id: '24')
queue.poll do |test|
sleep 1 # timeout
queue.acknowledge(test.id)
end
end
end
threads.each { |t| t.join(3) }
threads.each { |t| refute_predicate t, :alive? }
queue = worker(12, build_id: '24')
assert_predicate queue, :queue_initialized?
assert_predicate queue, :exhausted?
ensure
threads.each(&:kill)
end
end
end
def test_initialise_from_redis_uri
queue = CI::Queue.from_uri('redis://localhost:6379/0', config)
assert_instance_of CI::Queue::Redis::Worker, queue
end
def test_initialise_from_rediss_uri
queue = CI::Queue.from_uri('rediss://localhost:6379/0', config)
assert_instance_of CI::Queue::Redis::Worker, queue
end
def test_first_reserve_at_is_set_on_first_reserve
queue = worker(1)
assert_nil queue.first_reserve_at
queue.poll do |_test|
assert queue.first_reserve_at, "first_reserve_at should be set after first reserve"
break
end
end
def test_first_reserve_at_does_not_change_on_subsequent_reserves
queue = worker(1)
first_value = nil
count = 0
queue.poll do |_test|
first_value ||= queue.first_reserve_at
assert_equal first_value, queue.first_reserve_at
count += 1
break if count >= 3
end
assert_operator count, :>=, 2, "Should have reserved multiple tests"
end
def test_record_and_read_worker_profiles
queue = worker(1)
profile = {
'worker_id' => '1',
'mode' => 'lazy',
'role' => 'leader',
'total_wall_clock' => 12.34,
'time_to_first_test' => 1.23,
'memory_rss_kb' => 512_000,
}
queue.build.record_worker_profile(profile)
profiles = queue.build.worker_profiles
assert_equal 1, profiles.size
assert_equal profile, profiles['1']
end
def test_worker_profiles_aggregates_multiple_workers
q1 = worker(1)
q2 = worker(2)
q1.build.record_worker_profile({ 'worker_id' => '1', 'role' => 'leader' })
q2.build.record_worker_profile({ 'worker_id' => '2', 'role' => 'non-leader' })
profiles = q1.build.worker_profiles
assert_equal 2, profiles.size
assert_equal 'leader', profiles['1']['role']
assert_equal 'non-leader', profiles['2']['role']
end
def test_worker_does_not_pick_up_its_own_requeued_test_when_others_are_available
@redis.flushdb
test_list = TEST_LIST.first(3)
w1 = worker(1, tests: test_list, build_id: 'self-requeue', timeout: 10, max_requeues: 1, requeue_tolerance: 1.0)
w2 = worker(2, populate: false, build_id: 'self-requeue', timeout: 10, max_requeues: 1, requeue_tolerance: 1.0)
w3 = worker(3, populate: false, build_id: 'self-requeue', timeout: 10, max_requeues: 1, requeue_tolerance: 1.0)
w2.send(:register)
w3.send(:register)
id_for = ->(test) { test.respond_to?(:id) ? test.id : CI::Queue::QueueEntry.test_id(test) }
requeued_test_id = nil
picked_up_requeue = {}
worker_two_reserved = false
worker_three_reserved = false
release_other_workers = false
mon = Monitor.new
cond = mon.new_cond
threads = [
Thread.new do
w2.poll do |test|
test_id = id_for.call(test)
mon.synchronize do
worker_two_reserved = true
picked_up_requeue['2'] = true if test_id == requeued_test_id
cond.broadcast
cond.wait_until { release_other_workers }
end
w2.acknowledge(test_id)
end
end,
Thread.new do
w3.poll do |test|
test_id = id_for.call(test)
mon.synchronize do
worker_three_reserved = true
picked_up_requeue['3'] = true if test_id == requeued_test_id
cond.broadcast
cond.wait_until { release_other_workers }
end
w3.acknowledge(test_id)
end
end,
]
mon.synchronize do
cond.wait_until { worker_two_reserved && worker_three_reserved }
end
worker_one_picked_its_own_requeue = false
first_test = true
w1.poll do |test|
test_id = id_for.call(test)
if first_test
first_test = false
requeued_test_id = test_id
w1.report_failure!
assert_equal true, w1.requeue(test)
mon.synchronize do
release_other_workers = true
cond.broadcast
end
else
worker_one_picked_its_own_requeue = true if test_id == requeued_test_id
w1.acknowledge(test_id)
end
end
threads.each { |t| t.join(5) }
assert_equal false, worker_one_picked_its_own_requeue
assert_equal true, picked_up_requeue.values.any?
ensure
threads&.each(&:kill)
end
private
def shuffled_test_list
CI::Queue.shuffle(TEST_LIST, Random.new(0)).freeze
end
def build_queue
worker(1, max_requeues: 1, requeue_tolerance: 0.1, populate: false, max_consecutive_failures: 10)
end
def populate(worker, tests: TEST_LIST.dup)
worker.populate(tests, random: Random.new(0))
end
def worker(id, **args)
tests = args.delete(:tests) || TEST_LIST.dup
skip_populate = args.delete(:populate) == false
queue = CI::Queue::Redis.new(
@redis_url,
CI::Queue::Configuration.new(
build_id: '42',
worker_id: id.to_s,
timeout: 0.2,
**args,
)
)
if skip_populate
return queue
else
populate(queue, tests: tests)
end
end
end