From c75f6c39f982b0385dc97ca3c6cb97d931d193d6 Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Thu, 3 Sep 2026 07:05:06 -0700 Subject: [PATCH 1/2] fix: split broadcast polling queries The pending and stale predicates previously fed one OR query, forcing eligible rows through a sort before LIMIT 1. Lock one ordered candidate from each state and choose the earlier row so each probe follows the existing polling index without changing delivery order or claimant safety. --- lib/solid_objects/broadcast_executor.rb | 19 +++- .../lib/solid_objects/broadcast_executor.rbs | 3 + test/integration/broadcasts_test.rb | 107 ++++++++++++++++++ test/models/schema_constraints_test.rb | 14 +++ 4 files changed, 138 insertions(+), 5 deletions(-) diff --git a/lib/solid_objects/broadcast_executor.rb b/lib/solid_objects/broadcast_executor.rb index dd70af5..1aaa574 100644 --- a/lib/solid_objects/broadcast_executor.rb +++ b/lib/solid_objects/broadcast_executor.rb @@ -112,11 +112,15 @@ def claim_next database_adapter.transaction do now = database_adapter.database_now stale_at = now - SolidObjects.configuration.process_alive_threshold - relation = Broadcast - .where(status: "pending", available_at: ..now) - .or(Broadcast.where(status: "processing", claimed_at: ..stale_at)) - .order(:available_at, :id) - broadcast = database_adapter.lock_candidates(relation).first + pending_broadcast = claim_candidate( + Broadcast.where(status: "pending", available_at: ..now) + ) + stale_broadcast = claim_candidate( + Broadcast.where(status: "processing", claimed_at: ..stale_at) + ) + broadcast = [ pending_broadcast, stale_broadcast ] + .compact + .min_by { |candidate| [ candidate.available_at, candidate.id ] } next unless broadcast broadcast.update!( @@ -129,6 +133,11 @@ def claim_next end end + # @rbs (ActiveRecord::Relation[Broadcast]) -> Broadcast? + def claim_candidate(relation) + database_adapter.lock_candidates(relation.order(:available_at, :id)).first + end + # @rbs () -> Proc | ActionCableBroadcastAdapter def broadcast_adapter SolidObjects.configuration.broadcast_adapter || diff --git a/sig/generated/lib/solid_objects/broadcast_executor.rbs b/sig/generated/lib/solid_objects/broadcast_executor.rbs index e2776be..f69c94b 100644 --- a/sig/generated/lib/solid_objects/broadcast_executor.rbs +++ b/sig/generated/lib/solid_objects/broadcast_executor.rbs @@ -47,6 +47,9 @@ module SolidObjects # @rbs () -> Broadcast? def claim_next: () -> Broadcast? + # @rbs (ActiveRecord::Relation[Broadcast]) -> Broadcast? + def claim_candidate: (ActiveRecord::Relation[Broadcast]) -> Broadcast? + # @rbs () -> Proc | ActionCableBroadcastAdapter def broadcast_adapter: () -> Proc diff --git a/test/integration/broadcasts_test.rb b/test/integration/broadcasts_test.rb index 03f8beb..2abd4c1 100644 --- a/test/integration/broadcasts_test.rb +++ b/test/integration/broadcasts_test.rb @@ -149,4 +149,111 @@ def reveal(secret:) broadcast_executor&.stop worker&.stop end + + test "recovers the oldest broadcast across pending and stale work" do + pending_reference = PublicCounterActor.ref("pending").async.increment + stale_reference = PublicCounterActor.ref("stale").async.increment + worker = SolidObjects::Worker.new + worker.run_until_idle + stale_process_registry = SolidObjects::ProcessRegistry.new + stale_process = stale_process_registry.register(kind: "broadcast") + now = SolidObjects.database_adapter.database_now + pending_broadcast = SolidObjects::Broadcast.find_by!(message_id: pending_reference.id) + pending_broadcast.update!(available_at: now - 1.minute) + stale_broadcast = SolidObjects::Broadcast.find_by!(message_id: stale_reference.id) + stale_broadcast.update!( + status: "processing", + available_at: now - 2.minutes, + claimed_by: stale_process.id, + claimed_at: now - SolidObjects.configuration.process_alive_threshold - 1.second + ) + delivered = Queue.new + SolidObjects.configuration.broadcast_adapter = ->(broadcast) { delivered << broadcast.id } + broadcast_executor = SolidObjects::BroadcastExecutor.new + + assert broadcast_executor.run_once + + assert_equal stale_broadcast.id, delivered.pop + assert_equal "delivered", stale_broadcast.reload.status + assert_equal "pending", pending_broadcast.reload.status + ensure + broadcast_executor&.stop + stale_process_registry&.stop + worker&.stop + end + + test "concurrent executors claim different broadcasts" do + pending_reference = PublicCounterActor.ref("pending").async.increment + stale_reference = PublicCounterActor.ref("stale").async.increment + worker = SolidObjects::Worker.new + worker.run_until_idle + stale_process_registry = SolidObjects::ProcessRegistry.new + stale_process = stale_process_registry.register(kind: "broadcast") + now = SolidObjects.database_adapter.database_now + pending_broadcast = SolidObjects::Broadcast.find_by!(message_id: pending_reference.id) + pending_broadcast.update!(available_at: now - 1.minute) + stale_broadcast = SolidObjects::Broadcast.find_by!(message_id: stale_reference.id) + stale_broadcast.update!( + status: "processing", + available_at: now - 2.minutes, + claimed_by: stale_process.id, + claimed_at: now - SolidObjects.configuration.process_alive_threshold - 1.second + ) + claims = Queue.new + release = Queue.new + SolidObjects.configuration.broadcast_adapter = lambda do |broadcast| + claims << broadcast.id + release.pop + end + executor_a = SolidObjects::BroadcastExecutor.new + executor_b = SolidObjects::BroadcastExecutor.new + + thread_a = Thread.new { executor_a.run_once } + assert_equal stale_broadcast.id, Timeout.timeout(5) { claims.pop } + thread_b = Thread.new { executor_b.run_once } + assert_equal pending_broadcast.id, Timeout.timeout(5) { claims.pop } + 2.times { release << true } + + assert thread_a.value + assert thread_b.value + assert_equal %w[delivered delivered], SolidObjects::Broadcast.order(:id).pluck(:status) + ensure + 2.times { release << true } if release + thread_a&.join(2) + thread_b&.join(2) + executor_a&.stop + executor_b&.stop + stale_process_registry&.stop + worker&.stop + end + + test "polls pending and stale broadcasts separately" do + PublicCounterActor.ref("one").async.increment + worker = SolidObjects::Worker.new + worker.run_until_idle + SolidObjects.configuration.broadcast_adapter = ->(broadcast) { broadcast } + broadcast_executor = SolidObjects::BroadcastExecutor.new + polling_queries = [] + subscription = ActiveSupport::Notifications.subscribe("sql.active_record") do |event| + query = event.payload.fetch(:sql).squish + if query.match?(/SELECT .* FROM ["`]solid_objects_broadcasts["`]/) && + query.match?(/ORDER BY .*available_at.*id.*LIMIT/i) + polling_queries << query + end + end + + assert broadcast_executor.run_once + + assert_equal 2, polling_queries.length + assert polling_queries.one? { |query| !query.include?("claimed_at") } + assert polling_queries.one? { |query| query.include?("claimed_at") } + polling_queries.each { |query| refute_match(/\sOR\s/i, query) } + if database_family != :sqlite + polling_queries.each { |query| assert_match(/FOR UPDATE SKIP LOCKED\z/i, query) } + end + ensure + ActiveSupport::Notifications.unsubscribe(subscription) if subscription + broadcast_executor&.stop + worker&.stop + end end diff --git a/test/models/schema_constraints_test.rb b/test/models/schema_constraints_test.rb index ef06190..b6814b3 100644 --- a/test/models/schema_constraints_test.rb +++ b/test/models/schema_constraints_test.rb @@ -48,6 +48,20 @@ class SchemaConstraintsTest < ActiveSupport::TestCase assert indexes.all? { |index| index.where.nil? } end + test "indexes each polling query in delivery order" do + expected_indexes = { + "solid_objects_effects" => %w[status available_at id], + "solid_objects_broadcasts" => %w[status available_at id], + "solid_objects_reminders" => %w[status next_run_at id] + } + + expected_indexes.each do |table, columns| + indexes = ActiveRecord::Base.connection.indexes(table).map(&:columns) + + assert_includes indexes, columns + end + end + test "links every runtime claim owner to the process registry" do expected_claim_foreign_keys = { "solid_objects_claimed_messages" => "process_id", From ac0df54152c184f87a327715a6777c8b9b5e0abd Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Thu, 3 Sep 2026 07:17:11 -0700 Subject: [PATCH 2/2] chore: bump version to 0.14.5 Prepare the patch release for the broadcast polling query fix and record that existing installations need no migration or new index. --- CHANGELOG.md | 9 +++++++++ Gemfile.lock | 4 ++-- lib/solid_objects/version.rb | 2 +- 3 files changed, 12 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a7e801..09254f3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,14 @@ # Changelog +## 0.14.5 - 2026-09-03 + +- Split broadcast claiming into separate pending and stale-processing probes, + then choose the oldest locked candidate across both. The old `OR` query made + MySQL, PostgreSQL, and SQLite collect and sort eligible rows before applying + `LIMIT 1`; each probe now follows the existing + `(status, available_at, id)` index while preserving delivery order, recovery, + and concurrent claimant safety. No migration or new index is required. + ## 0.14.4 - 2026-08-30 - Reuse the encoding the after image already built. `State#to_h` copies the diff --git a/Gemfile.lock b/Gemfile.lock index ed845a2..1c8a95e 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -1,7 +1,7 @@ PATH remote: . specs: - solid_objects (0.14.4) + solid_objects (0.14.5) 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.14.4) + solid_objects (0.14.5) 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 22aeb34..f90d21e 100644 --- a/lib/solid_objects/version.rb +++ b/lib/solid_objects/version.rb @@ -1,5 +1,5 @@ # rbs_inline: enabled module SolidObjects - VERSION = "0.14.4" + VERSION = "0.14.5" end