Skip to content

Commit 09e05f9

Browse files
committed
fix: Bound recovery scans and verify crash safety
1 parent 6c94353 commit 09e05f9

8 files changed

Lines changed: 272 additions & 10 deletions

File tree

‎benchmark/support.rb‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -428,9 +428,11 @@ def migrate
428428
require_relative "../db/migrate/20260805000000_create_solid_objects_tables"
429429
require_relative "../db/migrate/20260806000000_add_state_revision_to_solid_objects_instances"
430430
require_relative "../db/migrate/20260813000000_rename_message_dispatch_columns"
431+
require_relative "../db/migrate/20260915000000_add_solid_objects_effect_recoveries"
431432
CreateSolidObjectsTables.new.migrate(:up)
432433
AddStateRevisionToSolidObjectsInstances.new.migrate(:up)
433434
RenameMessageDispatchColumns.new.migrate(:up)
435+
AddSolidObjectsEffectRecoveries.new.migrate(:up)
434436
end
435437

436438
# @rbs () -> void
@@ -444,6 +446,7 @@ def load_models
444446
claimed_message
445447
reminder
446448
effect
449+
effect_recovery
447450
broadcast
448451
dead_letter
449452
].each do |model|

‎docs/effect-recovery.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,13 @@ Mailbox insertion reuses the instance lock already held by the decision.
149149
Multiple checks in one actor commit lock all their effects and bindings before
150150
locking any processes. Unlocked candidate reads are hints, never decisions.
151151

152+
Automatic passes prefilter owner freshness and the effective per-effect timeout
153+
using database time, and visit at most `claim_scan_limit` stale candidates. Fresh
154+
actors and owners are not locked, including owners protected by extended grace.
155+
Remaining stale effects are revisited on later polls. Every candidate still
156+
undergoes the authoritative locked recheck, and successful retirement announces
157+
the committed mailbox work through the existing wake-up mechanism.
158+
152159
The decision samples database wall time after obtaining the owner lock. The
153160
effective timeout is the larger of the runtime's `process_alive_threshold` and
154161
the effect's persisted `recovery_timeout`, in seconds. A heartbeat newer than

‎examples/at_least_once/boot.rb‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ def self.call(database_path)
2323
require "solid_objects/database_adapter"
2424
%w[
2525
record process instance message ready_message claimed_message
26-
reminder effect broadcast dead_letter
26+
reminder effect effect_recovery broadcast dead_letter
2727
].each { |model| require File.join(ROOT, "app/models/solid_objects", model) }
2828

2929
SolidObjects.configuration.authorize_message = ->(**) { true }
@@ -40,8 +40,10 @@ def self.migrate
4040
require File.join(ROOT, "db/migrate/20260805000000_create_solid_objects_tables")
4141
require File.join(ROOT, "db/migrate/20260806000000_add_state_revision_to_solid_objects_instances")
4242
require File.join(ROOT, "db/migrate/20260813000000_rename_message_dispatch_columns")
43+
require File.join(ROOT, "db/migrate/20260915000000_add_solid_objects_effect_recoveries")
4344
CreateSolidObjectsTables.new.migrate(:up)
4445
AddStateRevisionToSolidObjectsInstances.new.migrate(:up)
4546
RenameMessageDispatchColumns.new.migrate(:up)
47+
AddSolidObjectsEffectRecoveries.new.migrate(:up)
4648
end
4749
end

