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
8 changes: 6 additions & 2 deletions app/models/solid_queue/process.rb
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,14 @@ def self.register(**attributes)

def heartbeat
# Clear any previous changes before locking, for example, in case a previous heartbeat
# failed because of a DB issue (with SQLite depending on configuration, a BusyException
# is not rare) and we still have the unpersisted value
# failed because of a DB issue and we still have the unpersisted value
restore_attributes
with_lock { touch(:last_heartbeat_at) }
rescue
# touch writes the attribute before persisting; don't let a failed
# update leave this object claiming a heartbeat that was never persisted
restore_attributes
raise
end

def deregister(pruned: false)
Expand Down
20 changes: 20 additions & 0 deletions lib/solid_queue/processes/registrable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,26 @@ def stop_heartbeat
def heartbeat
process&.heartbeat
rescue ActiveRecord::RecordNotFound
# Our registration is gone: a supervisor pruned it
stop_to_be_replaced
rescue => error
# Errors like a dropped database connection prevent the
# heartbeat from going through, and even from finding out whether the
# registration is still there
stop_to_be_replaced if presumed_dead?
raise error
end

# Whether this process's registration is prunable: if the last heartbeat that
# we were able to persist is older than the alive threshold, supervisors
# consider this one dead and would have pruned its registration by now
def presumed_dead?
process && process.last_heartbeat_at <= SolidQueue.process_alive_threshold.ago
end

# Deregister locally and wake the run loop, which stops when
# unregistered, so the supervisor replaces this process
def stop_to_be_replaced
self.process = nil
wake_up
end
Expand Down
13 changes: 13 additions & 0 deletions test/models/solid_queue/process_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -80,4 +80,17 @@ class SolidQueue::ProcessTest < ActiveSupport::TestCase
process.heartbeat
end
end

test "a heartbeat that fails to persist doesn't leave a fresh in-memory timestamp" do
process = SolidQueue::Process.register(kind: "Worker", pid: 42, name: "worker-42")
persisted_heartbeat_at = process.last_heartbeat_at

process.stubs(:_update_row).raises(ActiveRecord::StatementInvalid.new("no connection"))

travel 1.minute do
assert_raises(ActiveRecord::StatementInvalid) { process.heartbeat }
end

assert_equal persisted_heartbeat_at, process.last_heartbeat_at
end
end
30 changes: 30 additions & 0 deletions test/unit/process_recovery_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,36 @@ class ProcessRecoveryTest < ActiveSupport::TestCase
JobResult.delete_all
end

test "alive scheduler whose registration is pruned is torn down and replaced" do
old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 1.second

@pid = run_supervisor_as_fork(skip_recurring: false)
wait_for_registered_processes(5, timeout: 3.seconds) # supervisor + 2 workers + dispatcher + scheduler

scheduler_process = SolidQueue::Process.find_by(kind: "Scheduler")
assert scheduler_process.present?

# Simulate another process's prune sweep removing the scheduler's registration
# while the scheduler itself is still alive
scheduler_process.delete

# The scheduler should notice its registration is gone on its next heartbeat,
# terminate, and be replaced by the supervisor with a fresh registration
wait_while_with_timeout(10.seconds) do
skip_active_record_query_cache do
SolidQueue::Process.where(kind: "Scheduler").where.not(id: scheduler_process.id).none?
end
end

skip_active_record_query_cache do
new_scheduler_process = SolidQueue::Process.where(kind: "Scheduler").last
assert new_scheduler_process.present?
assert_not_equal scheduler_process.id, new_scheduler_process.id
end
ensure
SolidQueue.process_heartbeat_interval = old_heartbeat_interval
end

test "supervisor handles missing process record and fails claimed executions properly" do
# Start a supervisor with one worker
@pid = run_supervisor_as_fork(workers: [ { queues: "*", polling_interval: 0.1, processes: 1 } ])
Expand Down
20 changes: 20 additions & 0 deletions test/unit/worker_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,26 @@ class WorkerTest < ActiveSupport::TestCase
SolidQueue.process_heartbeat_interval = old_heartbeat_interval
end

test "terminate when heartbeats have been failing for longer than the alive threshold" do
old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.2.seconds
old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.5.seconds

SolidQueue::Process.any_instance.stubs(:heartbeat).raises(ActiveRecord::StatementInvalid.new("connection lost"))

@worker.start
wait_for_registered_processes(1, timeout: 1.second)

assert_not @worker.pool.shutdown?

# Heartbeats keep failing without the registration being confirmed gone, so
# the worker should give up once the failures outlast the alive threshold
wait_while_with_timeout(3) { !@worker.pool.shutdown? }
assert @worker.pool.shutdown?
ensure
SolidQueue.process_heartbeat_interval = old_heartbeat_interval
SolidQueue.process_alive_threshold = old_alive_threshold
end

test "sleeps `10.minutes` if at capacity" do
3.times { |i| StoreResultJob.perform_later(i, pause: 5.seconds) }

Expand Down
Loading