Skip to content

Commit 1e9d9cc

Browse files
cardmagicclaude
andcommitted
feat: redrive dead effects and broadcasts in bulk
Retry was one row at a time, and an incident produces dead rows in the hundreds. A scope answers `redrive` now, which opens a durable task and returns at once: task = SolidObjects.dead_letters.effects.redrive( actor_type: "payments", failed_after: 6.hours.ago, limit: 5_000, authorization_context: current_admin ) task.cancel(authorization_context: current_admin) The task is idempotent over its scope and filters, which a dashboard button needs. A unique index on the active scope enforces that in the database rather than in a read followed by a write, so two processes that start the same redrive at the same time share one task. The index is total rather than partial: the column holds the scope while the task runs and NULL once it finishes, so a later redrive of the same scope starts a new task. The supervisor advances one bounded batch per pass and pauses between batches, so a redrive of thousands of rows never holds a transaction longer than one batch and shares the database with delivery. `SolidObjects.redrives` reads tasks back by id and by status. A running task reports what is left to move rather than a stored estimate, because rows die and are retried while it runs. Every transition writes one audit row under the identity that asked for it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent c36cfd0 commit 1e9d9cc

24 files changed

Lines changed: 817 additions & 17 deletions
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
# rbs_inline: enabled
2+
3+
module SolidObjects
4+
class Redrive < Record
5+
self.table_name = SolidObjects.table_name(:redrives)
6+
end
7+
end
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
# rbs_inline: enabled
2+
3+
class AddSolidObjectsRedrives < ActiveRecord::Migration[7.1]
4+
# @rbs () -> void
5+
def change
6+
create_table SolidObjects.table_name(:redrives), id: :string, limit: 64 do |definition|
7+
definition.string :kind, null: false, limit: 32
8+
json_column definition, :filters, null: false
9+
definition.string :status, null: false, default: "running", limit: 32
10+
definition.string :active_scope, limit: 191
11+
definition.integer :moved, null: false, default: 0
12+
definition.integer :move_limit
13+
definition.string :actor, limit: 255
14+
definition.datetime :started_at, null: false, precision: 6
15+
definition.datetime :finished_at, precision: 6
16+
definition.timestamps precision: 6, null: false
17+
18+
definition.index :active_scope, unique: true, name: "idx_so_redrives_active_scope"
19+
definition.index [ :status, :started_at, :id ], name: "idx_so_redrives_poll"
20+
definition.check_constraint "moved >= 0", name: "chk_so_redrives_moved"
21+
definition.check_constraint "move_limit IS NULL OR move_limit > 0", name: "chk_so_redrives_limit"
22+
definition.check_constraint "status IN ('running', 'completed', 'cancelled')", name: "chk_so_redrives_status"
23+
end
24+
end
25+
26+
private
27+
28+
# @rbs () -> Symbol
29+
def json_column_type
30+
connection.adapter_name.match?(/postgres/i) ? :jsonb : :json
31+
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
37+
end

‎lib/solid_objects.rb‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@
2828
require "solid_objects/reference"
2929
require "solid_objects/message_reference"
3030
require "solid_objects/administration_audit"
31+
require "solid_objects/redrive_task"
32+
require "solid_objects/redrive_manager"
33+
require "solid_objects/redrive_runner"
3134
require "solid_objects/dead_letter_scope"
3235
require "solid_objects/dead_letter_manager"
3336
require "solid_objects/message_pruner"
@@ -162,6 +165,11 @@ def dead_letters
162165
@dead_letters ||= DeadLetterManager.new
163166
end
164167

168+
# @rbs () -> RedriveManager
169+
def redrives
170+
@redrives ||= RedriveManager.new
171+
end
172+
165173
# @rbs () -> Administration
166174
def administration
167175
@administration ||= Administration.new
@@ -194,6 +202,7 @@ def reset!
194202
@effect_registry = EffectRegistry.new
195203
@commit_action_registry = CommitActionRegistry.new
196204
@dead_letters = nil
205+
@redrives = nil
197206
@administration = nil
198207
end
199208

‎lib/solid_objects/administration_audit.rb‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,14 +4,14 @@ module SolidObjects
44
module AdministrationAudit
55
module_function
66

7-
# @rbs (action: String, kind: String, ?subject_id: untyped, ?filters: Hash[Symbol | String, untyped]?, authorization_context: untyped) -> void
8-
def record(action:, kind:, authorization_context:, subject_id: nil, filters: nil)
7+
# @rbs (action: String, kind: String, ?subject_id: untyped, ?filters: Hash[Symbol | String, untyped]?, ?actor: String?) -> void
8+
def record(action:, kind:, subject_id: nil, filters: nil, actor: nil)
99
AdministrationEvent.create!(
1010
action:,
1111
kind:,
1212
subject_id: subject_id&.to_s,
1313
filters:,
14-
actor: identity(authorization_context),
14+
actor:,
1515
occurred_at: SolidObjects.database_adapter.database_now
1616
)
1717
nil

‎lib/solid_objects/configuration.rb‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,8 @@ class Configuration
3131
# @rbs @instance_retention_by_actor_type: Hash[String, Numeric]
3232
# @rbs @process_retention: Numeric
3333
# @rbs @prune_batch_size: Integer
34+
# @rbs @redrive_batch_size: Integer
35+
# @rbs @redrive_batch_pause: Float
3436
# @rbs @worker_count: Integer
3537
# @rbs @effect_worker_count: Integer
3638
# @rbs @broadcast_worker_count: Integer
@@ -81,6 +83,8 @@ class Configuration
8183
:instance_retention_by_actor_type,
8284
:process_retention,
8385
:prune_batch_size,
86+
:redrive_batch_size,
87+
:redrive_batch_pause,
8488
:worker_count,
8589
:effect_worker_count,
8690
:broadcast_worker_count,
@@ -136,6 +140,8 @@ def initialize
136140
@instance_retention_by_actor_type = {}
137141
@process_retention = 7.days
138142
@prune_batch_size = 1_000
143+
@redrive_batch_size = 100
144+
@redrive_batch_pause = 0.05
139145
@worker_count = 1
140146
@effect_worker_count = 1
141147
@broadcast_worker_count = 1
@@ -229,6 +235,9 @@ def validate!
229235
positive_values.each do |name, value|
230236
raise ArgumentError, "#{name} must be positive" unless value.positive?
231237
end
238+
if redrive_batch_pause.negative?
239+
raise ArgumentError, "redrive_batch_pause must not be negative"
240+
end
232241
if warn_state_bytes > max_state_bytes
233242
raise ArgumentError, "warn_state_bytes must not exceed max_state_bytes"
234243
end
@@ -301,7 +310,8 @@ def positive_values
301310
shutdown_timeout:,
302311
message_retention:,
303312
process_retention:,
304-
prune_batch_size:
313+
prune_batch_size:,
314+
redrive_batch_size:
305315
}
306316
end
307317

‎lib/solid_objects/dead_letter_manager.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ def retry(dead_letter_id, authorization_context: nil)
1616
action: "dead_letter.retry",
1717
kind: "message",
1818
subject_id: dead_letter.id,
19-
authorization_context:
19+
actor: AdministrationAudit.identity(authorization_context)
2020
)
2121
return MessageReference.from_message(Message.find(dead_letter.retried_message_id)) if dead_letter.retried_message_id
2222

‎lib/solid_objects/dead_letter_scope.rb‎

Lines changed: 40 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ class DeadLetterScope
99
# @rbs @resource: String
1010
# @rbs @identifier: Symbol
1111

12-
attr_reader :resource
12+
attr_reader :resource, :kind
1313

1414
# @rbs (model: untyped, resource: String, identifier: Symbol, kind: String) -> void
1515
def initialize(model:, resource:, identifier:, kind:)
@@ -19,6 +19,14 @@ def initialize(model:, resource:, identifier:, kind:)
1919
@kind = kind
2020
end
2121

