diff --git a/redis/push_queue.lua b/redis/push_queue.lua new file mode 100644 index 0000000..6f18844 --- /dev/null +++ b/redis/push_queue.lua @@ -0,0 +1,32 @@ +local master_status_key = KEYS[1] +local queue_key = KEYS[2] +local total_key = KEYS[3] +local current_generation_key = KEYS[4] + +local expected_lock = ARGV[1] +local generation = ARGV[2] +local total = ARGV[3] +local redis_ttl = tonumber(ARGV[4]) + +-- Fence a master that resumed after its lease expired and another worker won. +if redis.call('get', master_status_key) ~= expected_lock then + return 0 +end + +-- Publishing the queue and changing the status to ready must be atomic. No +-- worker can observe a ready queue before every test has been enqueued. +redis.call('del', queue_key) +for index = 5, #ARGV do + redis.call('lpush', queue_key, ARGV[index]) +end + +redis.call('set', total_key, total) +redis.call('set', current_generation_key, generation) +redis.call('set', master_status_key, 'ready') + +redis.call('expire', queue_key, redis_ttl) +redis.call('expire', total_key, redis_ttl) +redis.call('expire', current_generation_key, redis_ttl) +redis.call('expire', master_status_key, redis_ttl) + +return 1 diff --git a/redis/renew_master_lock.lua b/redis/renew_master_lock.lua new file mode 100644 index 0000000..f9f3103 --- /dev/null +++ b/redis/renew_master_lock.lua @@ -0,0 +1,10 @@ +local master_status_key = KEYS[1] +local expected_lock = ARGV[1] +local lock_ttl = tonumber(ARGV[2]) + +if redis.call('get', master_status_key) ~= expected_lock then + return 0 +end + +redis.call('expire', master_status_key, lock_ttl) +return 1 diff --git a/redis/store_chunk_metadata.lua b/redis/store_chunk_metadata.lua new file mode 100644 index 0000000..1150b28 --- /dev/null +++ b/redis/store_chunk_metadata.lua @@ -0,0 +1,36 @@ +local master_status_key = KEYS[1] +local chunks_key = KEYS[2] +local test_group_timeout_key = KEYS[3] + +local expected_lock = ARGV[1] +local master_lock_ttl = tonumber(ARGV[2]) +local redis_ttl = tonumber(ARGV[3]) + +-- Metadata writes are fenced too: an expired master must not overwrite the +-- metadata prepared by its replacement. +if redis.call('get', master_status_key) ~= expected_lock then + return 0 +end + +local argument_index = 4 +local chunk_key_index = 4 +while argument_index <= #ARGV do + local chunk_id = ARGV[argument_index] + local chunk_json = ARGV[argument_index + 1] + local chunk_timeout = ARGV[argument_index + 2] + local chunk_key = KEYS[chunk_key_index] + + redis.call('set', chunk_key, chunk_json) + redis.call('expire', chunk_key, redis_ttl) + redis.call('sadd', chunks_key, chunk_id) + redis.call('hset', test_group_timeout_key, chunk_id, chunk_timeout) + + argument_index = argument_index + 3 + chunk_key_index = chunk_key_index + 1 +end + +redis.call('expire', chunks_key, redis_ttl) +redis.call('expire', test_group_timeout_key, redis_ttl) +redis.call('expire', master_status_key, master_lock_ttl) + +return 1 diff --git a/ruby/lib/ci/queue/configuration.rb b/ruby/lib/ci/queue/configuration.rb index 9e9c133..8882f94 100644 --- a/ruby/lib/ci/queue/configuration.rb +++ b/ruby/lib/ci/queue/configuration.rb @@ -15,6 +15,7 @@ class Configuration attr_accessor :timing_redis_url attr_accessor :write_duration_averages attr_accessor :heartbeat_grace_period, :heartbeat_interval + attr_accessor :master_lock_ttl, :max_election_attempts attr_reader :circuit_breakers attr_writer :seed, :build_id attr_writer :queue_init_timeout, :report_timeout, :inactive_workers_timeout @@ -66,7 +67,9 @@ def initialize( branch: nil, timing_redis_url: nil, heartbeat_grace_period: 30, - heartbeat_interval: 10 + heartbeat_interval: 10, + master_lock_ttl: 30, + max_election_attempts: 3 ) @build_id = build_id @circuit_breakers = [CircuitBreaker::Disabled] @@ -105,6 +108,8 @@ def initialize( @write_duration_averages = false @heartbeat_grace_period = heartbeat_grace_period @heartbeat_interval = heartbeat_interval + @master_lock_ttl = master_lock_ttl + @max_election_attempts = max_election_attempts end def queue_init_timeout diff --git a/ruby/lib/ci/queue/redis.rb b/ruby/lib/ci/queue/redis.rb index c3876ef..96f7606 100644 --- a/ruby/lib/ci/queue/redis.rb +++ b/ruby/lib/ci/queue/redis.rb @@ -18,6 +18,7 @@ module Queue module Redis Error = Class.new(StandardError) LostMaster = Class.new(Error) + MasterDied = Class.new(Error) class << self diff --git a/ruby/lib/ci/queue/redis/base.rb b/ruby/lib/ci/queue/redis/base.rb index 9ebdd58..4016dc5 100644 --- a/ruby/lib/ci/queue/redis/base.rb +++ b/ruby/lib/ci/queue/redis/base.rb @@ -53,12 +53,19 @@ def progress total - size end - def wait_for_master(timeout: 120) + def wait_for_master(timeout: 120, fail_if_unclaimed: false) return true if master? + last_status = nil (timeout * 10 + 1).to_i.times do - return true if queue_initialized? + status = master_status + return true if %w[ready finished].include?(status) + + if status.nil? && (last_status == 'setup' || fail_if_unclaimed) + raise MasterDied, 'The master lease expired during queue setup.' + end + last_status = status sleep 0.1 end raise LostMaster, "The master worker (worker #{master_worker_id}) is still `#{master_status}` after #{timeout} seconds waiting." @@ -110,7 +117,10 @@ def build_id end def master_status - redis.get(key('master-status')) + status = redis.get(key('master-status')) + return 'setup' if status&.start_with?('setup:') + + status end def eval_script(script, *args) diff --git a/ruby/lib/ci/queue/redis/worker.rb b/ruby/lib/ci/queue/redis/worker.rb index 7de1f17..a951103 100644 --- a/ruby/lib/ci/queue/redis/worker.rb +++ b/ruby/lib/ci/queue/redis/worker.rb @@ -2,6 +2,7 @@ require 'ci/queue/static' require 'concurrent/set' +require 'securerandom' module CI module Queue @@ -38,19 +39,43 @@ def populate(tests, random: Random.new) @index = tests.map { |t| [t.id, t] }.to_h @total = tests.size - if acquire_master_role? - executables = reorder_tests(tests, random: random) + election_attempts = 0 - chunks = executables.select { |e| e.is_a?(CI::Queue::TestChunk) } - individual_tests = executables.reject { |e| e.is_a?(CI::Queue::TestChunk) } + loop do + begin + if acquire_master_role? + all_ids = with_master_lock_renewal do + executables = reorder_tests(tests, random: random) - store_chunk_metadata(chunks) if chunks.any? + chunks = executables.select { |e| e.is_a?(CI::Queue::TestChunk) } + individual_tests = executables.reject { |e| e.is_a?(CI::Queue::TestChunk) } - all_ids = chunks.map(&:id) + individual_tests.map(&:id) - push(all_ids) - end + store_chunk_metadata(chunks) if chunks.any? + + chunks.map(&:id) + individual_tests.map(&:id) + end + push(all_ids) + else + rescue_connection_errors do + wait_for_master(timeout: config.queue_init_timeout, fail_if_unclaimed: true) + end + end - register_worker_presence + register_worker_presence + break + rescue MasterDied => error + election_attempts += 1 + if election_attempts >= config.max_election_attempts + raise LostMaster, + "Failed to recover queue setup after #{election_attempts} election attempts: #{error.message}" + end + + warn 'Master worker died during setup; retrying election ' \ + "(#{election_attempts}/#{config.max_election_attempts})." + @master = nil + @generation = nil + end + end self end @@ -441,15 +466,23 @@ def push(tests) @total = tests.size if @master - redis.multi do |transaction| - transaction.lpush(key('queue'), tests) unless tests.empty? - transaction.set(key('total'), @total) - transaction.set(key('master-status'), 'ready') - - transaction.expire(key('queue'), config.redis_ttl) - transaction.expire(key('total'), config.redis_ttl) - transaction.expire(key('master-status'), config.redis_ttl) - end + result = eval_script( + :push_queue, + keys: [ + key('master-status'), + key('queue'), + key('total'), + key('current-generation') + ], + argv: [ + master_lock_value, + @generation, + @total, + config.redis_ttl, + *tests + ] + ) + raise MasterDied, 'The master lease was lost before the queue could be published.' unless result == 1 end rescue *CONNECTION_ERRORS raise if @master @@ -462,24 +495,89 @@ def register def acquire_master_role? return true if @master - @master = redis.setnx(key('master-status'), 'setup') + @generation = SecureRandom.uuid + @master = redis.set( + key('master-status'), + master_lock_value, + nx: true, + ex: config.master_lock_ttl + ) if @master begin - redis.set(key('master-worker-id'), worker_id) - redis.expire(key('master-worker-id'), config.redis_ttl) - warn "Worker #{worker_id} elected as master" + redis.multi do |transaction| + transaction.set(key('master-worker-id'), worker_id) + transaction.expire(key('master-worker-id'), config.redis_ttl) + end + warn "Worker #{worker_id} elected as master (generation #{@generation})" rescue *CONNECTION_ERRORS # If setting master-worker-id fails, we still have master status # Log but don't lose master role warn("Failed to set master-worker-id: #{$!.message}") end + else + @generation = nil end @master rescue *CONNECTION_ERRORS @master = nil + @generation = nil false end + def master_lock_value + "setup:#{@generation}" + end + + def renew_master_lock! + result = eval_script( + :renew_master_lock, + keys: [key('master-status')], + argv: [master_lock_value, config.master_lock_ttl] + ) + raise MasterDied, 'The master lease expired while the queue was being populated.' unless result == 1 + end + + def with_master_lock_renewal + renewal_interval = [config.master_lock_ttl / 3.0, 0.1].max + lock_lost = false + stop_renewal = false + script = read_script(:renew_master_lock) + status_key = key('master-status') + expected_lock = master_lock_value + + renewal_thread = Thread.new do + renewal_redis = ::Redis.new(url: redis_url) + until stop_renewal + sleep renewal_interval + break if stop_renewal + + result = renewal_redis.eval( + script, + keys: [status_key], + argv: [expected_lock, config.master_lock_ttl] + ) + unless result == 1 + lock_lost = true + break + end + end + rescue *CONNECTION_ERRORS => error + warn "Failed to renew the master lease: #{error.message}" + ensure + renewal_redis&.close + end + + result = yield + raise MasterDied, 'The master lease expired while the queue was being populated.' if lock_lost + + renew_master_lock! + result + ensure + stop_renewal = true + renewal_thread&.kill + renewal_thread&.join + end + def register_worker_presence register redis.expire(key('workers'), config.redis_ttl) @@ -493,31 +591,29 @@ def store_chunk_metadata(chunks) batch_size = 5 # 5 chunks = 20 commands + 2 expires = 22 commands per batch chunks.each_slice(batch_size) do |chunk_batch| - redis.multi do |transaction| - chunk_batch.each do |chunk| - # Store chunk metadata with TTL - transaction.set( - key('chunk', chunk.id), - chunk.to_json - ) - transaction.expire(key('chunk', chunk.id), config.redis_ttl) - - # Track all chunks for cleanup - transaction.sadd(key('chunks'), chunk.id) - - # Store dynamic timeout for this chunk - # Timeout = estimated_duration (in ms) converted to seconds + buffer - # estimated_duration is in milliseconds, convert to seconds and add 10% buffer - buffer_percent = 10 - estimated_duration_seconds = chunk.estimated_duration / 1000.0 - chunk_timeout = (estimated_duration_seconds * (1 + buffer_percent / 100.0)).round(2) - # Format to string to avoid floating point precision issues in Redis - # Use %g to remove trailing zeros - transaction.hset(key('test-group-timeout'), chunk.id, format('%g', chunk_timeout)) - end - transaction.expire(key('chunks'), config.redis_ttl) - transaction.expire(key('test-group-timeout'), config.redis_ttl) + chunk_keys = chunk_batch.map { |chunk| key('chunk', chunk.id) } + chunk_data = chunk_batch.flat_map do |chunk| + # Timeout = estimated duration in seconds plus a 10% buffer. + chunk_timeout = (chunk.estimated_duration / 1000.0 * 1.1).round(2) + [chunk.id, chunk.to_json, format('%g', chunk_timeout)] end + + result = eval_script( + :store_chunk_metadata, + keys: [ + key('master-status'), + key('chunks'), + key('test-group-timeout'), + *chunk_keys + ], + argv: [ + master_lock_value, + config.master_lock_ttl, + config.redis_ttl, + *chunk_data + ] + ) + raise MasterDied, 'The master lease was lost while storing chunk metadata.' unless result == 1 end end diff --git a/ruby/test/ci/queue/configuration_test.rb b/ruby/test/ci/queue/configuration_test.rb index 64d7410..a371234 100644 --- a/ruby/test/ci/queue/configuration_test.rb +++ b/ruby/test/ci/queue/configuration_test.rb @@ -73,6 +73,13 @@ def test_redis_ttl_defaults assert_equal(28_800, config.redis_ttl) end + def test_master_election_defaults + config = Configuration.new + + assert_equal 30, config.master_lock_ttl + assert_equal 3, config.max_election_attempts + end + def test_redis_ttl_from_env config = Configuration.from_env( "CI_QUEUE_REDIS_TTL" => "14400" diff --git a/ruby/test/ci/queue/redis_master_recovery_test.rb b/ruby/test/ci/queue/redis_master_recovery_test.rb new file mode 100644 index 0000000..39d0c47 --- /dev/null +++ b/ruby/test/ci/queue/redis_master_recovery_test.rb @@ -0,0 +1,117 @@ +# frozen_string_literal: true + +require 'test_helper' + +class CI::Queue::Redis::MasterRecoveryTest < Minitest::Test + include QueueHelper + + BUILD_ID = 'master-recovery' + TEST_IDS = %w[ + ATest#test_foo + ATest#test_bar + BTest#test_foo + BTest#test_bar + ].freeze + + MockTest = Struct.new(:id) do + def <=>(other) + id <=> other.id + end + + def flaky? + false + end + end + + def setup + @redis_url = ENV.fetch('REDIS_URL', 'redis://localhost:6379/0') + @redis = ::Redis.new(url: @redis_url) + @redis.flushdb + end + + def teardown + @redis.flushdb + end + + def test_surviving_worker_repopulates_after_master_dies_during_setup + # Model a worker that won election and was terminated before it could + # publish the queue. The lease expires without any process cleaning it up. + @redis.set(master_status_key, 'setup:dead-generation', ex: 1) + + survivor = worker('survivor') + survivor.populate(tests, random: Random.new(0)) + + assert_predicate survivor, :master? + refute_equal 'dead-generation', @redis.get(current_generation_key) + + executed = poll(survivor).map(&:id) + assert_equal TEST_IDS.sort, executed.sort + assert_equal TEST_IDS.length, executed.length + end + + def test_stale_master_cannot_overwrite_replacement_queue + dead_master = worker('dead-master', master_lock_ttl: 1) + assert dead_master.send(:acquire_master_role?) + dead_generation = dead_master.instance_variable_get(:@generation) + + sleep 1.1 + + survivor = worker('survivor', master_lock_ttl: 1) + survivor.populate(tests, random: Random.new(0)) + replacement_generation = @redis.get(current_generation_key) + queue_before_stale_push = @redis.lrange("build:#{BUILD_ID}:queue", 0, -1) + + assert_raises(CI::Queue::Redis::MasterDied) do + dead_master.send(:push, ['StaleTest#test_should_not_run']) + end + + refute_equal dead_generation, replacement_generation + assert_equal replacement_generation, @redis.get(current_generation_key) + assert_equal queue_before_stale_push, @redis.lrange("build:#{BUILD_ID}:queue", 0, -1) + refute_includes queue_before_stale_push, 'StaleTest#test_should_not_run' + end + + def test_master_renews_lease_during_slow_population + master = worker('slow-master', master_lock_ttl: 1) + + slow_reorder = lambda do |passed_tests, **_args| + sleep 1.2 + passed_tests + end + master.stub(:reorder_tests, slow_reorder) do + master.populate(tests, random: Random.new(0)) + end + + assert_predicate master, :master? + assert_equal 'ready', @redis.get(master_status_key) + refute_nil @redis.get(current_generation_key) + end + + private + + def current_generation_key + "build:#{BUILD_ID}:current-generation" + end + + def master_status_key + "build:#{BUILD_ID}:master-status" + end + + def tests + TEST_IDS.map { |id| MockTest.new(id) } + end + + def worker(id, **options) + CI::Queue::Redis.new( + @redis_url, + CI::Queue::Configuration.new( + build_id: BUILD_ID, + worker_id: id, + timeout: 0.2, + queue_init_timeout: 3, + timing_redis_url: @redis_url, + **options, + ), + ) + end +end