‎lib/solid_objects/actor.rb‎

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -185,12 +185,7 @@ def emit(name, on_success: nil, on_failure: nil, on_recovery: nil, on_status: ni
185185
validate_effect_callback!(on_failure)
186186
validate_effect_callback!(on_recovery)
187187
validate_effect_callback!(on_status)
188-
unless recovery_timeout.nil?
189-
unless recovery_timeout.is_a?(Numeric) && recovery_timeout.real? && recovery_timeout.to_f.finite? && recovery_timeout.positive?
190-
raise ArgumentError, "recovery_timeout must be a positive finite duration in seconds"
191-
end
192-
raise ArgumentError, "recovery_timeout requires on_recovery" unless on_recovery
193-
end
188+
validate_recovery_timeout!(timeout: recovery_timeout, operation: on_recovery)
194189
effect_id = SecureRandom.uuid
195190
EffectIntent.new(
196191
effect_id:,
@@ -446,5 +441,15 @@ def validate_effect_callback!(operation)
446441

447442
raise UnknownMessage, "unknown effect callback operation #{operation.inspect}"
448443
end
444+
445+
# @rbs (timeout: Numeric?, operation: String | Symbol?) -> void
446+
def validate_recovery_timeout!(timeout:, operation:)
447+
return if timeout.nil?
448+
449+
unless timeout.is_a?(Numeric) && timeout.real? && timeout.to_f.finite? && timeout.to_f.positive?
450+
raise ArgumentError, "recovery_timeout must be a positive finite duration in seconds"
451+
end
452+
raise ArgumentError, "recovery_timeout requires on_recovery" unless operation
453+
end
449454
end
450455
end

‎lib/solid_objects/effect_recovery_coordinator.rb‎

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -35,9 +35,7 @@ def check(instance:, intents:)
3535

3636
# @rbs () -> void
3737
def recover_available
38-
candidates = EffectRecovery.where(retired_at: nil).where.not(recovery_operation: nil)
39-
.where(effect_id: Effect.where(status: "processing").select(:effect_id))
40-
candidates.find_each do |candidate|
38+
recovery_candidates.each do |candidate|
4139
notification = SolidObjects.database_adapter.transaction do
4240
instance = Instance.lock.find_by(id: candidate.instance_id)
4341
next unless instance
@@ -60,6 +58,26 @@ def recover_available
6058

6159
private
6260

61+
# @rbs () -> ActiveRecord::Relation[EffectRecovery]
62+
def recovery_candidates
63+
effects = Effect.table_name
64+
owners = Process.table_name
65+
bindings = EffectRecovery.table_name
66+
heartbeat = case DatabaseAdapter.family(Record.connection)
67+
when :postgresql then "EXTRACT(EPOCH FROM #{owners}.last_heartbeat_at)"
68+
when :mysql then "UNIX_TIMESTAMP(#{owners}.last_heartbeat_at)"
69+
else "CAST(STRFTIME('%s', #{owners}.last_heartbeat_at) AS REAL)"
70+
end
71+
now = SolidObjects.database_adapter.database_clock_now.to_f
72+
threshold = SolidObjects.configuration.process_alive_threshold
73+
EffectRecovery.joins("INNER JOIN #{effects} ON #{effects}.effect_id = #{bindings}.effect_id")
74+
.joins("LEFT JOIN #{owners} ON #{owners}.id = #{effects}.claimed_by")
75+
.where(retired_at: nil).where.not(recovery_operation: nil)
76+
.where("#{effects}.status = ?", "processing")
77+
.where("#{owners}.id IS NULL OR #{heartbeat} <= ? - CASE WHEN #{bindings}.recovery_timeout > ? THEN #{bindings}.recovery_timeout ELSE ? END", now, threshold, threshold)
78+
.order(:effect_id).limit(SolidObjects.configuration.claim_scan_limit)
79+
end
80+
6381
# @rbs (instance: Instance, intent: Actor::EffectRecoveryIntent, recovery: EffectRecovery, effect: Effect?, owners: Hash[String, Process], now: Time) -> void
6482
def check_one(instance:, intent:, recovery:, effect:, owners:, now:)
6583
key = "effect:#{intent.effect_id}:check:#{intent.request_id}"

‎sig/generated/lib/solid_objects/actor.rbs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -277,5 +277,8 @@ module SolidObjects
277277

278278
# @rbs (Symbol | String?) -> void
279279
def validate_effect_callback!: (Symbol | String?) -> void
280+
281+
# @rbs (timeout: Numeric?, operation: String | Symbol?) -> void
282+
def validate_recovery_timeout!: (timeout: Numeric?, operation: String | Symbol?) -> void
280283
end
281284
end

‎sig/generated/lib/solid_objects/effect_recovery_coordinator.rbs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,9 @@ module SolidObjects
2626

2727
private
2828

29+
# @rbs () -> ActiveRecord::Relation[EffectRecovery]
30+
def recovery_candidates: () -> ActiveRecord::Relation[EffectRecovery]
31+
2932
# @rbs (instance: Instance, intent: Actor::EffectRecoveryIntent, recovery: EffectRecovery, effect: Effect?, owners: Hash[String, Process], now: Time) -> void
3033
def check_one: (instance: Instance, intent: Actor::EffectRecoveryIntent, recovery: EffectRecovery, effect: Effect?, owners: Hash[String, Process], now: Time) -> void
3134

‎test/integration/effect_recovery_test.rb‎

Lines changed: 221 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -286,6 +286,57 @@ def checked(effect_id:, outcome:, arguments: nil, result: nil)
286286
worker&.stop
287287
end
288288

289+
test "one remaining mailbox slot cannot partially commit retirement" do
290+
worker, effect_executor, effect = processing_export("full-mailbox")
291+
original_limit = SolidObjects.configuration.max_mailbox_length
292+
SolidObjects.configuration.max_mailbox_length = 1
293+
assert_raises(SolidObjects::MailboxFull) do
294+
SolidObjects.database_adapter.transaction do
295+
instance = SolidObjects::Instance.lock.find(effect.instance_id)
296+
SolidObjects::EffectRecoveryCoordinator.new.check(instance:, intents: [ SolidObjects::Actor::EffectRecoveryIntent.new(effect_id: effect.effect_id, request_id: "full") ])
297+
end
298+
end
299+
assert_equal "processing", effect.reload.status
300+
assert_nil SolidObjects::EffectRecovery.find(effect.effect_id).retired_at
301+
assert_empty SolidObjects::Message.where(operation: [ "retired", "checked" ])
302+
ensure
303+
SolidObjects.configuration.max_mailbox_length = original_limit if original_limit
304+
effect_executor&.stop
305+
worker&.stop
306+
end
307+
308+
test "runtime floor changes and missing owners are rechecked for existing effects" do
309+
worker, effect_executor, effect = processing_export("runtime-floor")
310+
SolidObjects::EffectRecovery.find(effect.effect_id).update!(recovery_timeout: 1)
311+
SolidObjects::Process.find(effect.claimed_by).update!(last_heartbeat_at: SolidObjects.database_adapter.database_clock_now - 70)
312+
SolidObjects.configuration.process_alive_threshold = 120
313+
SolidObjects::EffectRecoveryCoordinator.new.recover_available
314+
assert_equal "processing", effect.reload.status
315+
assert_empty SolidObjects::Message.where(operation: "retired")
316+
effect.update!(claimed_by: nil)
317+
SolidObjects::EffectRecoveryCoordinator.new.recover_available
318+
assert_equal 1, SolidObjects::Message.where(operation: "retired").count
319+
ensure
320+
effect_executor&.stop
321+
worker&.stop
322+
end
323+
324+
test "database time lookup failure rolls back without an abandonment outcome" do
325+
worker, effect_executor, effect = processing_export("lookup-error")
326+
adapter = SolidObjects.database_adapter
327+
adapter.define_singleton_method(:database_clock_now) do
328+
SolidObjects::Record.connection.select_value("SELECT absent_recovery_column FROM #{SolidObjects.table_name(:processes)}")
329+
end
330+
assert_raises(ActiveRecord::StatementInvalid) { SolidObjects::EffectRecoveryCoordinator.new.recover_available }
331+
assert_equal "processing", effect.reload.status
332+
assert_nil SolidObjects::EffectRecovery.find(effect.effect_id).retired_at
333+
assert_empty SolidObjects::Message.where(operation: [ "retired", "checked" ])
334+
ensure
335+
adapter&.singleton_class&.remove_method(:database_clock_now)
336+
effect_executor&.stop
337+
worker&.stop
338+
end
339+
289340
test "a changed claimant is rechecked after the effect lock becomes available" do
290341
skip "PostgreSQL lock observation" unless database_family == :postgresql
291342

@@ -310,6 +361,176 @@ def checked(effect_id:, outcome:, arguments: nil, result: nil)
310361
worker&.stop
311362
end
312363

364+
test "pending work does not wait on a recovery candidate within its extended grace" do
365+
skip "PostgreSQL independent claims" unless database_family == :postgresql
366+
367+
worker, effect_executor, effect = processing_export("fresh-candidate")
368+
SolidObjects::EffectRecovery.find(effect.effect_id).update!(recovery_timeout: 120)
369+
SolidObjects::Process.find(effect.claimed_by).update!(last_heartbeat_at: SolidObjects.database_adapter.database_clock_now - 75)
370+
ExportActor.ref("other").async.start_export
371+
worker.run_until_idle
372+
claimant = SolidObjects::EffectExecutor.new
373+
results = Queue.new
374+
claim_thread = nil
375+
SolidObjects.database_adapter.transaction do
376+
SolidObjects::Instance.lock.find(effect.instance_id)
377+
claim_thread = Thread.new do
378+
SolidObjects::Record.connection_pool.with_connection do
379+
results << claimant.send(:claim_next)
380+
end
381+
end
382+
selected = Timeout.timeout(2) { results.pop }
383+
assert_equal "other", selected.instance.actor_id
384+
end
385+
ensure
386+
claim_thread&.join(5)
387+
claimant&.stop
388+
effect_executor&.stop
389+
worker&.stop
390+
end
391+
392+
test "recovery notification survives a process crash after retirement" do
393+
worker, effect_executor, effect = processing_export("crash")
394+
worker.stop
395+
configuration = SolidObjects::Record.connection_db_config.configuration_hash
396+
script = <<~RUBY
397+
require "solid_objects"
398+
ActiveRecord::Base.establish_connection(JSON.parse(ENV.fetch("RECOVERY_DATABASE_CONFIGURATION")))
399+
%w[record process instance message ready_message claimed_message reminder effect effect_recovery broadcast dead_letter].each do |model|
400+
require File.expand_path("app/models/solid_objects/\#{model}")
401+
end
402+
class RecoveryExport < SolidObjects::Actor
403+
actor_type "effect-recovery-export"
404+
def retired(effect_id:, arguments:, outcome:)
405+
end
406+
end
407+
SolidObjects::EffectRecoveryCoordinator.new.recover_available
408+
::Process.kill("KILL", ::Process.pid)
409+
RUBY
410+
process_id = ::Process.spawn({ "RECOVERY_DATABASE_CONFIGURATION" => JSON.generate(configuration) }, Gem.ruby, "-Ilib", "-e", script)
411+
_, status = ::Process.wait2(process_id)
412+
assert status.signaled?
413+
assert_equal Signal.list.fetch("KILL"), status.termsig
414+
assert_equal 1, SolidObjects::Message.where(operation: "retired").count
415+
assert_empty effect.instance.reload.state.fetch("notifications")
416+
worker = SolidObjects::Worker.new
417+
worker.run_until_idle
418+
assert_equal 1, effect.instance.reload.state.fetch("notifications").length
419+
assert_nil effect_executor.send(:claim_next)
420+
ensure
421+
effect_executor&.stop
422+
worker&.stop
423+
end
424+
425+
test "a pending claimant wins before the explicit check and is rechecked" do
426+
skip "PostgreSQL independent claims" unless database_family == :postgresql
427+
428+
ExportActor.ref("claim-first").async.start_checked_export
429+
worker = SolidObjects::Worker.new
430+
worker.run_until_idle
431+
effect = SolidObjects::Effect.find_by!(name: "build_report")
432+
claimant = SolidObjects::EffectExecutor.new
433+
adapter = SolidObjects.database_adapter
434+
original_lock = adapter.method(:lock_candidates)
435+
locked = Queue.new
436+
release = Queue.new
437+
results = Queue.new
438+
origin_locked = Queue.new
439+
adapter.define_singleton_method(:lock_candidates) do |scope|
440+
relation = original_lock.call(scope).load
441+
locked << true
442+
release.pop
443+
relation
444+
end
445+
claim_thread = Thread.new do
446+
SolidObjects::Record.connection_pool.with_connection { results << claimant.send(:claim_next) }
447+
end
448+
Timeout.timeout(5) { locked.pop }
449+
check_thread = Thread.new do
450+
SolidObjects::Record.connection_pool.with_connection do
451+
SolidObjects.database_adapter.transaction do
452+
instance = SolidObjects::Instance.lock.find(effect.instance_id)
453+
origin_locked << true
454+
SolidObjects::EffectRecoveryCoordinator.new.check(instance:, intents: [ SolidObjects::Actor::EffectRecoveryIntent.new(effect_id: effect.effect_id, request_id: "claim-first") ])
455+
end
456+
end
457+
end
458+
Timeout.timeout(5) { origin_locked.pop }
459+
release << true
460+
assert_equal effect.id, Timeout.timeout(5) { results.pop }.id
461+
check_thread.join(5)
462+
assert_equal "deferred", SolidObjects::Message.find_by!(operation: "checked").arguments.fetch("outcome")
463+
assert_empty SolidObjects::Message.where(operation: "retired")
464+
ensure
465+
release&.push(true)
466+
claim_thread&.join(5)
467+
check_thread&.join(5)
468+
adapter&.singleton_class&.remove_method(:lock_candidates)
469+
claimant&.stop
470+
worker&.stop
471+
end
472+
473+
test "an explicit check can win and leave pending work for the claimant" do
474+
skip "PostgreSQL independent claims" unless database_family == :postgresql
475+
476+
ExportActor.ref("check-first").async.start_checked_export
477+
worker = SolidObjects::Worker.new
478+
worker.run_until_idle
479+
effect = SolidObjects::Effect.find_by!(name: "build_report")
480+
claimant = SolidObjects::EffectExecutor.new
481+
results = Queue.new
482+
claim_thread = nil
483+
SolidObjects.database_adapter.transaction do
484+
instance = SolidObjects::Instance.lock.find(effect.instance_id)
485+
SolidObjects::EffectRecoveryCoordinator.new.check(instance:, intents: [ SolidObjects::Actor::EffectRecoveryIntent.new(effect_id: effect.effect_id, request_id: "check-first") ])
486+
claim_thread = Thread.new do
487+
SolidObjects::Record.connection_pool.with_connection { results << claimant.send(:claim_next) }
488+
end
489+
assert_nil Timeout.timeout(5) { results.pop }
490+
end
491+
assert_equal effect.id, claimant.send(:claim_next).id
492+
assert_equal "pending", SolidObjects::Message.find_by!(operation: "checked").arguments.fetch("outcome")
493+
assert_empty SolidObjects::Message.where(operation: "retired")
494+
ensure
495+
claim_thread&.join(5)
496+
claimant&.stop
497+
worker&.stop
498+
end
499+
500+
test "concurrent explicit checks and a replay share one retirement" do
501+
skip "PostgreSQL independent checks" unless database_family == :postgresql
502+
503+
worker, effect_executor, effect = processing_export("explicit-race")
504+
process_ids = Queue.new
505+
check = ->(request_id) do
506+
SolidObjects.database_adapter.transaction do
507+
instance = SolidObjects::Instance.lock.find(effect.instance_id)
508+
SolidObjects::EffectRecoveryCoordinator.new.check(instance:, intents: [ SolidObjects::Actor::EffectRecoveryIntent.new(effect_id: effect.effect_id, request_id:) ])
509+
end
510+
end
511+
threads = []
512+
SolidObjects.database_adapter.transaction do
513+
SolidObjects::Instance.lock.find(effect.instance_id)
514+
%w[one two].each do |request_id|
515+
threads << Thread.new do
516+
SolidObjects::Record.connection_pool.with_connection do |connection|
517+
process_ids << connection.select_value("SELECT pg_backend_pid()").to_i
518+
check.call(request_id)
519+
end
520+
end
521+
end
522+
2.times { wait_for_blocked_process(Timeout.timeout(5) { process_ids.pop }) }
523+
end
524+
threads.each { |thread| thread.join(5) }
525+
check.call("one")
526+
assert_equal 1, SolidObjects::Message.where(operation: "retired").count
527+
assert_equal %w[retired already_retired], SolidObjects::Message.where(operation: "checked").order(:sequence).map { |message| message.arguments.fetch("outcome") }
528+
ensure
529+
threads&.each { |thread| thread.join(5) }
530+
effect_executor&.stop
531+
worker&.stop
532+
end
533+
313534
private
314535

315536
def processing_export(actor_id)

0 commit comments

Comments
 (0)