Skip to content

Commit 106dbf0

Browse files
authored
Merge pull request #82 from cardmagic/fix/sqlite-claim-join-order
fix: SQLite join order for polling scans (0.16.1)
2 parents c18b82b + f02ec5d commit 106dbf0

8 files changed

Lines changed: 266 additions & 6 deletions

‎CHANGELOG.md‎

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

3+
## 0.16.1 - 2026-10-02
4+
5+
- Fix the SQLite join order of the claimed-message scan. The query in
6+
`ActivationManager#claimed_instance_ids` has no condition on a
7+
claimed-message column. The claimed-messages table is a short queue, so it is
8+
often empty when `ANALYZE` or `PRAGMA optimize` runs, and `sqlite_stat1` then
9+
has no row for it. Without statistics, SQLite assumed that the table was
10+
large, and it scanned every instance on each worker poll. The scan now uses
11+
`CROSS JOIN`, which SQLite keeps as a fixed join order, so it starts from the
12+
claimed messages. PostgreSQL and MySQL treat `CROSS JOIN` with an equality as
13+
an inner join and keep their plans. The ids and their order do not change.
14+
- Find SQLite effect recovery candidates through the processing effects.
15+
`sqlite_stat1` records only the average row count for each effect status.
16+
When most effects are complete, SQLite estimated that `status = 'processing'`
17+
matched most of the effects table. It then read every recovery row in key
18+
order to skip a sort, on each effect poll. An index on the recovery filter
19+
does not help, because a recovery row keeps `retired_at` empty after a normal
20+
completion. On SQLite the status test now carries
21+
`likelihood(..., 0.000001)`, so the plan starts from `idx_so_effects_poll`.
22+
The PostgreSQL and MySQL queries do not change.
23+
324
## 0.16.0 - 2026-09-23
425

526
- Find a message whose reference a caller lost.

‎Gemfile.lock‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
PATH
22
remote: .
33
specs:
4-
solid_objects (0.16.0)
4+
solid_objects (0.16.1)
55
actioncable (>= 7.1)
66
actionpack (>= 7.1)
77
actionview (>= 7.1)
@@ -384,7 +384,7 @@ CHECKSUMS
384384
rubocop-rails-omakase (1.1.0) sha256=2af73ac8ee5852de2919abbd2618af9c15c19b512c4cfc1f9a5d3b6ef009109d
385385
ruby-progressbar (1.13.0) sha256=80fc9c47a9b640d6834e0dc7b3c94c9df37f08cb072b7761e4a71e22cff29b33
386386
securerandom (0.4.1) sha256=cc5193d414a4341b6e225f0cb4446aceca8e50d5e1888743fac16987638ea0b1
387-
solid_objects (0.16.0)
387+
solid_objects (0.16.1)
388388
sqlite3 (2.9.5-aarch64-linux-gnu) sha256=78075b6337d3d182c6d2b4691049ed45cd220826160c9ea18946bf6a1de200dc
389389
sqlite3 (2.9.5-aarch64-linux-musl) sha256=18c801185deb4adc01ddb281e8f672a39e3d1729979ca91e39439cd3eac0402d
390390
sqlite3 (2.9.5-arm-linux-gnu) sha256=1bdfca0c7d63998c60b0f4a8e3c8df2d33800ccc4abd2d612eddbbbc92a4c48b

‎lib/solid_objects/activation_manager.rb‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,8 @@ def ready_instance_ids(now)
8686
# @rbs (Time) -> Array[Integer]
8787
def claimed_instance_ids(now)
8888
ClaimedMessage
89-
.joins(:instance)
89+
.joins("CROSS JOIN #{Instance.table_name}")
90+
.where("#{Instance.table_name}.id = #{ClaimedMessage.table_name}.instance_id")
9091
.where("#{Instance.table_name}.paused_at IS NULL")
9192
.where(available_lease_sql, now)
9293
.group(:instance_id)

‎lib/solid_objects/effect_recovery_coordinator.rb‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -63,17 +63,22 @@ def recovery_candidates
6363
effects = Effect.table_name
6464
owners = Process.table_name
6565
bindings = EffectRecovery.table_name
66-
heartbeat = case DatabaseAdapter.family(Record.connection)
66+
family = DatabaseAdapter.family(Record.connection)
67+
heartbeat = case family
6768
when :postgresql then "EXTRACT(EPOCH FROM #{owners}.last_heartbeat_at)"
6869
when :mysql then "UNIX_TIMESTAMP(#{owners}.last_heartbeat_at)"
6970
else "CAST(STRFTIME('%s', #{owners}.last_heartbeat_at) AS REAL)"
7071
end
72+
processing = case family
73+
when :postgresql, :mysql then "#{effects}.status = ?"
74+
else "likelihood(#{effects}.status = ?, 0.000001)"
75+
end
7176
now = SolidObjects.database_adapter.database_clock_now.to_f
7277
threshold = SolidObjects.configuration.process_alive_threshold
7378
EffectRecovery.joins("INNER JOIN #{effects} ON #{effects}.effect_id = #{bindings}.effect_id")
7479
.joins("LEFT JOIN #{owners} ON #{owners}.id = #{effects}.claimed_by")
7580
.where(retired_at: nil).where.not(recovery_operation: nil)
76-
.where("#{effects}.status = ?", "processing")
81+
.where(processing, "processing")
7782
.where("#{owners}.id IS NULL OR #{heartbeat} <= ? - CASE WHEN #{bindings}.recovery_timeout > ? THEN #{bindings}.recovery_timeout ELSE ? END", now, threshold, threshold)
7883
.order(:effect_id).limit(SolidObjects.configuration.claim_scan_limit)
7984
end

‎lib/solid_objects/version.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
# rbs_inline: enabled
22

33
module SolidObjects
4-
VERSION = "0.16.0"
4+
VERSION = "0.16.1"
55
end

‎test/database_test_helper.rb‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,30 @@ def suspend_sqlite_busy_wait(connection)
8383
connection.raw_connection.busy_handler_timeout = configured_sqlite_busy_handler_timeout
8484
end
8585

86+
def restoring_sqlite_statistics(connection)
87+
saved = sqlite_statistics_tables(connection).to_h { |table| [ table, connection.select_all("SELECT * FROM #{table}") ] }
88+
yield
89+
ensure
90+
restore_sqlite_statistics(connection, saved) if saved
91+
end
92+
93+
def restore_sqlite_statistics(connection, saved)
94+
database = connection.raw_connection
95+
saved.each do |table, statistics|
96+
placeholders = Array.new(statistics.columns.length, "?").join(", ")
97+
database.execute("DELETE FROM #{table}")
98+
statistics.rows.each do |row|
99+
database.execute("INSERT INTO #{table} (#{statistics.columns.join(", ")}) VALUES (#{placeholders})", row)
100+
end
101+
end
102+
database.execute("ANALYZE sqlite_schema")
103+
(sqlite_statistics_tables(connection) - saved.keys).each { |table| database.execute("DROP TABLE #{table}") }
104+
end
105+
106+
def sqlite_statistics_tables(connection)
107+
connection.select_values("SELECT name FROM sqlite_schema WHERE type = 'table' AND name LIKE 'sqlite\\_stat%' ESCAPE '\\'")
108+
end
109+
86110
def configured_sqlite_busy_handler_timeout
87111
SolidObjects::Record
88112
.connection_pool
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
# frozen_string_literal: true
2+
3+
require "database_test_helper"
4+
5+
class ActivationCandidatesTest < ActiveSupport::TestCase
6+
class CandidateActor < SolidObjects::Actor
7+
actor_type "activation-candidates"
8+
9+
def run
10+
end
11+
end
12+
13+
test "claimed candidates skip paused and live leases in claim order" do
14+
now = SolidObjects.database_adapter.database_now
15+
owner = create_process
16+
expired = claim_message("expired", claimed_at: now - 40.seconds, owner:, lease_expires_at: now - 1.second)
17+
unleased = claim_message("unleased", claimed_at: now - 30.seconds)
18+
first_tie = claim_message("first-tie", claimed_at: now - 20.seconds)
19+
second_tie = claim_message("second-tie", claimed_at: now - 20.seconds)
20+
unexpiring = claim_message("unexpiring", claimed_at: now - 10.seconds, owner:)
21+
claim_message("paused", claimed_at: now - 50.seconds, paused_at: now - 1.minute)
22+
claim_message("leased", claimed_at: now - 60.seconds, owner:, lease_expires_at: now + 1.minute)
23+
24+
assert_equal [ expired, unleased, first_tie, second_tie, unexpiring ], claimed_instance_ids(now)
25+
end
26+
27+
test "claimed candidates stop at the claim scan limit" do
28+
now = SolidObjects.database_adapter.database_now
29+
SolidObjects.configuration.claim_scan_limit = 2
30+
oldest = claim_message("oldest", claimed_at: now - 30.seconds)
31+
older = claim_message("older", claimed_at: now - 20.seconds)
32+
claim_message("newest", claimed_at: now - 10.seconds)
33+
34+
assert_equal [ oldest, older ], claimed_instance_ids(now)
35+
end
36+
37+
test "SQLite reads claimed candidates from the claimed messages when they have no statistics" do
38+
skip "requires a SQLite query plan" unless database_family == :sqlite
39+
40+
now = SolidObjects.database_adapter.database_now
41+
connection = SolidObjects::Record.connection
42+
restoring_sqlite_statistics(connection) do
43+
SolidObjects::Instance.insert_all!(Array.new(3_000) { |index|
44+
{ actor_type: "activation-candidates", actor_id: "idle-#{index}", state: {}, created_at: now, updated_at: now }
45+
})
46+
connection.execute("ANALYZE")
47+
analyzed_tables = connection.select_values("SELECT DISTINCT tbl FROM sqlite_stat1")
48+
49+
assert_includes analyzed_tables, SolidObjects::Instance.table_name
50+
refute_includes analyzed_tables, SolidObjects::ClaimedMessage.table_name
51+
52+
plan = sqlite_query_plan(connection, SolidObjects::ClaimedMessage.table_name) { claimed_instance_ids(now) }
53+
54+
assert_match(/\A(SCAN|SEARCH) #{SolidObjects::ClaimedMessage.table_name}\b/, plan.first, plan.join("\n"))
55+
assert plan.none? { |step| step.start_with?("SCAN #{SolidObjects::Instance.table_name}") }, plan.join("\n")
56+
end
57+
end
58+
59+
private
60+
61+
def claimed_instance_ids(now)
62+
SolidObjects::ActivationManager.new(owner_id: SecureRandom.uuid).send(:claimed_instance_ids, now)
63+
end
64+
65+
def claim_message(actor_id, claimed_at:, owner: nil, lease_expires_at: nil, paused_at: nil)
66+
message = SolidObjects::Message.find(CandidateActor.ref(actor_id).async.run.id)
67+
SolidObjects::ReadyMessage.where(message:).delete_all
68+
message.instance.update!(
69+
activation_owner_id: owner&.id,
70+
activation_token: owner && SecureRandom.uuid,
71+
activation_expires_at: lease_expires_at,
72+
paused_at:
73+
)
74+
SolidObjects::ClaimedMessage.create!(message:, instance: message.instance, activation_generation: 1, claimed_at:)
75+
message.instance_id
76+
end
77+
78+
def create_process
79+
SolidObjects::Process.create!(
80+
id: SecureRandom.uuid,
81+
kind: "worker",
82+
hostname: "test-host",
83+
pid: ::Process.pid,
84+
started_at: Time.current,
85+
last_heartbeat_at: Time.current,
86+
metadata: {}
87+
)
88+
end
89+
90+
def sqlite_query_plan(connection, table_name)
91+
statements = []
92+
subscriber = ->(*arguments) { statements << arguments.last }
93+
ActiveSupport::Notifications.subscribed(subscriber, "sql.active_record") { yield }
94+
statement = statements.find { |payload| payload[:sql].match?(/\ASELECT .* FROM "#{table_name}"/) }
95+
assert statement, "the #{table_name} query was not captured"
96+
97+
connection.select_all("EXPLAIN QUERY PLAN #{statement[:sql]}", "SQL", statement[:binds]).map { |row| row["detail"] }
98+
end
99+
end
Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
# frozen_string_literal: true
2+
3+
require "database_test_helper"
4+
5+
class EffectRecoveryCandidatesTest < ActiveSupport::TestCase
6+
class CandidateActor < SolidObjects::Actor
7+
actor_type "effect-recovery-candidates"
8+
9+
def run
10+
end
11+
end
12+
13+
setup do
14+
@message = SolidObjects::Message.find(CandidateActor.ref("one").async.run.id)
15+
end
16+
17+
test "recovery candidates are abandoned processing effects with a recovery operation in effect order" do
18+
now = SolidObjects.database_adapter.database_clock_now
19+
threshold = SolidObjects.configuration.process_alive_threshold
20+
live_owner = create_process(last_heartbeat_at: now)
21+
stale_owner = create_process(last_heartbeat_at: now - (threshold * 3))
22+
unowned = create_effect("00000000-0000-4000-8000-000000000009", status: "processing")
23+
create_effect("00000000-0000-4000-8000-000000000001", status: "processing", owner: live_owner)
24+
abandoned = create_effect("00000000-0000-4000-8000-000000000002", status: "processing", owner: stale_owner)
25+
create_effect("00000000-0000-4000-8000-000000000003", status: "processing", owner: stale_owner, recovery_timeout: threshold * 10)
26+
create_effect("00000000-0000-4000-8000-000000000004", status: "processing", retired_at: now)
27+
create_effect("00000000-0000-4000-8000-000000000005", status: "processing", recovery_operation: nil)
28+
create_effect("00000000-0000-4000-8000-000000000006", status: "pending")
29+
create_effect("00000000-0000-4000-8000-000000000007", status: "completed")
30+
31+
assert_equal [ abandoned, unowned ], recovery_candidates.map(&:effect_id)
32+
end
33+
34+
test "recovery candidates stop at the claim scan limit" do
35+
SolidObjects.configuration.claim_scan_limit = 1
36+
create_effect("00000000-0000-4000-8000-000000000002", status: "processing")
37+
first = create_effect("00000000-0000-4000-8000-000000000001", status: "processing")
38+
39+
assert_equal [ first ], recovery_candidates.map(&:effect_id)
40+
end
41+
42+
test "SQLite finds recovery candidates through processing effects when most effects are complete" do
43+
skip "requires a SQLite query plan" unless database_family == :sqlite
44+
45+
now = Time.current
46+
effect_ids = Array.new(3_000) { SecureRandom.uuid }
47+
connection = SolidObjects::Record.connection
48+
restoring_sqlite_statistics(connection) do
49+
SolidObjects::Effect.insert_all!(effect_ids.map { |effect_id|
50+
{ message_id: @message.id, instance_id: @message.instance_id, effect_id:, name: "work", arguments: {},
51+
status: "completed", max_attempts: 3, available_at: now, completed_at: now, created_at: now, updated_at: now }
52+
})
53+
SolidObjects::EffectRecovery.insert_all!(effect_ids.last(300).map { |effect_id|
54+
{ effect_id:, instance_id: @message.instance_id, recovery_operation: "recover", status_operation: "status",
55+
created_at: now, updated_at: now }
56+
})
57+
connection.execute("ANALYZE")
58+
poll_statistics = connection.select_value("SELECT stat FROM sqlite_stat1 WHERE idx = 'idx_so_effects_poll'")
59+
60+
assert_equal 3_000, poll_statistics.split[1].to_i
61+
62+
plan = connection.select_all("EXPLAIN QUERY PLAN #{recovery_candidates.to_sql}").map { |row| row["detail"] }
63+
64+
assert_match(/\ASEARCH #{SolidObjects::Effect.table_name} USING INDEX idx_so_effects_poll \(status=\?\)/, plan.first, plan.join("\n"))
65+
assert plan.none? { |step| step.start_with?("SCAN ") }, plan.join("\n")
66+
end
67+
end
68+
69+
private
70+
71+
def recovery_candidates
72+
SolidObjects::EffectRecoveryCoordinator.new.send(:recovery_candidates)
73+
end
74+
75+
def create_effect(effect_id, status:, owner: nil, recovery_operation: "recover", recovery_timeout: nil, retired_at: nil)
76+
SolidObjects::Effect.create!(
77+
message: @message,
78+
instance: @message.instance,
79+
effect_id:,
80+
name: "work",
81+
arguments: {},
82+
status:,
83+
max_attempts: 3,
84+
available_at: Time.current,
85+
claimed_by: owner&.id,
86+
claimed_at: owner && Time.current
87+
)
88+
SolidObjects::EffectRecovery.create!(
89+
effect_id:,
90+
instance: @message.instance,
91+
recovery_operation:,
92+
status_operation: "status",
93+
recovery_timeout:,
94+
retired_at:
95+
)
96+
effect_id
97+
end
98+
99+
def create_process(last_heartbeat_at:)
100+
SolidObjects::Process.create!(
101+
id: SecureRandom.uuid,
102+
kind: "effect",
103+
hostname: "test-host",
104+
pid: ::Process.pid,
105+
started_at: last_heartbeat_at,
106+
last_heartbeat_at:,
107+
metadata: {}
108+
)
109+
end
110+
end

0 commit comments

Comments
 (0)