Skip to content
Merged
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
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
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.14.4)
solid_objects (0.14.5)
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.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
Expand Down
19 changes: 14 additions & 5 deletions lib/solid_objects/broadcast_executor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand All @@ -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 ||
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.14.4"
VERSION = "0.14.5"
end
3 changes: 3 additions & 0 deletions sig/generated/lib/solid_objects/broadcast_executor.rbs
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
107 changes: 107 additions & 0 deletions test/integration/broadcasts_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
14 changes: 14 additions & 0 deletions test/models/schema_constraints_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading