Skip to content

Commit 275f213

Browse files
committed
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.
1 parent 99c41e1 commit 275f213

10 files changed

Lines changed: 244 additions & 8 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,15 @@
11
# Changelog
22

3+
## Unreleased
4+
5+
- Find the actor instance before the insert when an enqueue starts, and lock
6+
that row by its primary key. A steady-state enqueue now writes no instance
7+
row and issues 10 statements instead of 12.
8+
- Stop the deadlock between concurrent enqueues that create the same actor
9+
inside a transaction that already wrote. MySQL keeps the shared lock of a
10+
failed insert across a savepoint rollback, so the mailbox reads the winning
11+
row in shared mode and never asks to upgrade that lock.
12+
313
## 0.15.1 - 2026-09-16
414

515
- Use the existing cleanup index when finding expired actor instances. Preserve

‎docs/roadmap.md‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,12 @@
66
- Explicit actor registry, references, JSON state, and state migrations
77
- Fluent direct synchronous RPC, configured `sync`, and durable `async`
88
- Durable message history plus ready/claimed membership tables
9-
- Concurrent sequence allocation and actor creation
9+
- Concurrent sequence allocation and actor creation. An enqueue finds the
10+
instance row with an unlocked read, then locks that row by its primary key.
11+
A steady-state enqueue writes no instance row, and issues 10 statements
12+
instead of 12. Concurrent creation causes no deadlock on SQLite, PostgreSQL,
13+
or MySQL. MySQL needs a shared read after a duplicate key, because it uses
14+
repeatable read. The tests count statements, and do not measure latency.
1015
- Activation leases, renewal, unique activation tokens, generations, and
1116
fenced commits
1217
- Bounded activation passes, idle cache, hot-actor yield, and process records

‎lib/solid_objects/database_adapter.rb‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,11 @@ def claim_lock
9999
nil
100100
end
101101

102+
# @rbs () -> String?
103+
def shared_lock
104+
nil
105+
end
106+
102107
# @rbs () -> String
103108
def current_time_expression
104109
"CURRENT_TIMESTAMP"
@@ -163,6 +168,11 @@ def lock_candidates(relation)
163168
claim_lock ? relation.lock(claim_lock) : relation
164169
end
165170

171+
# @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped]
172+
def share_locked(relation)
173+
shared_lock ? relation.lock(shared_lock) : relation
174+
end
175+
166176
private
167177

168178
attr_reader :connection_pool, :fixed_connection

‎lib/solid_objects/database_adapters/mysql.rb‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,11 @@ def claim_lock
1515
"FOR UPDATE SKIP LOCKED"
1616
end
1717

18+
# @rbs () -> String
19+
def shared_lock
20+
"FOR SHARE"
21+
end
22+
1823
# A non-transactional engine would silently break fenced commits, so the
1924
# storage engine is verified rather than assumed.
2025
# @rbs () -> Array[String]

‎lib/solid_objects/mailbox.rb‎

Lines changed: 39 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,6 @@ def enqueue_in_transaction(
5252
max_bytes: SolidObjects.configuration.max_payload_bytes
5353
)
5454
instance = find_or_create_instance(reference, actor_class)
55-
instance.lock!
5655

5756
existing = find_idempotent_message(instance, idempotency_key)
5857
if existing
@@ -108,13 +107,46 @@ def with_instance_retry
108107

109108
# @rbs (Reference, Class) -> Instance
110109
def find_or_create_instance(reference, actor_class)
111-
Instance.create_or_find_by!(
112-
actor_type: reference.actor_type,
113-
actor_id: reference.actor_id
114-
) do |instance|
115-
instance.state = {}
116-
instance.state_version = actor_class.state_version
110+
identifier = instance_identifier(reference)
111+
return lock_instance!(identifier) if identifier
112+
113+
create_locked_instance(reference, actor_class)
114+
end
115+
116+
# @rbs (Reference) -> Integer?
117+
def instance_identifier(reference)
118+
Instance
119+
.where(actor_type: reference.actor_type, actor_id: reference.actor_id)
120+
.pick(:id)
121+
end
122+
123+
# @rbs (Reference, Class) -> Instance
124+
def create_locked_instance(reference, actor_class)
125+
Instance.transaction(requires_new: true) do
126+
Instance.create!(
127+
actor_type: reference.actor_type,
128+
actor_id: reference.actor_id,
129+
state: {},
130+
state_version: actor_class.state_version
131+
)
117132
end
133+
rescue ActiveRecord::RecordNotUnique
134+
lock_instance!(committed_instance_identifier(reference))
135+
end
136+
137+
# @rbs (Reference) -> Integer?
138+
def committed_instance_identifier(reference)
139+
database_adapter.share_locked(
140+
Instance.where(actor_type: reference.actor_type, actor_id: reference.actor_id)
141+
).pick(:id)
142+
end
143+
144+
# @rbs (Integer?) -> Instance
145+
def lock_instance!(identifier)
146+
instance = identifier && Instance.lock.find_by(id: identifier)
147+
return instance if instance
148+
149+
raise ActiveRecord::RecordNotFound, "actor instance disappeared while enqueueing"
118150
end
119151

