From 275f213cd705076994422d0a24b2134725f3aff9 Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 10:15:42 -0700 Subject: [PATCH 1/5] fix: lock the actor instance by primary key An enqueue used create_or_find_by!, which inserts first. For an actor that already exists, every enqueue paid an insert that failed on the unique key, a transaction restart, and two locking reads. Find the row first, then lock it by its primary key. A steady-state enqueue now issues 10 statements instead of 12, and holds the row for one statement less. This also removes a deadlock. A failed insert leaves a shared lock on the identity index, and MySQL keeps that lock across a savepoint rollback. The old code escaped this only when the insert was the first statement of the transaction, because Active Record then restarts the transaction instead of the savepoint. Callers that already wrote, such as the executor and the reminder scheduler, got a real savepoint and deadlocked when they created the same actor at the same time. Four of eight concurrent callers failed that way before this change. The mailbox now reads the winning row in shared mode after a duplicate key, and locks it by primary key. MySQL needs the shared read because it defaults to repeatable read, and a consistent read cannot see the winning row. PostgreSQL and SQLite take no shared lock, because a share lock there creates the same upgrade deadlock it prevents on MySQL. Validate with bundle exec rake test on SQLite, PostgreSQL, mysql2, and Trilogy, and with bundle exec rake standard rubocop rbs steep security. --- CHANGELOG.md | 10 +++ docs/roadmap.md | 7 +- lib/solid_objects/database_adapter.rb | 10 +++ lib/solid_objects/database_adapters/mysql.rb | 5 ++ lib/solid_objects/mailbox.rb | 46 +++++++++-- .../lib/solid_objects/database_adapter.rbs | 6 ++ .../solid_objects/database_adapters/mysql.rbs | 3 + sig/generated/lib/solid_objects/mailbox.rbs | 12 +++ .../enqueue_statement_count_test.rb | 78 +++++++++++++++++++ test/integration/enqueue_test.rb | 75 ++++++++++++++++++ 10 files changed, 244 insertions(+), 8 deletions(-) create mode 100644 test/integration/enqueue_statement_count_test.rb diff --git a/CHANGELOG.md b/CHANGELOG.md index 7216aa8..9b67088 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,15 @@ # Changelog +## Unreleased + +- Find the actor instance before the insert when an enqueue starts, and lock + that row by its primary key. A steady-state enqueue now writes no instance + row and issues 10 statements instead of 12. +- Stop the deadlock between concurrent enqueues that create the same actor + inside a transaction that already wrote. MySQL keeps the shared lock of a + failed insert across a savepoint rollback, so the mailbox reads the winning + row in shared mode and never asks to upgrade that lock. + ## 0.15.1 - 2026-09-16 - Use the existing cleanup index when finding expired actor instances. Preserve diff --git a/docs/roadmap.md b/docs/roadmap.md index 0261cfa..cb42e67 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -6,7 +6,12 @@ - Explicit actor registry, references, JSON state, and state migrations - Fluent direct synchronous RPC, configured `sync`, and durable `async` - Durable message history plus ready/claimed membership tables -- Concurrent sequence allocation and actor creation +- Concurrent sequence allocation and actor creation. An enqueue finds the + instance row with an unlocked read, then locks that row by its primary key. + A steady-state enqueue writes no instance row, and issues 10 statements + instead of 12. Concurrent creation causes no deadlock on SQLite, PostgreSQL, + or MySQL. MySQL needs a shared read after a duplicate key, because it uses + repeatable read. The tests count statements, and do not measure latency. - Activation leases, renewal, unique activation tokens, generations, and fenced commits - Bounded activation passes, idle cache, hot-actor yield, and process records diff --git a/lib/solid_objects/database_adapter.rb b/lib/solid_objects/database_adapter.rb index 6ea5a3a..746852d 100644 --- a/lib/solid_objects/database_adapter.rb +++ b/lib/solid_objects/database_adapter.rb @@ -99,6 +99,11 @@ def claim_lock nil end + # @rbs () -> String? + def shared_lock + nil + end + # @rbs () -> String def current_time_expression "CURRENT_TIMESTAMP" @@ -163,6 +168,11 @@ def lock_candidates(relation) claim_lock ? relation.lock(claim_lock) : relation end + # @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped] + def share_locked(relation) + shared_lock ? relation.lock(shared_lock) : relation + end + private attr_reader :connection_pool, :fixed_connection diff --git a/lib/solid_objects/database_adapters/mysql.rb b/lib/solid_objects/database_adapters/mysql.rb index d91c0ec..8d6b929 100644 --- a/lib/solid_objects/database_adapters/mysql.rb +++ b/lib/solid_objects/database_adapters/mysql.rb @@ -15,6 +15,11 @@ def claim_lock "FOR UPDATE SKIP LOCKED" end + # @rbs () -> String + def shared_lock + "FOR SHARE" + end + # A non-transactional engine would silently break fenced commits, so the # storage engine is verified rather than assumed. # @rbs () -> Array[String] diff --git a/lib/solid_objects/mailbox.rb b/lib/solid_objects/mailbox.rb index b1db68c..6c630bb 100644 --- a/lib/solid_objects/mailbox.rb +++ b/lib/solid_objects/mailbox.rb @@ -52,7 +52,6 @@ def enqueue_in_transaction( max_bytes: SolidObjects.configuration.max_payload_bytes ) instance = find_or_create_instance(reference, actor_class) - instance.lock! existing = find_idempotent_message(instance, idempotency_key) if existing @@ -108,13 +107,46 @@ def with_instance_retry # @rbs (Reference, Class) -> Instance def find_or_create_instance(reference, actor_class) - Instance.create_or_find_by!( - actor_type: reference.actor_type, - actor_id: reference.actor_id - ) do |instance| - instance.state = {} - instance.state_version = actor_class.state_version + identifier = instance_identifier(reference) + return lock_instance!(identifier) if identifier + + create_locked_instance(reference, actor_class) + end + + # @rbs (Reference) -> Integer? + def instance_identifier(reference) + Instance + .where(actor_type: reference.actor_type, actor_id: reference.actor_id) + .pick(:id) + end + + # @rbs (Reference, Class) -> Instance + def create_locked_instance(reference, actor_class) + Instance.transaction(requires_new: true) do + Instance.create!( + actor_type: reference.actor_type, + actor_id: reference.actor_id, + state: {}, + state_version: actor_class.state_version + ) end + rescue ActiveRecord::RecordNotUnique + lock_instance!(committed_instance_identifier(reference)) + end + + # @rbs (Reference) -> Integer? + def committed_instance_identifier(reference) + database_adapter.share_locked( + Instance.where(actor_type: reference.actor_type, actor_id: reference.actor_id) + ).pick(:id) + end + + # @rbs (Integer?) -> Instance + def lock_instance!(identifier) + instance = identifier && Instance.lock.find_by(id: identifier) + return instance if instance + + raise ActiveRecord::RecordNotFound, "actor instance disappeared while enqueueing" end # @rbs (Instance, String?) -> Message? diff --git a/sig/generated/lib/solid_objects/database_adapter.rbs b/sig/generated/lib/solid_objects/database_adapter.rbs index f429e0b..6e7ba54 100644 --- a/sig/generated/lib/solid_objects/database_adapter.rbs +++ b/sig/generated/lib/solid_objects/database_adapter.rbs @@ -49,6 +49,9 @@ module SolidObjects # @rbs () -> String? def claim_lock: () -> String? + # @rbs () -> String? + def shared_lock: () -> String? + # @rbs () -> String def current_time_expression: () -> String @@ -70,6 +73,9 @@ module SolidObjects # @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped] def lock_candidates: (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped] + # @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped] + def share_locked: (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped] + private attr_reader connection_pool: untyped diff --git a/sig/generated/lib/solid_objects/database_adapters/mysql.rbs b/sig/generated/lib/solid_objects/database_adapters/mysql.rbs index 6ad899e..e53b81f 100644 --- a/sig/generated/lib/solid_objects/database_adapters/mysql.rbs +++ b/sig/generated/lib/solid_objects/database_adapters/mysql.rbs @@ -11,6 +11,9 @@ module SolidObjects # @rbs () -> String def claim_lock: () -> String + # @rbs () -> String + def shared_lock: () -> String + # A non-transactional engine would silently break fenced commits, so the # storage engine is verified rather than assumed. # @rbs () -> Array[String] diff --git a/sig/generated/lib/solid_objects/mailbox.rbs b/sig/generated/lib/solid_objects/mailbox.rbs index 7f84e24..bbfe3ac 100644 --- a/sig/generated/lib/solid_objects/mailbox.rbs +++ b/sig/generated/lib/solid_objects/mailbox.rbs @@ -28,6 +28,18 @@ module SolidObjects # @rbs (Reference, Class) -> Instance def find_or_create_instance: (Reference, Class) -> Instance + # @rbs (Reference) -> Integer? + def instance_identifier: (Reference) -> Integer? + + # @rbs (Reference, Class) -> Instance + def create_locked_instance: (Reference, Class) -> Instance + + # @rbs (Reference) -> Integer? + def committed_instance_identifier: (Reference) -> Integer? + + # @rbs (Integer?) -> Instance + def lock_instance!: (Integer?) -> Instance + # @rbs (Instance, String?) -> Message? def find_idempotent_message: (Instance, String?) -> Message? diff --git a/test/integration/enqueue_statement_count_test.rb b/test/integration/enqueue_statement_count_test.rb new file mode 100644 index 0000000..2059031 --- /dev/null +++ b/test/integration/enqueue_statement_count_test.rb @@ -0,0 +1,78 @@ +# frozen_string_literal: true + +require "database_test_helper" + +class EnqueueStatementCountTest < ActiveSupport::TestCase + class CartActor < SolidObjects::Actor + actor_type "enqueue-statement-count-cart" + + attribute :items, default: -> { [] } + + def add(product_id:) + self.items += [ product_id ] + end + end + + STEADY_STATE_STATEMENT_COUNT = 10 + + setup { CartActor.ensure_registered! } + + test "a steady-state enqueue never inserts the instance row" do + reference = CartActor.ref("alice") + reference.async.add(product_id: "shirt") + + statements = capture_statements { reference.async.add(product_id: "pants") } + + assert_empty instance_statements(statements).grep(/\AINSERT/i) + end + + test "a steady-state enqueue touches the instance row three times" do + reference = CartActor.ref("alice") + reference.async.add(product_id: "shirt") + + statements = instance_statements(capture_statements { reference.async.add(product_id: "pants") }) + + assert_equal 2, statements.grep(/\ASELECT/i).length, statements.inspect + assert_equal 1, statements.grep(/\AUPDATE/i).length, statements.inspect + assert_equal 3, statements.length, statements.inspect + end + + test "a steady-state enqueue opens one transaction and never restarts it" do + reference = CartActor.ref("alice") + reference.async.add(product_id: "shirt") + + statements = capture_statements { reference.async.add(product_id: "pants") } + control = statements.grep(/\A(?:BEGIN|COMMIT|ROLLBACK|SAVEPOINT|RELEASE)/i) + + assert_empty control.grep(/ROLLBACK|SAVEPOINT/i), statements.inspect + assert_equal 2, control.length, statements.inspect + end + + test "a steady-state enqueue issues a fixed number of statements" do + reference = CartActor.ref("alice") + reference.async.add(product_id: "shirt") + + statements = capture_statements { reference.async.add(product_id: "pants") } + + assert_equal STEADY_STATE_STATEMENT_COUNT, statements.length, statements.inspect + end + + private + + def instance_statements(statements) + statements.grep(/solid_objects_instances/) + end + + def capture_statements + statements = [] + subscriber = lambda do |*arguments| + payload = arguments.last + next if payload[:name] == "SCHEMA" + next if payload[:cached] + + statements << payload.fetch(:sql).to_s.strip + end + ActiveSupport::Notifications.subscribed(subscriber, "sql.active_record") { yield } + statements + end +end diff --git a/test/integration/enqueue_test.rb b/test/integration/enqueue_test.rb index 12ef508..17c3e3a 100644 --- a/test/integration/enqueue_test.rb +++ b/test/integration/enqueue_test.rb @@ -58,6 +58,81 @@ def add(product_id:) assert_empty errors.size.times.map { errors.pop } assert_equal (1..8).to_a, results.size.times.map { results.pop }.sort + assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count + end + + test "allocates unique sequences under concurrent enqueue to an existing actor" do + reference = CartActor.ref("alice") + reference.async.add(product_id: "first") + start = Queue.new + results = Queue.new + errors = Queue.new + + threads = 8.times.map do |index| + Thread.new do + SolidObjects::Record.connection_pool.with_connection do + start.pop + results << reference.async.add(product_id: "product-#{index}").sequence + rescue => error + errors << error + end + end + end + + threads.length.times { start << true } + threads.each(&:join) + + assert_empty errors.size.times.map { errors.pop } + assert_equal (2..9).to_a, results.size.times.map { results.pop }.sort + assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count + end + + test "creates the instance once when concurrent callers already hold a dirty transaction" do + CartActor.ensure_registered! + reference = SolidObjects::Reference.new(actor_type: "enqueue-carts", actor_id: "alice") + mailbox = SolidObjects::Mailbox.new + start = Queue.new + sequences = Queue.new + errors = Queue.new + + threads = 8.times.map do |index| + Thread.new do + SolidObjects::Record.connection_pool.with_connection do + start.pop + SolidObjects.database_adapter.transaction do + SolidObjectsTestDomainRecord.create!(name: "dirty-#{index}") + sequences << mailbox.enqueue_in_transaction( + reference:, + operation: :add, + arguments: { product_id: "product-#{index}" }, + delivery_mode: "async", + idempotency_key: nil + ).sequence + end + rescue => error + errors << error + end + end + end + + threads.length.times { start << true } + threads.each { |thread| thread.join(30) } + + assert_empty errors.size.times.map { errors.pop } + assert_equal (1..8).to_a, sequences.size.times.map { sequences.pop }.sort + assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count + end + + test "gives up when the instance keeps disappearing between the lookup and the insert" do + SolidObjects::Instance.singleton_class.define_method(:create!) do |*, **| + raise ActiveRecord::RecordNotUnique, "simulated create race" + end + + assert_raises(SolidObjects::ActorDestroyed) do + CartActor.ref("ghost").async.add(product_id: "shirt") + end + ensure + SolidObjects::Instance.singleton_class.send(:remove_method, :create!) end test "deduplicates the same idempotent enqueue" do From 8ceb2431ab6f9b19c8d531ba02752f2c490c9bfc Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 10:22:04 -0700 Subject: [PATCH 2/5] chore: prepare version 0.15.2 --- CHANGELOG.md | 2 +- Gemfile.lock | 4 ++-- lib/solid_objects/version.rb | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9b67088..c87adbe 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,6 @@ # Changelog -## Unreleased +## 0.15.2 - 2026-09-21 - Find the actor instance before the insert when an enqueue starts, and lock that row by its primary key. A steady-state enqueue now writes no instance diff --git a/Gemfile.lock b/Gemfile.lock index e4f588d..0c0f2e8 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -1,7 +1,7 @@ PATH remote: . specs: - solid_objects (0.15.1) + solid_objects (0.15.2) actioncable (>= 7.1) actionpack (>= 7.1) actionview (>= 7.1) @@ -384,7 +384,7 @@ CHECKSUMS rubocop-rails-omakase (1.1.0) sha256=2af73ac8ee5852de2919abbd2618af9c15c19b512c4cfc1f9a5d3b6ef009109d ruby-progressbar (1.13.0) sha256=80fc9c47a9b640d6834e0dc7b3c94c9df37f08cb072b7761e4a71e22cff29b33 securerandom (0.4.1) sha256=cc5193d414a4341b6e225f0cb4446aceca8e50d5e1888743fac16987638ea0b1 - solid_objects (0.15.1) + solid_objects (0.15.2) sqlite3 (2.9.5-aarch64-linux-gnu) sha256=78075b6337d3d182c6d2b4691049ed45cd220826160c9ea18946bf6a1de200dc sqlite3 (2.9.5-aarch64-linux-musl) sha256=18c801185deb4adc01ddb281e8f672a39e3d1729979ca91e39439cd3eac0402d sqlite3 (2.9.5-arm-linux-gnu) sha256=1bdfca0c7d63998c60b0f4a8e3c8df2d33800ccc4abd2d612eddbbbc92a4c48b diff --git a/lib/solid_objects/version.rb b/lib/solid_objects/version.rb index fad4efd..aaaac93 100644 --- a/lib/solid_objects/version.rb +++ b/lib/solid_objects/version.rb @@ -1,5 +1,5 @@ # rbs_inline: enabled module SolidObjects - VERSION = "0.15.1" + VERSION = "0.15.2" end From cfe6042be959255d6352d55f2f8508bb65118753 Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 10:35:06 -0700 Subject: [PATCH 3/5] refactor: inline the instance lookup queries Put each lookup at its only call site, and drop the two helpers that forwarded a single query. Assert that every worker in the create race test finishes, and kill any worker that outlives its timeout, so a worker cannot hold a pooled connection while the assertions run. --- lib/solid_objects/mailbox.rb | 19 +++++-------------- sig/generated/lib/solid_objects/mailbox.rbs | 6 ------ test/integration/enqueue_test.rb | 4 +++- 3 files changed, 8 insertions(+), 21 deletions(-) diff --git a/lib/solid_objects/mailbox.rb b/lib/solid_objects/mailbox.rb index 6c630bb..0a8e671 100644 --- a/lib/solid_objects/mailbox.rb +++ b/lib/solid_objects/mailbox.rb @@ -107,19 +107,14 @@ def with_instance_retry # @rbs (Reference, Class) -> Instance def find_or_create_instance(reference, actor_class) - identifier = instance_identifier(reference) + identifier = Instance + .where(actor_type: reference.actor_type, actor_id: reference.actor_id) + .pick(:id) return lock_instance!(identifier) if identifier create_locked_instance(reference, actor_class) end - # @rbs (Reference) -> Integer? - def instance_identifier(reference) - Instance - .where(actor_type: reference.actor_type, actor_id: reference.actor_id) - .pick(:id) - end - # @rbs (Reference, Class) -> Instance def create_locked_instance(reference, actor_class) Instance.transaction(requires_new: true) do @@ -131,14 +126,10 @@ def create_locked_instance(reference, actor_class) ) end rescue ActiveRecord::RecordNotUnique - lock_instance!(committed_instance_identifier(reference)) - end - - # @rbs (Reference) -> Integer? - def committed_instance_identifier(reference) - database_adapter.share_locked( + identifier = database_adapter.share_locked( Instance.where(actor_type: reference.actor_type, actor_id: reference.actor_id) ).pick(:id) + lock_instance!(identifier) end # @rbs (Integer?) -> Instance diff --git a/sig/generated/lib/solid_objects/mailbox.rbs b/sig/generated/lib/solid_objects/mailbox.rbs index bbfe3ac..88e2a91 100644 --- a/sig/generated/lib/solid_objects/mailbox.rbs +++ b/sig/generated/lib/solid_objects/mailbox.rbs @@ -28,15 +28,9 @@ module SolidObjects # @rbs (Reference, Class) -> Instance def find_or_create_instance: (Reference, Class) -> Instance - # @rbs (Reference) -> Integer? - def instance_identifier: (Reference) -> Integer? - # @rbs (Reference, Class) -> Instance def create_locked_instance: (Reference, Class) -> Instance - # @rbs (Reference) -> Integer? - def committed_instance_identifier: (Reference) -> Integer? - # @rbs (Integer?) -> Instance def lock_instance!: (Integer?) -> Instance diff --git a/test/integration/enqueue_test.rb b/test/integration/enqueue_test.rb index 17c3e3a..083c838 100644 --- a/test/integration/enqueue_test.rb +++ b/test/integration/enqueue_test.rb @@ -116,8 +116,10 @@ def add(product_id:) end threads.length.times { start << true } - threads.each { |thread| thread.join(30) } + unfinished = threads.reject { |thread| thread.join(30) } + unfinished.each(&:kill).each(&:join) + assert_empty unfinished, "an enqueue was still running after its timeout" assert_empty errors.size.times.map { errors.pop } assert_equal (1..8).to_a, sequences.size.times.map { sequences.pop }.sort assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count From 38d30b19ebc46e8dca370f6dfe9b840ca97cc1a3 Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 11:40:35 -0700 Subject: [PATCH 4/5] refactor: inline the create race into the lookup Put the savepoint and the duplicate-key recovery in the one method that uses them. Bind the identity attributes once, so the same pair no longer repeats across three queries. --- lib/solid_objects/mailbox.rb | 22 ++++----------------- sig/generated/lib/solid_objects/mailbox.rbs | 3 --- 2 files changed, 4 insertions(+), 21 deletions(-) diff --git a/lib/solid_objects/mailbox.rb b/lib/solid_objects/mailbox.rb index 0a8e671..65c2167 100644 --- a/lib/solid_objects/mailbox.rb +++ b/lib/solid_objects/mailbox.rb @@ -107,29 +107,15 @@ def with_instance_retry # @rbs (Reference, Class) -> Instance def find_or_create_instance(reference, actor_class) - identifier = Instance - .where(actor_type: reference.actor_type, actor_id: reference.actor_id) - .pick(:id) + identity = { actor_type: reference.actor_type, actor_id: reference.actor_id } + identifier = Instance.where(identity).pick(:id) return lock_instance!(identifier) if identifier - create_locked_instance(reference, actor_class) - end - - # @rbs (Reference, Class) -> Instance - def create_locked_instance(reference, actor_class) Instance.transaction(requires_new: true) do - Instance.create!( - actor_type: reference.actor_type, - actor_id: reference.actor_id, - state: {}, - state_version: actor_class.state_version - ) + Instance.create!(**identity, state: {}, state_version: actor_class.state_version) end rescue ActiveRecord::RecordNotUnique - identifier = database_adapter.share_locked( - Instance.where(actor_type: reference.actor_type, actor_id: reference.actor_id) - ).pick(:id) - lock_instance!(identifier) + lock_instance!(database_adapter.share_locked(Instance.where(identity)).pick(:id)) end # @rbs (Integer?) -> Instance diff --git a/sig/generated/lib/solid_objects/mailbox.rbs b/sig/generated/lib/solid_objects/mailbox.rbs index 88e2a91..0de38fb 100644 --- a/sig/generated/lib/solid_objects/mailbox.rbs +++ b/sig/generated/lib/solid_objects/mailbox.rbs @@ -28,9 +28,6 @@ module SolidObjects # @rbs (Reference, Class) -> Instance def find_or_create_instance: (Reference, Class) -> Instance - # @rbs (Reference, Class) -> Instance - def create_locked_instance: (Reference, Class) -> Instance - # @rbs (Integer?) -> Instance def lock_instance!: (Integer?) -> Instance From 9cbd293b8eecb95b5aa860ea86e95de576f6f10c Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 11:53:44 -0700 Subject: [PATCH 5/5] fix: raise the worker hang budget in the create race test The create race test capped each worker at 30 seconds. SQLite retries a busy write up to lock_retry_attempts times, and every retry waits out the 5 second busy handler, so a starved worker can run for about a minute before it does any real work. The cap sat below that legitimate worst case, and CI failed on Rails 7.1 and 7.2 while the same jobs passed in a parallel run of the same commit. Raise the budget to 180 seconds, so it reports a worker that hangs rather than one that waits. Confirmed by forcing a worker to sleep past a shortened budget and watching the assertion fail. --- test/integration/enqueue_test.rb | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/test/integration/enqueue_test.rb b/test/integration/enqueue_test.rb index 083c838..543df0a 100644 --- a/test/integration/enqueue_test.rb +++ b/test/integration/enqueue_test.rb @@ -3,6 +3,8 @@ require "database_test_helper" class EnqueueTest < ActiveSupport::TestCase + WORKER_HANG_TIMEOUT = 180 + class CartActor < SolidObjects::Actor actor_type "enqueue-carts" @@ -116,10 +118,10 @@ def add(product_id:) end threads.length.times { start << true } - unfinished = threads.reject { |thread| thread.join(30) } + unfinished = threads.reject { |thread| thread.join(WORKER_HANG_TIMEOUT) } unfinished.each(&:kill).each(&:join) - assert_empty unfinished, "an enqueue was still running after its timeout" + assert_empty unfinished, "an enqueue hung for over #{WORKER_HANG_TIMEOUT} seconds" assert_empty errors.size.times.map { errors.pop } assert_equal (1..8).to_a, sequences.size.times.map { sequences.pop }.sort assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count