diff --git a/CHANGELOG.md b/CHANGELOG.md index 3c049c8..126b6a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/Gemfile.lock b/Gemfile.lock index 42ee6f8..8491504 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -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) @@ -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 diff --git a/lib/solid_objects/activation_manager.rb b/lib/solid_objects/activation_manager.rb index bc0a624..12d6c23 100644 --- a/lib/solid_objects/activation_manager.rb +++ b/lib/solid_objects/activation_manager.rb @@ -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) diff --git a/lib/solid_objects/effect_recovery_coordinator.rb b/lib/solid_objects/effect_recovery_coordinator.rb index 9d0e398..2277498 100644 --- a/lib/solid_objects/effect_recovery_coordinator.rb +++ b/lib/solid_objects/effect_recovery_coordinator.rb @@ -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 diff --git a/lib/solid_objects/version.rb b/lib/solid_objects/version.rb index d9c362b..dd4bb70 100644 --- a/lib/solid_objects/version.rb +++ b/lib/solid_objects/version.rb @@ -1,5 +1,5 @@ # rbs_inline: enabled module SolidObjects - VERSION = "0.16.0" + VERSION = "0.16.1" end diff --git a/test/database_test_helper.rb b/test/database_test_helper.rb index 66d8db5..9ef67a0 100644 --- a/test/database_test_helper.rb +++ b/test/database_test_helper.rb @@ -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 diff --git a/test/integration/activation_candidates_test.rb b/test/integration/activation_candidates_test.rb new file mode 100644 index 0000000..54a5628 --- /dev/null +++ b/test/integration/activation_candidates_test.rb @@ -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 diff --git a/test/integration/effect_recovery_candidates_test.rb b/test/integration/effect_recovery_candidates_test.rb new file mode 100644 index 0000000..9a33af8 --- /dev/null +++ b/test/integration/effect_recovery_candidates_test.rb @@ -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