diff --git a/redis/requeue.lua b/redis/requeue.lua index 00d54e40..a0f09342 100644 --- a/redis/requeue.lua +++ b/redis/requeue.lua @@ -11,10 +11,30 @@ local leases_key = KEYS[9] local max_requeues = tonumber(ARGV[1]) local global_max_requeues = tonumber(ARGV[2]) local entry = ARGV[3] -local offset = ARGV[4] +local offset = tonumber(ARGV[4]) or 0 local ttl = tonumber(ARGV[5]) local lease_id = ARGV[6] +-- Inserts `entry` behind the `offset` + 1 entries nearest the tail, the end RPOP reserves +-- from: where `LINSERT BEFORE ` would put it. LINSERT finds its +-- pivot by scanning from the head, which is O(queue length) and blocks Redis on large queues; +-- this only touches the tail, so it is O(offset). Queues of at most `offset` + 1 entries, and +-- offsets of zero or less, push to the head instead. +-- Keep in sync with reserve.lua: the Python client does not resolve `-- @include`. +local function insert_with_offset(queue_key, entry, offset) + if offset <= 0 or redis.call('llen', queue_key) <= offset + 1 then + redis.call('lpush', queue_key, entry) + return + end + + local ahead = redis.call('lrange', queue_key, -1 - offset, -1) + redis.call('ltrim', queue_key, 0, -2 - offset) + redis.call('rpush', queue_key, entry) + for _, ahead_entry in ipairs(ahead) do + redis.call('rpush', queue_key, ahead_entry) + end +end + -- Only the current lease holder can requeue a test. -- If the lease was transferred (e.g. via reserve_lost), reject the stale -- worker's requeue so the running entry stays intact for the new holder. @@ -41,12 +61,7 @@ redis.call('hincrby', requeues_count_key, entry, 1) redis.call('hdel', error_reports_key, entry) -local pivot = redis.call('lrange', queue_key, -1 - offset, 0 - offset)[1] -if pivot then - redis.call('linsert', queue_key, 'BEFORE', pivot, entry) -else - redis.call('lpush', queue_key, entry) -end +insert_with_offset(queue_key, entry, offset) redis.call('hset', requeued_by_key, entry, worker_queue_key) if ttl and ttl > 0 then diff --git a/redis/reserve.lua b/redis/reserve.lua index 928dd4d0..71f8ba8e 100644 --- a/redis/reserve.lua +++ b/redis/reserve.lua @@ -12,12 +12,23 @@ local current_time = ARGV[1] local defer_offset = tonumber(ARGV[2]) or 0 local max_skip_attempts = 4 -local function insert_with_offset(test) - local pivot = redis.call('lrange', queue_key, -1 - defer_offset, 0 - defer_offset)[1] - if pivot then - redis.call('linsert', queue_key, 'BEFORE', pivot, test) - else - redis.call('lpush', queue_key, test) +-- Inserts `entry` behind the `offset` + 1 entries nearest the tail, the end RPOP reserves +-- from: where `LINSERT BEFORE ` would put it. LINSERT finds its +-- pivot by scanning from the head, which is O(queue length) and blocks Redis on large queues; +-- this only touches the tail, so it is O(offset). Queues of at most `offset` + 1 entries, and +-- offsets of zero or less, push to the head instead. +-- Keep in sync with requeue.lua: the Python client does not resolve `-- @include`. +local function insert_with_offset(queue_key, entry, offset) + if offset <= 0 or redis.call('llen', queue_key) <= offset + 1 then + redis.call('lpush', queue_key, entry) + return + end + + local ahead = redis.call('lrange', queue_key, -1 - offset, -1) + redis.call('ltrim', queue_key, 0, -2 - offset) + redis.call('rpush', queue_key, entry) + for _, ahead_entry in ipairs(ahead) do + redis.call('rpush', queue_key, ahead_entry) end end @@ -44,7 +55,7 @@ for attempt = 1, max_skip_attempts do return claim_test(test) end - insert_with_offset(test) + insert_with_offset(queue_key, test, defer_offset) -- If this worker only finds its own requeued tests, defer once by returning nil, -- then allow pickup on a subsequent reserve attempt. diff --git a/ruby/test/ci/queue/redis_test.rb b/ruby/test/ci/queue/redis_test.rb index ca87b506..e970b96d 100644 --- a/ruby/test/ci/queue/redis_test.rb +++ b/ruby/test/ci/queue/redis_test.rb @@ -43,6 +43,46 @@ def test_requeue # redefine the shared one CI::Queue::Redis.requeue_offset = previous_offset end + def test_requeue_runs_the_test_again_after_the_next_offset_plus_one_tests + with_requeue_offset(2) do + tests = numbered_tests(8) + pop_order = CI::Queue.shuffle(tests, Random.new(0)) + queue = worker(1, tests: tests, build_id: 'requeue-offset', max_requeues: 1, requeue_tolerance: 1.0) + + test_order = poll_failing_once(queue, pop_order.first) + + assert_equal [pop_order[0], *pop_order[1..3], pop_order[0], *pop_order[4..]], test_order + end + end + + def test_requeue_runs_the_test_last_and_keeps_the_queue_ttl_when_only_offset_plus_one_tests_remain + with_requeue_offset(2) do + tests = numbered_tests(4) + pop_order = CI::Queue.shuffle(tests, Random.new(0)) + queue = worker(1, tests: tests, build_id: 'requeue-short-queue', max_requeues: 1, requeue_tolerance: 1.0) + queue_ttl_after_requeue = nil + + test_order = poll_failing_once(queue, pop_order.first) do |test| + queue_ttl_after_requeue = @redis.ttl('build:requeue-short-queue:queue') if test == pop_order[1] + end + + assert_equal [*pop_order, pop_order[0]], test_order + assert_operator queue_ttl_after_requeue, :>, 0 + end + end + + def test_requeue_with_a_zero_offset_runs_the_test_last + with_requeue_offset(0) do + tests = numbered_tests(4) + pop_order = CI::Queue.shuffle(tests, Random.new(0)) + queue = worker(1, tests: tests, build_id: 'requeue-zero-offset', max_requeues: 1, requeue_tolerance: 1.0) + + test_order = poll_failing_once(queue, pop_order.first) + + assert_equal [*pop_order, pop_order[0]], test_order + end + end + def test_retry_queue_with_all_tests_passing poll(@queue) retry_queue = @queue.retry_queue @@ -467,6 +507,24 @@ def test_reserve_defers_own_requeued_test_once assert_equal entry, second_try[0] end + def test_reserve_defers_own_requeued_test_behind_the_next_offset_plus_one_tests + with_requeue_offset(2) do + tests = numbered_tests(10) + pop_order = CI::Queue.shuffle(tests, Random.new(0)) + queue = worker(1, tests: tests, build_id: 'defer-offset', max_requeues: 1, requeue_tolerance: 1.0) + worker(2, tests: tests, build_id: 'defer-offset') + assert_equal 2, queue.workers_count + queue_after_deferral = nil + + test_order = poll_failing_once(queue, pop_order.first) do |test| + queue_after_deferral = queue.to_a if test == pop_order[4] + end + + assert_equal [*pop_order[4..6], pop_order[0], *pop_order[7..]], queue_after_deferral + assert_equal [*pop_order, pop_order[0]], test_order + end + end + def test_heartbeat_only_checks_lease queue = worker(1, populate: false) entry = CI::Queue::QueueEntry.format('ATest#test_foo', '/tmp/a_test.rb') @@ -1111,6 +1169,29 @@ def worker(id, **args) end end + def numbered_tests(count) + Array.new(count) { |index| SharedTestCases::TestCase.new("ATest#test_#{index}") } + end + + def poll_failing_once(queue, failing_test, &block) + failed = false + passes = ->(test) do + next true if failed || test != failing_test + + failed = true + false + end + poll(queue, passes, &block) + end + + def with_requeue_offset(offset) + previous_offset = CI::Queue::Redis.requeue_offset + CI::Queue::Redis.requeue_offset = offset + yield + ensure + CI::Queue::Redis.requeue_offset = previous_offset + end + # A failed TLS handshake is otherwise retried for ~10 seconds. def without_reconnect_attempts original = ENV['CI_QUEUE_DISABLE_RECONNECT_ATTEMPTS']