Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion redis/acknowledge.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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
5 changes: 4 additions & 1 deletion redis/heartbeat.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
3 changes: 3 additions & 0 deletions redis/release.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions redis/requeue.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
Expand All @@ -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
9 changes: 9 additions & 0 deletions redis/reserve.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions redis/reserve_lost.lua
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,19 @@ 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
if redis.call('sismember', processed_key, test) == 0 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) -- 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
Expand Down
15 changes: 12 additions & 3 deletions ruby/lib/ci/queue/redis/base.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -261,6 +268,7 @@ def boot!
@processed_key,
@owners_key,
@worker_queue_key,
@redis_ttl.to_s,
in: child_read,
out: child_write,
)
Expand Down Expand Up @@ -335,6 +343,7 @@ def heartbeat_process
key('processed'),
key('owners'),
key('worker', worker_id, 'queue'),
config.redis_ttl,
)
end

Expand Down
5 changes: 4 additions & 1 deletion ruby/lib/ci/queue/redis/build_record.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
8 changes: 5 additions & 3 deletions ruby/lib/ci/queue/redis/monitor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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(".")
Expand Down
10 changes: 5 additions & 5 deletions ruby/lib/ci/queue/redis/worker.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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

Expand All @@ -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
Expand Down
36 changes: 31 additions & 5 deletions ruby/lib/rspec/queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down
39 changes: 39 additions & 0 deletions ruby/test/ci/queue/redis_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading