Skip to content
Open
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
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,26 @@
# Changelog

## 0.16.1 - 2026-10-02

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

## 0.16.0 - 2026-09-23

- Find a message whose reference a caller lost.
Expand Down
4 changes: 2 additions & 2 deletions Gemfile.lock
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
PATH
remote: .
specs:
solid_objects (0.16.0)
solid_objects (0.16.1)
actioncable (>= 7.1)
actionpack (>= 7.1)
actionview (>= 7.1)
Expand Down Expand Up @@ -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.16.0)
solid_objects (0.16.1)
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
Expand Down
3 changes: 2 additions & 1 deletion lib/solid_objects/activation_manager.rb
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,8 @@ def ready_instance_ids(now)
# @rbs (Time) -> Array[Integer]
def claimed_instance_ids(now)
ClaimedMessage
.joins(:instance)
.joins("CROSS JOIN #{Instance.table_name}")
.where("#{Instance.table_name}.id = #{ClaimedMessage.table_name}.instance_id")
.where("#{Instance.table_name}.paused_at IS NULL")
.where(available_lease_sql, now)
.group(:instance_id)
Expand Down
9 changes: 7 additions & 2 deletions lib/solid_objects/effect_recovery_coordinator.rb
Original file line number Diff line number Diff line change
Expand Up @@ -63,17 +63,22 @@ def recovery_candidates
effects = Effect.table_name
owners = Process.table_name
bindings = EffectRecovery.table_name
heartbeat = case DatabaseAdapter.family(Record.connection)
family = DatabaseAdapter.family(Record.connection)
heartbeat = case family
when :postgresql then "EXTRACT(EPOCH FROM #{owners}.last_heartbeat_at)"
when :mysql then "UNIX_TIMESTAMP(#{owners}.last_heartbeat_at)"
else "CAST(STRFTIME('%s', #{owners}.last_heartbeat_at) AS REAL)"
end
processing = case family
when :postgresql, :mysql then "#{effects}.status = ?"
else "likelihood(#{effects}.status = ?, 0.000001)"
end
now = SolidObjects.database_adapter.database_clock_now.to_f
threshold = SolidObjects.configuration.process_alive_threshold
EffectRecovery.joins("INNER JOIN #{effects} ON #{effects}.effect_id = #{bindings}.effect_id")
.joins("LEFT JOIN #{owners} ON #{owners}.id = #{effects}.claimed_by")
.where(retired_at: nil).where.not(recovery_operation: nil)
.where("#{effects}.status = ?", "processing")
.where(processing, "processing")
.where("#{owners}.id IS NULL OR #{heartbeat} <= ? - CASE WHEN #{bindings}.recovery_timeout > ? THEN #{bindings}.recovery_timeout ELSE ? END", now, threshold, threshold)
.order(:effect_id).limit(SolidObjects.configuration.claim_scan_limit)
end
Expand Down
2 changes: 1 addition & 1 deletion lib/solid_objects/version.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
# rbs_inline: enabled

module SolidObjects
VERSION = "0.16.0"
VERSION = "0.16.1"
end
24 changes: 24 additions & 0 deletions test/database_test_helper.rb
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,30 @@ def suspend_sqlite_busy_wait(connection)
connection.raw_connection.busy_handler_timeout = configured_sqlite_busy_handler_timeout
end

def restoring_sqlite_statistics(connection)
saved = sqlite_statistics_tables(connection).to_h { |table| [ table, connection.select_all("SELECT * FROM #{table}") ] }
yield
ensure
restore_sqlite_statistics(connection, saved) if saved
end

def restore_sqlite_statistics(connection, saved)
database = connection.raw_connection
saved.each do |table, statistics|
placeholders = Array.new(statistics.columns.length, "?").join(", ")
database.execute("DELETE FROM #{table}")
statistics.rows.each do |row|
database.execute("INSERT INTO #{table} (#{statistics.columns.join(", ")}) VALUES (#{placeholders})", row)
end
end
database.execute("ANALYZE sqlite_schema")
(sqlite_statistics_tables(connection) - saved.keys).each { |table| database.execute("DROP TABLE #{table}") }
end

def sqlite_statistics_tables(connection)
connection.select_values("SELECT name FROM sqlite_schema WHERE type = 'table' AND name LIKE 'sqlite\\_stat%' ESCAPE '\\'")
end

def configured_sqlite_busy_handler_timeout
SolidObjects::Record
.connection_pool
Expand Down
99 changes: 99 additions & 0 deletions test/integration/activation_candidates_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
# frozen_string_literal: true

require "database_test_helper"

class ActivationCandidatesTest < ActiveSupport::TestCase
class CandidateActor < SolidObjects::Actor
actor_type "activation-candidates"

def run
end
end

test "claimed candidates skip paused and live leases in claim order" do
now = SolidObjects.database_adapter.database_now
owner = create_process
expired = claim_message("expired", claimed_at: now - 40.seconds, owner:, lease_expires_at: now - 1.second)
unleased = claim_message("unleased", claimed_at: now - 30.seconds)
first_tie = claim_message("first-tie", claimed_at: now - 20.seconds)
second_tie = claim_message("second-tie", claimed_at: now - 20.seconds)
unexpiring = claim_message("unexpiring", claimed_at: now - 10.seconds, owner:)
claim_message("paused", claimed_at: now - 50.seconds, paused_at: now - 1.minute)
claim_message("leased", claimed_at: now - 60.seconds, owner:, lease_expires_at: now + 1.minute)

assert_equal [ expired, unleased, first_tie, second_tie, unexpiring ], claimed_instance_ids(now)
end

test "claimed candidates stop at the claim scan limit" do
now = SolidObjects.database_adapter.database_now
SolidObjects.configuration.claim_scan_limit = 2
oldest = claim_message("oldest", claimed_at: now - 30.seconds)
older = claim_message("older", claimed_at: now - 20.seconds)
claim_message("newest", claimed_at: now - 10.seconds)

assert_equal [ oldest, older ], claimed_instance_ids(now)
end

test "SQLite reads claimed candidates from the claimed messages when they have no statistics" do
skip "requires a SQLite query plan" unless database_family == :sqlite

now = SolidObjects.database_adapter.database_now
connection = SolidObjects::Record.connection
restoring_sqlite_statistics(connection) do
SolidObjects::Instance.insert_all!(Array.new(3_000) { |index|
{ actor_type: "activation-candidates", actor_id: "idle-#{index}", state: {}, created_at: now, updated_at: now }
})
connection.execute("ANALYZE")
analyzed_tables = connection.select_values("SELECT DISTINCT tbl FROM sqlite_stat1")

assert_includes analyzed_tables, SolidObjects::Instance.table_name
refute_includes analyzed_tables, SolidObjects::ClaimedMessage.table_name

plan = sqlite_query_plan(connection, SolidObjects::ClaimedMessage.table_name) { claimed_instance_ids(now) }