120152
# @rbs (Instance, String?) -> Message?

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,9 @@ module SolidObjects
4949
# @rbs () -> String?
5050
def claim_lock: () -> String?
5151

52+
# @rbs () -> String?
53+
def shared_lock: () -> String?
54+
5255
# @rbs () -> String
5356
def current_time_expression: () -> String
5457

@@ -70,6 +73,9 @@ module SolidObjects
7073
# @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped]
7174
def lock_candidates: (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped]
7275

76+
# @rbs (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped]
77+
def share_locked: (ActiveRecord::Relation[untyped]) -> ActiveRecord::Relation[untyped]
78+
7379
private
7480

7581
attr_reader connection_pool: untyped

‎sig/generated/lib/solid_objects/database_adapters/mysql.rbs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,9 @@ module SolidObjects
1111
# @rbs () -> String
1212
def claim_lock: () -> String
1313

14+
# @rbs () -> String
15+
def shared_lock: () -> String
16+
1417
# A non-transactional engine would silently break fenced commits, so the
1518
# storage engine is verified rather than assumed.
1619
# @rbs () -> Array[String]

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

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,18 @@ module SolidObjects
2828
# @rbs (Reference, Class) -> Instance
2929
def find_or_create_instance: (Reference, Class) -> Instance
3030

31+
# @rbs (Reference) -> Integer?
32+
def instance_identifier: (Reference) -> Integer?
33+
34+
# @rbs (Reference, Class) -> Instance
35+
def create_locked_instance: (Reference, Class) -> Instance
36+
37+
# @rbs (Reference) -> Integer?
38+
def committed_instance_identifier: (Reference) -> Integer?
39+
40+
# @rbs (Integer?) -> Instance
41+
def lock_instance!: (Integer?) -> Instance
42+
3143
# @rbs (Instance, String?) -> Message?
3244
def find_idempotent_message: (Instance, String?) -> Message?
3345

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
# frozen_string_literal: true
2+
3+
require "database_test_helper"
4+
5+
class EnqueueStatementCountTest < ActiveSupport::TestCase
6+
class CartActor < SolidObjects::Actor
7+
actor_type "enqueue-statement-count-cart"
8+
9+
attribute :items, default: -> { [] }
10+
11+
def add(product_id:)
12+
self.items += [ product_id ]
13+
end
14+
end
15+
16+
STEADY_STATE_STATEMENT_COUNT = 10
17+
18+
setup { CartActor.ensure_registered! }
19+
20+
test "a steady-state enqueue never inserts the instance row" do
21+
reference = CartActor.ref("alice")
22+
reference.async.add(product_id: "shirt")
23+
24+
statements = capture_statements { reference.async.add(product_id: "pants") }
25+
26+
assert_empty instance_statements(statements).grep(/\AINSERT/i)
27+
end
28+
29+
test "a steady-state enqueue touches the instance row three times" do
30+
reference = CartActor.ref("alice")
31+
reference.async.add(product_id: "shirt")
32+
33+
statements = instance_statements(capture_statements { reference.async.add(product_id: "pants") })
34+
35+
assert_equal 2, statements.grep(/\ASELECT/i).length, statements.inspect
36+
assert_equal 1, statements.grep(/\AUPDATE/i).length, statements.inspect
37+
assert_equal 3, statements.length, statements.inspect
38+
end
39+
40+
test "a steady-state enqueue opens one transaction and never restarts it" do
41+
reference = CartActor.ref("alice")
42+
reference.async.add(product_id: "shirt")
43+
44+
statements = capture_statements { reference.async.add(product_id: "pants") }
45+
control = statements.grep(/\A(?:BEGIN|COMMIT|ROLLBACK|SAVEPOINT|RELEASE)/i)
46+
47+
assert_empty control.grep(/ROLLBACK|SAVEPOINT/i), statements.inspect
48+
assert_equal 2, control.length, statements.inspect
49+
end
50+
51+
test "a steady-state enqueue issues a fixed number of statements" do
52+
reference = CartActor.ref("alice")
53+
reference.async.add(product_id: "shirt")
54+
55+
statements = capture_statements { reference.async.add(product_id: "pants") }
56+
57+
assert_equal STEADY_STATE_STATEMENT_COUNT, statements.length, statements.inspect
58+
end
59+
60+
private
61+
62+
def instance_statements(statements)
63+
statements.grep(/solid_objects_instances/)
64+
end
65+
66+
def capture_statements
67+
statements = []
68+
subscriber = lambda do |*arguments|
69+
payload = arguments.last
70+
next if payload[:name] == "SCHEMA"
71+
next if payload[:cached]
72+
73+
statements << payload.fetch(:sql).to_s.strip
74+
end
75+
ActiveSupport::Notifications.subscribed(subscriber, "sql.active_record") { yield }
76+
statements
77+
end
78+
end

