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/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