assert_match(/\A(SCAN|SEARCH) #{SolidObjects::ClaimedMessage.table_name}\b/, plan.first, plan.join("\n"))
assert plan.none? { |step| step.start_with?("SCAN #{SolidObjects::Instance.table_name}") }, plan.join("\n")
end
end

private

def claimed_instance_ids(now)
SolidObjects::ActivationManager.new(owner_id: SecureRandom.uuid).send(:claimed_instance_ids, now)
end

def claim_message(actor_id, claimed_at:, owner: nil, lease_expires_at: nil, paused_at: nil)
message = SolidObjects::Message.find(CandidateActor.ref(actor_id).async.run.id)
SolidObjects::ReadyMessage.where(message:).delete_all
message.instance.update!(
activation_owner_id: owner&.id,
activation_token: owner && SecureRandom.uuid,
activation_expires_at: lease_expires_at,
paused_at:
)
SolidObjects::ClaimedMessage.create!(message:, instance: message.instance, activation_generation: 1, claimed_at:)
message.instance_id
end

def create_process
SolidObjects::Process.create!(
id: SecureRandom.uuid,
kind: "worker",
hostname: "test-host",
pid: ::Process.pid,
started_at: Time.current,
last_heartbeat_at: Time.current,
metadata: {}
)
end

def sqlite_query_plan(connection, table_name)
statements = []
subscriber = ->(*arguments) { statements << arguments.last }
ActiveSupport::Notifications.subscribed(subscriber, "sql.active_record") { yield }
statement = statements.find { |payload| payload[:sql].match?(/\ASELECT .* FROM "#{table_name}"/) }
assert statement, "the #{table_name} query was not captured"

connection.select_all("EXPLAIN QUERY PLAN #{statement[:sql]}", "SQL", statement[:binds]).map { |row| row["detail"] }
end
end
110 changes: 110 additions & 0 deletions test/integration/effect_recovery_candidates_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
# frozen_string_literal: true

require "database_test_helper"

class EffectRecoveryCandidatesTest < ActiveSupport::TestCase
class CandidateActor < SolidObjects::Actor
actor_type "effect-recovery-candidates"

def run
end
end

setup do
@message = SolidObjects::Message.find(CandidateActor.ref("one").async.run.id)
end

test "recovery candidates are abandoned processing effects with a recovery operation in effect order" do
now = SolidObjects.database_adapter.database_clock_now
threshold = SolidObjects.configuration.process_alive_threshold
live_owner = create_process(last_heartbeat_at: now)
stale_owner = create_process(last_heartbeat_at: now - (threshold * 3))
unowned = create_effect("00000000-0000-4000-8000-000000000009", status: "processing")
create_effect("00000000-0000-4000-8000-000000000001", status: "processing", owner: live_owner)
abandoned = create_effect("00000000-0000-4000-8000-000000000002", status: "processing", owner: stale_owner)
create_effect("00000000-0000-4000-8000-000000000003", status: "processing", owner: stale_owner, recovery_timeout: threshold * 10)
create_effect("00000000-0000-4000-8000-000000000004", status: "processing", retired_at: now)
create_effect("00000000-0000-4000-8000-000000000005", status: "processing", recovery_operation: nil)
create_effect("00000000-0000-4000-8000-000000000006", status: "pending")
create_effect("00000000-0000-4000-8000-000000000007", status: "completed")

assert_equal [ abandoned, unowned ], recovery_candidates.map(&:effect_id)
end

test "recovery candidates stop at the claim scan limit" do
SolidObjects.configuration.claim_scan_limit = 1
create_effect("00000000-0000-4000-8000-000000000002", status: "processing")
first = create_effect("00000000-0000-4000-8000-000000000001", status: "processing")

assert_equal [ first ], recovery_candidates.map(&:effect_id)
end

test "SQLite finds recovery candidates through processing effects when most effects are complete" do
skip "requires a SQLite query plan" unless database_family == :sqlite

now = Time.current
effect_ids = Array.new(3_000) { SecureRandom.uuid }
connection = SolidObjects::Record.connection
restoring_sqlite_statistics(connection) do
SolidObjects::Effect.insert_all!(effect_ids.map { |effect_id|
{ message_id: @message.id, instance_id: @message.instance_id, effect_id:, name: "work", arguments: {},
status: "completed", max_attempts: 3, available_at: now, completed_at: now, created_at: now, updated_at: now }
})
SolidObjects::EffectRecovery.insert_all!(effect_ids.last(300).map { |effect_id|
{ effect_id:, instance_id: @message.instance_id, recovery_operation: "recover", status_operation: "status",
created_at: now, updated_at: now }
})
connection.execute("ANALYZE")
poll_statistics = connection.select_value("SELECT stat FROM sqlite_stat1 WHERE idx = 'idx_so_effects_poll'")

assert_equal 3_000, poll_statistics.split[1].to_i

plan = connection.select_all("EXPLAIN QUERY PLAN #{recovery_candidates.to_sql}").map { |row| row["detail"] }

assert_match(/\ASEARCH #{SolidObjects::Effect.table_name} USING INDEX idx_so_effects_poll \(status=\?\)/, plan.first, plan.join("\n"))
assert plan.none? { |step| step.start_with?("SCAN ") }, plan.join("\n")
end
end

private

def recovery_candidates
SolidObjects::EffectRecoveryCoordinator.new.send(:recovery_candidates)
end

def create_effect(effect_id, status:, owner: nil, recovery_operation: "recover", recovery_timeout: nil, retired_at: nil)
SolidObjects::Effect.create!(
message: @message,
instance: @message.instance,
effect_id:,
name: "work",
arguments: {},
status:,
max_attempts: 3,
available_at: Time.current,
claimed_by: owner&.id,
claimed_at: owner && Time.current
)
SolidObjects::EffectRecovery.create!(
effect_id:,
instance: @message.instance,
recovery_operation:,
status_operation: "status",
recovery_timeout:,
retired_at:
)
effect_id
end

def create_process(last_heartbeat_at:)
SolidObjects::Process.create!(
id: SecureRandom.uuid,
kind: "effect",
hostname: "test-host",
pid: ::Process.pid,
started_at: last_heartbeat_at,
last_heartbeat_at:,
metadata: {}
)
end
end
Loading