22+
# @rbs (String) -> DeadLetterScope
23+
def self.for_kind(kind)
24+
return SolidObjects.dead_letters.effects if kind == "effect"
25+
return SolidObjects.dead_letters.broadcasts if kind == "broadcast"
26+
27+
raise ArgumentError, "unknown dead letter kind #{kind.inspect}"
28+
end
29+
2230
# @rbs (?authorization_context: untyped) -> ActiveRecord::Relation[untyped]
2331
def all(authorization_context: nil)
2432
authorize!(:inspect, authorization_context:)
@@ -33,19 +41,48 @@ def retry(identifier_value, authorization_context: nil)
3341
action: "dead_letter.retry",
3442
kind: kind,
3543
subject_id: identifier_value,
36-
authorization_context:
44+
actor: AdministrationAudit.identity(authorization_context)
3745
)
3846
return row unless row.status == DEAD
3947

4048
revive(row)
4149
row
4250
end
4351

52+
# @rbs (?actor_type: String?, ?failed_after: untyped, ?limit: Integer?, ?authorization_context: untyped) -> RedriveTask
53+
def redrive(actor_type: nil, failed_after: nil, limit: nil, authorization_context: nil)
54+
authorize!(:redrive, authorization_context:)
55+
SolidObjects.redrives.start(
56+
scope: self,
57+
filters: {
58+
"actor_type" => actor_type,
59+
"failed_after" => failed_after&.utc&.iso8601(6),
60+
"limit" => limit
61+
},
62+
authorization_context:
63+
)
64+
end
65+
4466
# @rbs () -> ActiveRecord::Relation[untyped]
4567
def dead
4668
model.where(status: DEAD)
4769
end
4870

71+
# @rbs (Hash[String, untyped]) -> ActiveRecord::Relation[untyped]
72+
def matching(filters)
73+
relation = dead
74+
actor_type = filters["actor_type"]
75+
failed_after = filters["failed_after"]
76+
relation = relation.joins(:instance).where(Instance.table_name => { actor_type: }) if actor_type
77+
relation = relation.where(updated_at: Time.parse(failed_after)..) if failed_after
78+
relation
79+
end
80+
81+
# @rbs (Array[untyped]) -> Integer
82+
def revive_all(identifiers)
83+
model.where(id: identifiers, status: DEAD).update_all(revival_attributes)
84+
end
85+
4986
# @rbs (untyped) -> Integer
5087
def revive(row)
5188
model.where(id: row.id, status: DEAD).update_all(revival_attributes).tap do
@@ -68,7 +105,7 @@ def authorize!(action, authorization_context:, resource_id: nil)
68105

69106
private
70107

71-
attr_reader :model, :identifier, :kind
108+
attr_reader :model, :identifier
72109

73110
# @rbs () -> Hash[Symbol, untyped]
74111
def revival_attributes

‎lib/solid_objects/doctor.rb‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,7 +81,8 @@ def to_s
8181
effect_recoveries: %w[effect_id instance_id recovery_operation status_operation recovery_timeout retired_at],
8282
broadcasts: %w[id message_id instance_id broadcast_id status available_at],
8383
dead_letters: %w[id message_id instance_id actor_type actor_id attempts],
84-
administration_events: %w[id action kind subject_id actor occurred_at]
84+
administration_events: %w[id action kind subject_id actor occurred_at],
85+
redrives: %w[id kind filters status active_scope moved move_limit started_at finished_at]
8586
}.freeze
8687

