From 3fded8a99770688857114d06c264d7b83633ac62 Mon Sep 17 00:00:00 2001 From: Paul Simpson Date: Thu, 10 Sep 2026 18:43:10 -0500 Subject: [PATCH 1/2] fix: expire every build key at write time The EXPIRE for `processed` and `worker::queue` lived at the tail of Worker#poll, and `running`, `owners`, `requeues-count`, `warnings`, `test_failed_count` and `created-at` never got one at all. Any worker that is SIGKILLed, cancelled, OOMs or loses its Redis connection never reaches the end of poll, so those keys leaked permanently. On Kajabi's ci-queue node that reached 464k permanent keys holding 9.5 GiB: `processed` 4.6 GiB, `worker::queue` 3.8 GiB, `requeues-count` 1.0 GiB. Because the node runs `volatile-lru`, the leaked keys were the only ones ineligible for eviction, so memory pressure fell entirely on the live queue keys of in-flight builds -- producing red spec lanes with no failing spec. Every key is now expired at the point of write, so no key's lifetime depends on a clean shutdown. The TTL is threaded through the Lua scripts as an argument and through the spawned heartbeat monitor process. Also: `rspec-queue --report` exited 0 against a queue whose keys were gone. An evicted queue reads as exhausted and its error-report hash reads as empty, so the report certified a run whose results it no longer had. It now refuses to report a result when `total` is 0 or no progress was recorded, and explains that this is eviction rather than a test failure. --- redis/acknowledge.lua | 7 ++++- redis/heartbeat.lua | 5 +++- redis/release.lua | 3 ++ redis/requeue.lua | 3 ++ redis/reserve.lua | 9 ++++++ redis/reserve_lost.lua | 6 ++++ ruby/lib/ci/queue/redis/base.rb | 15 ++++++++-- ruby/lib/ci/queue/redis/build_record.rb | 5 +++- ruby/lib/ci/queue/redis/monitor.rb | 8 +++-- ruby/lib/ci/queue/redis/worker.rb | 10 +++---- ruby/lib/ci/queue/version.rb | 2 +- ruby/lib/rspec/queue.rb | 36 +++++++++++++++++++---- ruby/test/ci/queue/redis_test.rb | 39 +++++++++++++++++++++++++ 13 files changed, 128 insertions(+), 20 deletions(-) diff --git a/redis/acknowledge.lua b/redis/acknowledge.lua index 2da9f6e7..e484f378 100644 --- a/redis/acknowledge.lua +++ b/redis/acknowledge.lua @@ -3,7 +3,12 @@ local processed_key = KEYS[2] local owners_key = KEYS[3] local test = ARGV[1] +local ttl = ARGV[2] redis.call('zrem', zset_key, test) redis.call('hdel', owners_key, test) -- Doesn't matter if it was reclaimed by another workers -return redis.call('sadd', processed_key, test) +local acknowledged = redis.call('sadd', processed_key, test) + +redis.call('expire', processed_key, ttl) + +return acknowledged diff --git a/redis/heartbeat.lua b/redis/heartbeat.lua index 6c9c1e9d..730c1ba7 100644 --- a/redis/heartbeat.lua +++ b/redis/heartbeat.lua @@ -5,6 +5,7 @@ local worker_queue_key = KEYS[4] local current_time = ARGV[1] local test = ARGV[2] +local ttl = ARGV[3] -- already processed, we do not need to bump the timestamp if redis.call('sismember', processed_key, test) == 1 then @@ -13,5 +14,7 @@ end -- we're still the owner of the test, we can bump the timestamp if redis.call('hget', owners_key, test) == worker_queue_key then - return redis.call('zadd', zset_key, current_time, test) + local result = redis.call('zadd', zset_key, current_time, test) + redis.call('expire', zset_key, ttl) + return result end diff --git a/redis/release.lua b/redis/release.lua index 4bcc5929..7e1ed218 100644 --- a/redis/release.lua +++ b/redis/release.lua @@ -2,12 +2,15 @@ local zset_key = KEYS[1] local worker_queue_key = KEYS[2] local owners_key = KEYS[3] +local ttl = ARGV[1] + -- owned_tests = {"SomeTest", "worker:1", "SomeOtherTest", "worker:2", ...} local owned_tests = redis.call('hgetall', owners_key) for index, owner_or_test in ipairs(owned_tests) do if owner_or_test == worker_queue_key then -- If we owned a test local test = owned_tests[index - 1] redis.call('zadd', zset_key, "0", test) -- We expire the lease immediately + redis.call('expire', zset_key, ttl) return nil end end diff --git a/redis/requeue.lua b/redis/requeue.lua index a1fcc25c..91e6cc69 100644 --- a/redis/requeue.lua +++ b/redis/requeue.lua @@ -9,6 +9,7 @@ local max_requeues = tonumber(ARGV[1]) local global_max_requeues = tonumber(ARGV[2]) local test = ARGV[3] local offset = ARGV[4] +local ttl = ARGV[5] if redis.call('hget', owners_key, test) == worker_queue_key then redis.call('hdel', owners_key, test) @@ -30,6 +31,7 @@ end redis.call('hincrby', requeues_count_key, '___total___', 1) redis.call('hincrby', requeues_count_key, test, 1) +redis.call('expire', requeues_count_key, ttl) local pivot = redis.call('lrange', queue_key, -1 - offset, 0 - offset)[1] if pivot then @@ -38,6 +40,7 @@ else redis.call('lpush', queue_key, test) end +redis.call('expire', queue_key, ttl) redis.call('zrem', zset_key, test) return true diff --git a/redis/reserve.lua b/redis/reserve.lua index d5951146..187ee910 100644 --- a/redis/reserve.lua +++ b/redis/reserve.lua @@ -5,12 +5,21 @@ local worker_queue_key = KEYS[4] local owners_key = KEYS[5] local current_time = ARGV[1] +local ttl = ARGV[2] local test = redis.call('rpop', queue_key) if test then redis.call('zadd', zset_key, current_time, test) redis.call('lpush', worker_queue_key, test) redis.call('hset', owners_key, test, worker_queue_key) + + -- Expire at write time. These keys used to rely on an EXPIRE at the tail of + -- CI::Queue::Redis::Worker#poll, which never runs for a worker that is killed, + -- cancelled, OOMs or loses its Redis connection -- leaking the key forever. + redis.call('expire', zset_key, ttl) + redis.call('expire', worker_queue_key, ttl) + redis.call('expire', owners_key, ttl) + return test else return nil diff --git a/redis/reserve_lost.lua b/redis/reserve_lost.lua index 9dfaa616..c33ca049 100644 --- a/redis/reserve_lost.lua +++ b/redis/reserve_lost.lua @@ -5,6 +5,7 @@ local owners_key = KEYS[4] local current_time = ARGV[1] local timeout = ARGV[2] +local ttl = ARGV[3] local lost_tests = redis.call('zrangebyscore', zset_key, 0, current_time - timeout) for _, test in ipairs(lost_tests) do @@ -12,6 +13,11 @@ for _, test in ipairs(lost_tests) do redis.call('zadd', zset_key, current_time, test) redis.call('lpush', worker_queue_key, test) redis.call('hset', owners_key, test, worker_queue_key) -- Take ownership + + redis.call('expire', zset_key, ttl) + redis.call('expire', worker_queue_key, ttl) + redis.call('expire', owners_key, ttl) + return test end end diff --git a/ruby/lib/ci/queue/redis/base.rb b/ruby/lib/ci/queue/redis/base.rb index 9f42c919..31d2d4e8 100644 --- a/ruby/lib/ci/queue/redis/base.rb +++ b/ruby/lib/ci/queue/redis/base.rb @@ -118,7 +118,10 @@ def expired? end def created_at=(timestamp) - redis.setnx(key('created-at'), timestamp) + redis.pipelined do |pipeline| + pipeline.setnx(key('created-at'), timestamp) + pipeline.expire(key('created-at'), config.redis_ttl) + end end def size @@ -182,7 +185,10 @@ def queue_initializing? end def increment_test_failed - redis.incr(key('test_failed_count')) + redis.pipelined do |pipeline| + pipeline.incr(key('test_failed_count')) + pipeline.expire(key('test_failed_count'), config.redis_ttl) + end end def test_failed @@ -241,7 +247,8 @@ def read_script(name) end class HeartbeatProcess - def initialize(redis_url, zset_key, processed_key, owners_key, worker_queue_key) + def initialize(redis_url, zset_key, processed_key, owners_key, worker_queue_key, redis_ttl) + @redis_ttl = redis_ttl @redis_url = redis_url @zset_key = zset_key @processed_key = processed_key @@ -261,6 +268,7 @@ def boot! @processed_key, @owners_key, @worker_queue_key, + @redis_ttl.to_s, in: child_read, out: child_write, ) @@ -335,6 +343,7 @@ def heartbeat_process key('processed'), key('owners'), key('worker', worker_id, 'queue'), + config.redis_ttl, ) end diff --git a/ruby/lib/ci/queue/redis/build_record.rb b/ruby/lib/ci/queue/redis/build_record.rb index b1329ff9..75f0946b 100644 --- a/ruby/lib/ci/queue/redis/build_record.rb +++ b/ruby/lib/ci/queue/redis/build_record.rb @@ -53,7 +53,10 @@ def pop_warnings end def record_warning(type, attributes) - redis.rpush(key('warnings'), Marshal.dump([type, attributes])) + redis.pipelined do |pipeline| + pipeline.rpush(key('warnings'), Marshal.dump([type, attributes])) + pipeline.expire(key('warnings'), config.redis_ttl) + end end def record_error(id, payload, stats: nil) diff --git a/ruby/lib/ci/queue/redis/monitor.rb b/ruby/lib/ci/queue/redis/monitor.rb index bdf1466e..56ec6c4e 100755 --- a/ruby/lib/ci/queue/redis/monitor.rb +++ b/ruby/lib/ci/queue/redis/monitor.rb @@ -13,7 +13,8 @@ class Monitor DEV_SCRIPTS_ROOT = ::File.expand_path('../../../../../../redis', __FILE__) RELEASE_SCRIPTS_ROOT = ::File.expand_path('../../redis', __FILE__) - def initialize(pipe, logger, redis_url, zset_key, processed_key, owners_key, worker_queue_key) + def initialize(pipe, logger, redis_url, zset_key, processed_key, owners_key, worker_queue_key, redis_ttl) + @redis_ttl = redis_ttl @zset_key = zset_key @processed_key = processed_key @owners_key = owners_key @@ -40,7 +41,7 @@ def process_tick!(id:) eval_script( :heartbeat, keys: [@zset_key, @processed_key, @owners_key, @worker_queue_key], - argv: [Time.now.to_f, id] + argv: [Time.now.to_f, id, @redis_ttl] ) rescue => error @logger.info(error) @@ -142,9 +143,10 @@ def monitor processed_key = ARGV[2] owners_key = ARGV[3] worker_queue_key = ARGV[4] +redis_ttl = ARGV[5] logger.debug("Starting monitor: #{redis_url} #{zset_key} #{processed_key}") -manager = CI::Queue::Redis::Monitor.new($stdin, logger, redis_url, zset_key, processed_key, owners_key, worker_queue_key) +manager = CI::Queue::Redis::Monitor.new($stdin, logger, redis_url, zset_key, processed_key, owners_key, worker_queue_key, redis_ttl) # Notify the parent we're ready $stdout.puts(".") diff --git a/ruby/lib/ci/queue/redis/worker.rb b/ruby/lib/ci/queue/redis/worker.rb index c5eb5448..702cee8a 100644 --- a/ruby/lib/ci/queue/redis/worker.rb +++ b/ruby/lib/ci/queue/redis/worker.rb @@ -103,7 +103,7 @@ def acknowledge(test) eval_script( :acknowledge, keys: [key('running'), key('processed'), key('owners')], - argv: [test_key], + argv: [test_key, config.redis_ttl], ) == 1 end @@ -122,7 +122,7 @@ def requeue(test, offset: Redis.requeue_offset) key('worker', worker_id, 'queue'), key('owners'), ], - argv: [config.max_requeues, global_max_requeues, test_key, offset], + argv: [config.max_requeues, global_max_requeues, test_key, offset, config.redis_ttl], ) == 1 @reserved_test = test_key unless requeued @@ -133,7 +133,7 @@ def release! eval_script( :release, keys: [key('running'), key('worker', worker_id, 'queue'), key('owners')], - argv: [], + argv: [config.redis_ttl], ) nil end @@ -173,7 +173,7 @@ def try_to_reserve_test key('worker', worker_id, 'queue'), key('owners'), ], - argv: [CI::Queue.time_now.to_f], + argv: [CI::Queue.time_now.to_f, config.redis_ttl], ) end @@ -188,7 +188,7 @@ def try_to_reserve_lost_test key('worker', worker_id, 'queue'), key('owners'), ], - argv: [CI::Queue.time_now.to_f, timeout], + argv: [CI::Queue.time_now.to_f, timeout, config.redis_ttl], ) if lost_test diff --git a/ruby/lib/ci/queue/version.rb b/ruby/lib/ci/queue/version.rb index 899ea4e4..133b0706 100644 --- a/ruby/lib/ci/queue/version.rb +++ b/ruby/lib/ci/queue/version.rb @@ -2,7 +2,7 @@ module CI module Queue - VERSION = '0.66.0' + VERSION = '0.66.1' DEV_SCRIPTS_ROOT = ::File.expand_path('../../../../../redis', __FILE__) RELEASE_SCRIPTS_ROOT = ::File.expand_path('../redis', __FILE__) end diff --git a/ruby/lib/rspec/queue.rb b/ruby/lib/rspec/queue.rb index 97616058..b60079cf 100644 --- a/ruby/lib/rspec/queue.rb +++ b/ruby/lib/rspec/queue.rb @@ -299,14 +299,33 @@ def call(options, stdout, stderr) step("Waiting for workers to complete") - unless supervisor.wait_for_workers - unless supervisor.queue_initialized? - abort! "No master was elected. Did all workers crash?" + begin + unless supervisor.wait_for_workers + unless supervisor.queue_initialized? + abort! "No master was elected. Did all workers crash?" + end + + unless supervisor.exhausted? + abort! "#{supervisor.size} tests weren't run." + end + end + + # PE-3857: an evicted queue is indistinguishable from a completed one -- + # the queue reads as exhausted and the error-report hash reads as empty, + # so we would report success for a run whose results we no longer have. + # Refuse to certify a build whose bookkeeping is gone. + total = supervisor.total + progress = supervisor.build.progress + + if total.zero? + abort! queue_vanished_message("`total` is 0 -- the queue is empty or was never populated") end - unless supervisor.exhausted? - abort! "#{supervisor.size} tests weren't run." + if progress <= 0 + abort! queue_vanished_message("no test progress was recorded (total=#{total}, progress=#{progress})") end + rescue CI::Queue::Redis::LostMaster + abort! queue_vanished_message("the master worker record is gone") end # TODO: better reporting @@ -324,6 +343,13 @@ def call(options, stdout, stderr) private + def queue_vanished_message(detail) + "Refusing to report a result: #{detail}.\n" \ + "The queue's Redis keys are missing, so this run cannot be certified as passing. " \ + "This is not a test failure -- it usually means the keys were evicted (check the " \ + "ci-queue Redis memory usage) or the build ran longer than CI_QUEUE_REDIS_TTL." + end + attr_reader :configuration def setup(options, out, err) diff --git a/ruby/test/ci/queue/redis_test.rb b/ruby/test/ci/queue/redis_test.rb index 5c5bd139..8c1e4b7c 100644 --- a/ruby/test/ci/queue/redis_test.rb +++ b/ruby/test/ci/queue/redis_test.rb @@ -270,8 +270,47 @@ def test_initialise_from_rediss_uri assert_instance_of CI::Queue::Redis::Worker, queue end + # PE-3857: every key ci-queue writes must carry a TTL at write time. + # + # These keys used to get their EXPIRE at the tail of Worker#poll, which never + # runs for a worker that is killed, cancelled, OOMs or loses its Redis + # connection -- so the key leaked forever. On tools-redis-001 that grew to + # 464k permanent keys / 9.5 GiB, which then forced the node to evict the live + # queue keys of in-flight builds under its volatile-lru policy. + # + # Deliberately drives reserve/requeue/acknowledge by hand rather than through + # #poll, so that the end-of-poll EXPIRE never runs -- that is the leak. + def test_every_build_key_has_a_ttl + ttl = 600 + # setup already populated a queue under the same build id with the default + # 8h TTL; start clean so this test only measures its own writes. + @redis.flushdb + queue = worker(1, max_requeues: 1, requeue_tolerance: 1.0, redis_ttl: ttl) + + first = test_for(queue.send(:reserve)) + queue.requeue(first) + second = test_for(queue.send(:reserve)) + queue.acknowledge(second) + + queue.build.record_warning(:some_warning, test: 'x', timeout: 1) + queue.increment_test_failed + queue.created_at = CI::Queue.time_now.to_f + queue.release! + + keys = @redis.keys('build:*') + refute_empty keys + + without_ttl = keys.reject { |key| @redis.ttl(key) > 0 }.sort + assert_equal [], without_ttl, "keys written without a TTL: #{without_ttl.join(', ')}" + assert keys.all? { |key| @redis.ttl(key) <= ttl } + end + private + def test_for(id) + TEST_LIST.find { |test| test.id == id } or raise "unknown test id #{id.inspect}" + end + def shuffled_test_list CI::Queue.shuffle(TEST_LIST, Random.new(0)).freeze end From 7e2564a30b7e9284a7de16f302c01f49c8ea78ff Mon Sep 17 00:00:00 2001 From: Paul Simpson Date: Thu, 10 Sep 2026 18:51:32 -0500 Subject: [PATCH 2/2] revert: keep VERSION at 0.66.0 This fork carries Kajabi-only patches on top of upstream 0.66.0 without bumping the version (see b304463), and the monolith's supply-chain cooldown check fails closed on a version it cannot find on rubygems. The Gemfile.lock revision is the real identifier for a git-sourced gem. --- ruby/lib/ci/queue/version.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/ruby/lib/ci/queue/version.rb b/ruby/lib/ci/queue/version.rb index 133b0706..899ea4e4 100644 --- a/ruby/lib/ci/queue/version.rb +++ b/ruby/lib/ci/queue/version.rb @@ -2,7 +2,7 @@ module CI module Queue - VERSION = '0.66.1' + VERSION = '0.66.0' DEV_SCRIPTS_ROOT = ::File.expand_path('../../../../../redis', __FILE__) RELEASE_SCRIPTS_ROOT = ::File.expand_path('../redis', __FILE__) end