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/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/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 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",