8788
class ProbeActor < Actor
Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,131 @@
1+
# rbs_inline: enabled
2+
3+
require "digest"
4+
5+
module SolidObjects
6+
class RedriveManager
7+
RUNNING = "running"
8+
COMPLETED = "completed"
9+
CANCELLED = "cancelled"
10+
11+
# @rbs (scope: DeadLetterScope, filters: Hash[String, untyped], authorization_context: untyped) -> RedriveTask
12+
def start(scope:, filters:, authorization_context:)
13+
active_scope = active_scope_for(kind: scope.kind, filters:)
14+
running = Redrive.find_by(active_scope:)
15+
return task_for(running) if running
16+
17+
task_for(open_task(scope:, filters:, active_scope:, authorization_context:))
18+
rescue ActiveRecord::RecordNotUnique
19+
task_for(Redrive.find_by!(active_scope:))
20+
end
21+
22+
# @rbs (String, ?authorization_context: untyped) -> RedriveTask
23+
def find(id, authorization_context: nil)
24+
authorize!(:inspect, authorization_context:, resource_id: id)
25+
task_for(Redrive.find(id))
26+
end
27+
28+
# @rbs (?status: Symbol | String | nil, ?authorization_context: untyped) -> Array[RedriveTask]
29+
def all(status: nil, authorization_context: nil)
30+
authorize!(:inspect, authorization_context:)
31+
relation = Redrive.order(started_at: :desc, id: :desc)
32+
relation = relation.where(status: status.to_s) if status
33+
relation.map { |record| task_for(record) }
34+
end
35+
36+
# @rbs (String, ?authorization_context: untyped) -> RedriveTask
37+
def cancel(id, authorization_context: nil)
38+
authorize!(:cancel, authorization_context:, resource_id: id)
39+
record = Redrive.find(id)
40+
return task_for(record) unless record.status == RUNNING
41+
42+
close(record, status: CANCELLED)
43+
audit(record, action: "redrive.cancel")
44+
task_for(record)
45+
end
46+
47+
# @rbs (Redrive, status: String) -> void
48+
def close(record, status:)
49+
record.update!(
50+
status:,
51+
active_scope: nil,
52+
finished_at: SolidObjects.database_adapter.database_now
53+
)
54+
end
55+
56+
# @rbs (Redrive, action: String) -> void
57+
def audit(record, action:)
58+
AdministrationAudit.record(
59+
action:,
60+
kind: record.kind,
61+
subject_id: record.id,
62+
filters: record.filters,
63+
actor: record.actor
64+
)
65+
end
66+
67+
# @rbs (Redrive) -> RedriveTask
68+
def task_for(record)
69+
RedriveTask.new(
70+
id: record.id,
71+
kind: record.kind,
72+
filters: Serialization.readonly_copy(record.filters),
73+
status: record.status,
74+
moved: record.moved,
75+
remaining: remaining_for(record),
76+
started_at: record.started_at,
77+
finished_at: record.finished_at
78+
)
79+
end
80+
81+
private
82+
83+
# @rbs (scope: DeadLetterScope, filters: Hash[String, untyped], active_scope: String, authorization_context: untyped) -> Redrive
84+
def open_task(scope:, filters:, active_scope:, authorization_context:)
85+
record = Redrive.create!(
86+
id: "redrive_#{SecureRandom.uuid}",
87+
kind: scope.kind,
88+
filters:,
89+
status: RUNNING,
90+
active_scope:,
91+
moved: 0,
92+
move_limit: filters["limit"],
93+
actor: AdministrationAudit.identity(authorization_context),
94+
started_at: SolidObjects.database_adapter.database_now
95+
)
96+
audit(record, action: "redrive.start")
97+
record
98+
end
99+
100+
# A running task reports what is left to move rather than a stored
101+
# estimate, because rows die and are retried while it runs.
102+
# @rbs (Redrive) -> Integer
103+
def remaining_for(record)
104+
return 0 unless record.status == RUNNING
105+
106+
matching = DeadLetterScope.for_kind(record.kind).matching(record.filters).count
107+
limit = record.move_limit
108+
return matching unless limit
109+
110+
[ matching, limit - record.moved ].min
111+
end
112+
113+
# @rbs (kind: String, filters: Hash[String, untyped]) -> String
114+
def active_scope_for(kind:, filters:)
115+
"#{kind}:#{Digest::SHA256.hexdigest(filters.to_json)}"
116+
end
117+
118+
# @rbs (Symbol, authorization_context: untyped, ?resource_id: String?) -> void
119+
def authorize!(action, authorization_context:, resource_id: nil)
120+
authorized = SolidObjects.configuration.authorize_administration.call(
121+
action: action.to_s,
122+
resource: "redrives",
123+
resource_id:,
124+
authorization_context:
125+
)
126+
return if authorized
127+
128+
raise Unauthorized, "actor administration is not authorized"
129+
end
130+
end
131+
end

0 commit comments

Comments
 (0)