Skip to content

Commit 5c7d4f0

Browse files
cardmagicclaude
andcommitted
fix: claim a redrive the way every other role claims work
`Redrive.lock` takes a plain `FOR UPDATE`, so with a redrive thread in every supervisor the second one waits on the first rather than taking the next task. It uses `lock_candidates` now, which is `FOR UPDATE SKIP LOCKED` where the database has it, as the actor, effect, broadcast, and reminder claims already do. `RedriveTask#cancel` moves out of the `Data.define` block, because rbs-inline does not read that block and the shipped signature therefore omitted a documented method. The migrations inline their one-use JSON helper, and four comments that restated the code they sat above are gone. The reasoning they carried is in the commit that introduced each one. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 94786d4 commit 5c7d4f0

12 files changed

Lines changed: 15 additions & 49 deletions

‎db/migrate/20260922000000_add_solid_objects_administration_events.rb‎

Lines changed: 2 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ def change
77
definition.string :action, null: false, limit: 64
88
definition.string :kind, null: false, limit: 32
99
definition.string :subject_id, limit: 191
10-
json_column definition, :filters
10+
definition.public_send(json_type, :filters)
1111
definition.string :actor, limit: 255
1212
definition.datetime :occurred_at, null: false, precision: 6
1313
definition.timestamps precision: 6, null: false
@@ -20,12 +20,7 @@ def change
2020
private
2121

2222
# @rbs () -> Symbol
23-
def json_column_type
23+
def json_type
2424
connection.adapter_name.match?(/postgres/i) ? :jsonb : :json
2525
end
26-
27-
# @rbs (untyped, Symbol, ?null: bool) -> void
28-
def json_column(definition, name, null: true)
29-
definition.public_send(json_column_type, name, null:)
30-
end
3126
end

‎db/migrate/20260922000001_add_solid_objects_redrives.rb‎

Lines changed: 2 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ class AddSolidObjectsRedrives < ActiveRecord::Migration[7.1]
55
def change
66
create_table SolidObjects.table_name(:redrives), id: :string, limit: 64 do |definition|
77
definition.string :kind, null: false, limit: 32
8-
json_column definition, :filters, null: false
8+
definition.public_send(json_type, :filters, null: false)
99
definition.string :status, null: false, default: "running", limit: 32
1010
definition.string :active_scope, limit: 191
1111
definition.integer :moved, null: false, default: 0
@@ -26,12 +26,7 @@ def change
2626
private
2727

2828
# @rbs () -> Symbol
29-
def json_column_type
29+
def json_type
3030
connection.adapter_name.match?(/postgres/i) ? :jsonb : :json
3131
end
32-
33-
# @rbs (untyped, Symbol, ?null: bool) -> void
34-
def json_column(definition, name, null: true)
35-
definition.public_send(json_column_type, name, null:)
36-
end
3732
end

‎lib/solid_objects/dead_letter_scope.rb‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -67,9 +67,6 @@ def dead
6767
model.where(status: DEAD)
6868
end
6969

70-
# A redrive moves what was already dead when it started. Without that bound
71-
# a row that fails again lands back in the same scope, and a task whose
72-
# handler is still broken would move it forever.
7370
# @rbs (Hash[String, untyped], ?dead_before: untyped) -> ActiveRecord::Relation[untyped]
7471
def matching(filters, dead_before: nil)
7572
relation = dead

‎lib/solid_objects/redrive_manager.rb‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -98,8 +98,6 @@ def open_task(scope:, filters:, active_scope:, authorization_context:)
9898
record
9999
end
100100

101-
# A running task reports what is left to move rather than a stored
102-
# estimate, because rows die and are retried while it runs.
103101
# @rbs (Redrive) -> Integer
104102
def remaining_for(record)
105103
return 0 unless record.status == RUNNING

‎lib/solid_objects/redrive_runner.rb‎

Lines changed: 3 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,10 @@
11
# rbs_inline: enabled
22

33
module SolidObjects
4-
# Moves dead rows back to pending for one redrive task at a time, in bounded
5-
# batches. Each batch is its own short transaction, so a redrive of thousands
6-
# of rows never holds a lock long enough to starve delivery.
74
class RedriveRunner
85
# @rbs () -> bool
96
def run_once
10-
database_adapter.transaction do
7+
SolidObjects.database_adapter.transaction do
118
record = claim
129
next false unless record
1310

@@ -19,7 +16,8 @@ def run_once
1916

2017
# @rbs () -> Redrive?
2118
def claim
22-
Redrive.lock.where(status: RedriveManager::RUNNING).order(:started_at, :id).first
19+
relation = Redrive.where(status: RedriveManager::RUNNING).order(:started_at, :id)
20+
SolidObjects.database_adapter.lock_candidates(relation).first
2321
end
2422

2523
# @rbs (Redrive) -> bool
@@ -64,10 +62,5 @@ def finish(record)
6462
manager.close(record, status: RedriveManager::COMPLETED)
6563
manager.audit(record, action: "redrive.finish")
6664
end
67-
68-
# @rbs () -> DatabaseAdapter
69-
def database_adapter
70-
SolidObjects.database_adapter
71-
end
7265
end
7366
end

‎lib/solid_objects/redrive_task.rb‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,9 @@
33
module SolidObjects
44
RedriveTask = Data.define(
55
:id, :kind, :filters, :status, :moved, :remaining, :started_at, :finished_at
6-
) do
6+
)
7+
8+
class RedriveTask
79
# @rbs (?authorization_context: untyped) -> RedriveTask
810
def cancel(authorization_context: nil)
911
SolidObjects.redrives.cancel(id, authorization_context:)

‎lib/solid_objects/supervisor.rb‎

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -207,10 +207,6 @@ def retention_pause(failures)
207207
[ backoff, interval ].min
208208
end
209209

210-
# An operator starts a redrive and expects it to move, so the supervisor
211-
# advances it rather than ask the application to schedule a job. Each pass
212-
# takes one bounded batch, and the loop pauses between batches so a large
213-
# redrive shares the database with delivery.
214210
# @rbs () -> void
215211
def redrive_loop
216212
while @started

‎sig/generated/lib/solid_objects/dead_letter_scope.rbs‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,6 @@ module SolidObjects
3434
# @rbs () -> ActiveRecord::Relation[untyped]
3535
def dead: () -> ActiveRecord::Relation[untyped]
3636

37-
# A redrive moves what was already dead when it started. Without that bound
38-
# a row that fails again lands back in the same scope, and a task whose
39-
# handler is still broken would move it forever.
4037
# @rbs (Hash[String, untyped], ?dead_before: untyped) -> ActiveRecord::Relation[untyped]
4138
def matching: (Hash[String, untyped], ?dead_before: untyped) -> ActiveRecord::Relation[untyped]
4239

‎sig/generated/lib/solid_objects/redrive_manager.rbs‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,8 +34,6 @@ module SolidObjects
3434
# @rbs (scope: DeadLetterScope, filters: Hash[String, untyped], active_scope: String, authorization_context: untyped) -> Redrive
3535
def open_task: (scope: DeadLetterScope, filters: Hash[String, untyped], active_scope: String, authorization_context: untyped) -> Redrive
3636

37-
# A running task reports what is left to move rather than a stored
38-
# estimate, because rows die and are retried while it runs.
3937
# @rbs (Redrive) -> Integer
4038
def remaining_for: (Redrive) -> Integer
4139

‎sig/generated/lib/solid_objects/redrive_runner.rbs‎

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,6 @@
11
# Generated from lib/solid_objects/redrive_runner.rb with RBS::Inline
22

33
module SolidObjects
4-
# Moves dead rows back to pending for one redrive task at a time, in bounded
5-
# batches. Each batch is its own short transaction, so a redrive of thousands
6-
# of rows never holds a lock long enough to starve delivery.
74
class RedriveRunner
85
# @rbs () -> bool
96
def run_once: () -> bool
@@ -24,8 +21,5 @@ module SolidObjects
2421

2522
# @rbs (Redrive) -> void
2623
def finish: (Redrive) -> void
27-
28-
# @rbs () -> DatabaseAdapter
29-
def database_adapter: () -> DatabaseAdapter
3024
end
3125
end

0 commit comments

Comments
 (0)