From 471e2abe8854d76524d12a2169982f4aeebde80e Mon Sep 17 00:00:00 2001 From: Ryan Hoffman Date: Thu, 1 Oct 2026 19:51:45 -0400 Subject: [PATCH] Requeue near the tail in O(offset) instead of scanning the queue with LINSERT requeue.lua and reserve.lua put a requeued test back `offset` + 1 entries from the tail with `LINSERT BEFORE `. LINSERT finds the pivot by scanning from the head, so every requeue costs O(queue length) on Redis's single command thread: about 10 ms per call on a 650k-entry queue, and about 40 ms once earlier requeues have piled up near the tail. A burst of requeues on a large queue can then block every other build that shares the Redis. Instead, read the `offset` + 1 tail entries with LRANGE, drop them with LTRIM, push the requeued entry, and push them back. That gives the same list in O(offset). Queues of at most `offset` + 1 entries, and offsets of zero or less, still push to the head. One intended difference: when the pivot's value also appears nearer the head (a duplicate entry), LINSERT inserted before that first copy, often sending the test to the back of the queue. It now lands at the offset. --- redis/requeue.lua | 29 +++++++++--- redis/reserve.lua | 25 +++++++--- ruby/test/ci/queue/redis_test.rb | 81 ++++++++++++++++++++++++++++++++ 3 files changed, 121 insertions(+), 14 deletions(-) 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']