‎test/integration/enqueue_test.rb‎

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,81 @@ def add(product_id:)
5858

5959
assert_empty errors.size.times.map { errors.pop }
6060
assert_equal (1..8).to_a, results.size.times.map { results.pop }.sort
61+
assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count
62+
end
63+
64+
test "allocates unique sequences under concurrent enqueue to an existing actor" do
65+
reference = CartActor.ref("alice")
66+
reference.async.add(product_id: "first")
67+
start = Queue.new
68+
results = Queue.new
69+
errors = Queue.new
70+
71+
threads = 8.times.map do |index|
72+
Thread.new do
73+
SolidObjects::Record.connection_pool.with_connection do
74+
start.pop
75+
results << reference.async.add(product_id: "product-#{index}").sequence
76+
rescue => error
77+
errors << error
78+
end
79+
end
80+
end
81+
82+
threads.length.times { start << true }
83+
threads.each(&:join)
84+
85+
assert_empty errors.size.times.map { errors.pop }
86+
assert_equal (2..9).to_a, results.size.times.map { results.pop }.sort
87+
assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count
88+
end
89+
90+
test "creates the instance once when concurrent callers already hold a dirty transaction" do
91+
CartActor.ensure_registered!
92+
reference = SolidObjects::Reference.new(actor_type: "enqueue-carts", actor_id: "alice")
93+
mailbox = SolidObjects::Mailbox.new
94+
start = Queue.new
95+
sequences = Queue.new
96+
errors = Queue.new
97+
98+
threads = 8.times.map do |index|
99+
Thread.new do
100+
SolidObjects::Record.connection_pool.with_connection do
101+
start.pop
102+
SolidObjects.database_adapter.transaction do
103+
SolidObjectsTestDomainRecord.create!(name: "dirty-#{index}")
104+
sequences << mailbox.enqueue_in_transaction(
105+
reference:,
106+
operation: :add,
107+
arguments: { product_id: "product-#{index}" },
108+
delivery_mode: "async",
109+
idempotency_key: nil
110+
).sequence
111+
end
112+
rescue => error
113+
errors << error
114+
end
115+
end
116+
end
117+
118+
threads.length.times { start << true }
119+
threads.each { |thread| thread.join(30) }
120+
121+
assert_empty errors.size.times.map { errors.pop }
122+
assert_equal (1..8).to_a, sequences.size.times.map { sequences.pop }.sort
123+
assert_equal 1, SolidObjects::Instance.where(actor_type: "enqueue-carts", actor_id: "alice").count
124+
end
125+
126+
test "gives up when the instance keeps disappearing between the lookup and the insert" do
127+
SolidObjects::Instance.singleton_class.define_method(:create!) do |*, **|
128+
raise ActiveRecord::RecordNotUnique, "simulated create race"
129+
end
130+
131+
assert_raises(SolidObjects::ActorDestroyed) do
132+
CartActor.ref("ghost").async.add(product_id: "shirt")
133+
end
134+
ensure
135+
SolidObjects::Instance.singleton_class.send(:remove_method, :create!)
61136
end
62137

63138
test "deduplicates the same idempotent enqueue" do

0 commit comments

Comments
 (0)