Skip to content

Commit 0630154

Browse files
authored
Merge pull request #68 from cardmagic/feat/effect-recovery
Coordinate abandoned effect recovery
2 parents 0b92ed1 + db97952 commit 0630154

41 files changed

Lines changed: 1452 additions & 27 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.rubocop.yml‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,9 @@
33

44
inherit_gem: { rubocop-rails-omakase: rubocop.yml }
55

6+
Layout/LeadingCommentSpace:
7+
AllowRBSInlineAnnotation: true
8+
69
AllCops:
710
TargetRubyVersion: 3.3
811
Exclude:

‎CHANGELOG.md‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,20 @@
11
# Changelog
22

3+
## 0.15.0 - 2026-09-15
4+
5+
- Maintain effect-owner heartbeats during long-running handlers, completion, and
6+
failure handling, so healthy external I/O cannot trigger abandoned recovery.
7+
Report failed heartbeat updates and retry on the next configured interval
8+
without consuming effect attempts.
9+
- Return a stable effect handle from every `emit`. Wrappers must return it;
10+
operations relying on an implicit `nil` result should return `nil` explicitly.
11+
- Add abandoned effect recovery with `on_recovery`, optional `on_status`, staged
12+
`request_effect_recovery`, and an extending `recovery_timeout` in seconds.
13+
Retirement and durable callbacks share the claim-locking transaction. Add
14+
frozen outcome constants and public RBS envelopes. Install the new recovery
15+
binding migration before upgrading runtime processes. External actions still
16+
require idempotency; retirement does not cancel an old handler or remote call.
17+
318
## 0.14.7 - 2026-09-14
419

520
- Publish RBS contracts for effect callback envelopes and Ruby error summaries.

‎Gemfile.lock‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
PATH
22
remote: .
33
specs:
4-
solid_objects (0.14.7)
4+
solid_objects (0.15.0)
55
actioncable (>= 7.1)
66
actionpack (>= 7.1)
77
actionview (>= 7.1)
@@ -384,7 +384,7 @@ CHECKSUMS
384384
rubocop-rails-omakase (1.1.0) sha256=2af73ac8ee5852de2919abbd2618af9c15c19b512c4cfc1f9a5d3b6ef009109d
385385
ruby-progressbar (1.13.0) sha256=80fc9c47a9b640d6834e0dc7b3c94c9df37f08cb072b7761e4a71e22cff29b33
386386
securerandom (0.4.1) sha256=cc5193d414a4341b6e225f0cb4446aceca8e50d5e1888743fac16987638ea0b1
387-
solid_objects (0.14.7)
387+
solid_objects (0.15.0)
388388
sqlite3 (2.9.5-aarch64-linux-gnu) sha256=78075b6337d3d182c6d2b4691049ed45cd220826160c9ea18946bf6a1de200dc
389389
sqlite3 (2.9.5-aarch64-linux-musl) sha256=18c801185deb4adc01ddb281e8f672a39e3d1729979ca91e39439cd3eac0402d
390390
sqlite3 (2.9.5-arm-linux-gnu) sha256=1bdfca0c7d63998c60b0f4a8e3c8df2d33800ccc4abd2d612eddbbbc92a4c48b

‎README.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,7 @@ SQL and should be allowed to enjoy that.
149149
- One successful turn commits actor state and staged reminders, messages, effects, commit actions, and broadcasts together.
150150
- Fencing prevents stale Ruby code from committing, but it cannot stop that code from continuing to run.
151151
- External effects can repeat and must deduplicate with the stable effect ID or another durable idempotency key.
152+
- [Effect recovery](docs/effect-recovery.md) uses stable `emit` handles and `on_recovery` to retire abandoned work atomically with a durable actor callback.
152153
- Actor handlers may read application records but cannot write them directly. Use `commit_action` for bounded same-database writes and `emit` for external I/O.
153154
- `async`, reminders, effects, and broadcasts need `bundle exec solid_objects start`. Pending work remains in SQL while it is down.
154155
- One hot identity is intentionally sequential. There are no transactions across actor identities.
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
# rbs_inline: enabled
2+
3+
module SolidObjects
4+
class EffectRecovery < Record
5+
self.table_name = SolidObjects.table_name(:effect_recoveries)
6+
self.primary_key = "effect_id"
7+
8+
belongs_to :instance, class_name: "SolidObjects::Instance"
9+
end
10+
end

‎benchmark/support.rb‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -428,9 +428,11 @@ def migrate
428428
require_relative "../db/migrate/20260805000000_create_solid_objects_tables"
429429
require_relative "../db/migrate/20260806000000_add_state_revision_to_solid_objects_instances"
430430
require_relative "../db/migrate/20260813000000_rename_message_dispatch_columns"
431+
require_relative "../db/migrate/20260915000000_add_solid_objects_effect_recoveries"
431432
CreateSolidObjectsTables.new.migrate(:up)
432433
AddStateRevisionToSolidObjectsInstances.new.migrate(:up)
433434
RenameMessageDispatchColumns.new.migrate(:up)
435+
AddSolidObjectsEffectRecoveries.new.migrate(:up)
434436
end
435437

436438
# @rbs () -> void
@@ -444,6 +446,7 @@ def load_models
444446
claimed_message
445447
reminder
446448
effect
449+
effect_recovery
447450
broadcast
448451
dead_letter
449452
].each do |model|
Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
# rbs_inline: enabled
2+
3+
class AddSolidObjectsEffectRecoveries < ActiveRecord::Migration[7.1]
4+
# @rbs () -> void
5+
def change
6+
create_table SolidObjects.table_name(:effect_recoveries), id: :string, limit: 36, primary_key: :effect_id do |definition|
7+
definition.references :instance, null: false,
8+
foreign_key: { to_table: SolidObjects.table_name(:instances), on_delete: :cascade, name: "fk_so_effect_recoveries_instance" }
9+
definition.string :recovery_operation, limit: 191
10+
definition.string :status_operation, limit: 191
11+
definition.float :recovery_timeout
12+
definition.datetime :retired_at, precision: 6
13+
definition.timestamps precision: 6, null: false
14+
definition.check_constraint "recovery_timeout IS NULL OR recovery_timeout > 0", name: "chk_so_effect_recoveries_timeout"
15+
end
16+
end
17+
end

‎docs/effect-recovery.md‎

Lines changed: 203 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,203 @@
1+
# Effect recovery coordination
2+
3+
`emit` returns a JSON-serializable handle containing the public `effect_id`.
4+
Registering `on_recovery` opts the effect into retirement when its owner has
5+
stopped heartbeating. `on_status` is optional and receives responses only to
6+
explicit `request_effect_recovery(handle)` intents. Normal success and failure
7+
retain their existing callbacks.
8+
9+
## Watchdog using supported APIs
10+
11+
```ruby
12+
class ReportExport < SolidObjects::Actor
13+
attribute :revision, default: 0
14+
attribute :export_effect, default: nil
15+
attribute :artifact_key, default: ""
16+
attribute :applied_effect_id, default: nil
17+
18+
def start
19+
self.revision += 1
20+
self.export_effect = emit(:build_report,
21+
revision: revision,
22+
on_success: :export_finished,
23+
on_failure: :export_failed,
24+
on_recovery: :recover_export,
25+
on_status: :inspect_export,
26+
recovery_timeout: 120)
27+
schedule(at: Time.now + 30, key: "export-watchdog").watchdog
28+
nil
29+
end
30+
31+
def watchdog
32+
request_effect_recovery(export_effect) if export_effect
33+
end
34+
35+
def recover_export(effect_id:, arguments:, outcome:)
36+
return unless effect_id == export_effect&.fetch("effect_id")
37+
return unless arguments.fetch("revision") == revision
38+
39+
start
40+
end
41+
42+
def export_finished(effect_id:, arguments:, result:)
43+
apply_export_result(effect_id:, arguments:, result:)
44+
end
45+
46+
def export_failed(effect_id:, arguments:, error:)
47+
end
48+
49+
def inspect_export(effect_id:, outcome:, arguments: nil, result: nil)
50+
return unless effect_id == export_effect&.fetch("effect_id")
51+
52+
case outcome
53+
when SolidObjects::EffectRecoveryOutcome::COMPLETED
54+
apply_export_result(effect_id:, arguments:, result:)
55+
when SolidObjects::EffectRecoveryOutcome::DEFERRED, SolidObjects::EffectRecoveryOutcome::PENDING
56+
schedule(at: Time.now + 30, key: "export-watchdog").watchdog
57+
end
58+
end
59+
60+
private
61+
62+
def apply_export_result(effect_id:, arguments:, result:)
63+
return unless effect_id == export_effect&.fetch("effect_id")
64+
return unless arguments.fetch("revision") == revision
65+
return if applied_effect_id == effect_id
66+
67+
self.artifact_key = result.fetch("artifact_key")
68+
self.applied_effect_id = effect_id
69+
end
70+
end
71+
```
72+
73+
Register `build_report` through the ordinary effect registry. Its successful
74+
result in this example is `{ "artifact_key" => "reports/example.pdf" }`.
75+
The library builds the retirement payload, including `"outcome" => "retired"`;
76+
the effect handler does not return that outcome itself. Ruby actor operations
77+
receive keywords. `recover_export` needs handle/revision guards but no outcome
78+
guard because only a new retirement invokes it. The same guarded result helper
79+
handles success and completed-status repair, preventing duplicate application.
80+
Only recovery emits a replacement; status observations never do.
81+
82+
The watchdog is optional: `on_recovery` alone enables automatic retirement.
83+
`on_status` alone does not enable retirement, polling, or subscriptions. Explicit
84+
checks require both bindings persisted by `emit` and cannot replace either.
85+
86+
## Public envelopes and timeout
87+
88+
`SolidObjects::effect_handle` describes `{ "effect_id" => String }`.
89+
`SolidObjects::effect_retired_payload[Arguments]` requires the effect ID, original
90+
arguments, and `"outcome" => "retired"`. Status uses
91+
`SolidObjects::effect_recovery_payload[Arguments, Result]`, a record union.
92+
Every variant has an effect ID; retired and completed require original arguments;
93+
only completed has a recorded result (including `nil`). Other Ruby observations
94+
contain only the effect ID and outcome.
95+
96+
Strict packaged-consumer tests verify the constant literals and individual
97+
records. Steep 2.0 does not narrow this string-keyed record union after comparing
98+
`payload["outcome"]` with `COMPLETED`; accessing `result` through the union still
99+
fails its return-type check. Use the concrete completed/retired record in typed
100+
helpers after validating the discriminator, with an explicit type assertion if
101+
needed. The library retains precise records rather than weakening them to an
102+
untyped hash. Ordinary Ruby keyword dispatch needs no payload hydration.
103+
104+
| Frozen `SolidObjects::EffectRecoveryOutcome` constant | Wire value | Meaning |
105+
| --- | --- | --- |
106+
| `RETIRED` | `"retired"` | This check retired the effect; separate recovery owns replacement. |
107+
| `DEFERRED` | `"deferred"` | Fresh owner; preserve its claim and attempts. |
108+
| `PENDING` | `"pending"` | Initial execution or retry remains with the scheduler. |
109+
| `COMPLETED` | `"completed"` | Original arguments and recorded result are available. |
110+
| `DEAD` | `"dead"` | Preserve terminal failure and its existing callback. |
111+
| `ALREADY_RETIRED` | `"already_retired"` | Earlier retirement; no additional recovery notification. |
112+
| `MISSING` | `"missing"` | Owned binding exists but effect data was pruned. |
113+
114+
`recovery_timeout` must be positive finite seconds and requires `on_recovery`.
115+
Fractional durations are allowed. Omission uses the current runtime
116+
`process_alive_threshold`, normally 60 seconds. Smaller positive values are
117+
floored at that runtime threshold; changing configuration changes the effective
118+
floor even for existing effects. Database lookup errors surface as errors,
119+
never as missing/stale observations.
120+
121+
Effect workers maintain their process heartbeat while the handler waits on
122+
external I/O and while committing success or failure. A long-running healthy
123+
handler therefore remains protected beyond the recovery timeout. This requires
124+
an available database connection for the heartbeat, as well as runtime threads
125+
that can continue running.
126+
127+
Failed updates emit `solid_objects.process.heartbeat_failed` and retry at the
128+
configured heartbeat interval without consuming effect attempts. If an outage
129+
lasts beyond the freshness window, recovery can still be permitted; retries do
130+
not cancel external work or extend the configured window.
131+
132+
## Compatibility and installation
133+
134+
Upgrade all effect workers and process cleanup roles before emitting effects
135+
with recovery enabled. Older runtimes do not honor the persisted bindings or
136+
the new lock protocol.
137+
138+
Run `solid_objects:install:migrations` and your application's normal migration
139+
process before starting upgraded workers. The additive migration creates the
140+
durable binding table; it does not change existing effect status constraints.
141+
142+
`emit` now returns its handle, including without recovery options. Callers may
143+
ignore it. Wrappers must return `super`; operations whose last expression used
144+
to be `emit` may now return the handle to callers. End those operations with
145+
`nil` if their previous result must remain unchanged. The return-value change is
146+
intentional and is not strictly backward compatible.
147+
148+
## Transaction and lock protocol
149+
150+
Emission, actor state, the effect, and its recovery binding share the actor's
151+
fenced commit. An explicit check executes on that same connection. Automatic
152+
recovery performs one independent library transaction per candidate, with no
153+
application transaction waiting on a second connection.
154+
155+
The lock order is originating instance, effect rows ordered by public effect
156+
ID, recovery binding rows in the same order, then owner processes ordered by ID.
157+
Completion and failure must acquire the instance before the effect. Pending
158+
claims lock only their effect and do not subsequently acquire an instance lock.
159+
Mailbox insertion reuses the instance lock already held by the decision.
160+
Multiple checks in one actor commit lock all their effects and bindings before
161+
locking any processes. Unlocked candidate reads are hints, never decisions.
162+
163+
Automatic passes prefilter owner freshness and the effective per-effect timeout
164+
using database time, and visit at most `claim_scan_limit` stale candidates. Fresh
165+
actors and owners are not locked, including owners protected by extended grace.
166+
Remaining stale effects are revisited on later polls. Every candidate still
167+
undergoes the authoritative locked recheck, and successful retirement announces
168+
the committed mailbox work through the existing wake-up mechanism.
169+
170+
The decision samples database wall time after obtaining the owner lock. The
171+
effective timeout is the larger of the runtime's `process_alive_threshold` and
172+
the effect's persisted `recovery_timeout`, in seconds. A heartbeat newer than
173+
the cutoff is fresh; equality is stale. A stopped or draining process with
174+
fresh heartbeat evidence still protects an opted-in effect until that timeout.
175+
Cleanup preserves opted-in claims, and process pruning excludes processes that
176+
still own effects. Later effect polling and process cleanup revisit deferred
177+
effects without resetting their last heartbeat.
178+
179+
Retirement stores a durable `retired_at` in `effect_recoveries` and moves the
180+
effect into the existing terminal `completed` storage state, clearing its
181+
claim. The recovery record distinguishes retirement from successful completion;
182+
no success callback is generated. All recovery observations consult that record
183+
before interpreting the effect row. A late completion or failure is rejected by
184+
the existing processing/claim fence. This representation avoids rewriting the
185+
existing effect-status constraint across adapters.
186+
187+
The retirement record and the recovery mailbox message commit atomically. A
188+
winning explicit check additionally enqueues its status response after the
189+
retirement notification. Failure to insert either message rolls back the whole
190+
decision. Retirement is deduplicated per effect; check responses use separate
191+
per-request idempotency keys. Wake-up signals are delivery hints after commit.
192+
193+
Recovery bindings survive effect/message pruning and remain until the originating
194+
instance is destroyed or pruned. They do not prevent normal message or instance
195+
retention. Within that lifetime, a removed non-retired effect reports `missing`
196+
and a retirement record reports `already_retired`. A handle without an owned
197+
binding raises an error, without disclosing another actor's state or recreating
198+
a destroyed actor. Status-response message idempotency follows normal mailbox
199+
retention; callers cannot supply or reuse internal check request IDs.
200+
201+
An owner heartbeat measures process liveness, not effect progress. Retirement
202+
does not cancel the old handler or prove a remote request stopped. External
203+
actions still require idempotency across retries and replacement generations.

‎docs/roadmap.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,11 @@
1414
dead letters, and tail retry
1515
- Transactional effects with success/failure actor messages carrying the
1616
originally staged arguments for callback correlation
17+
- Stable `emit` handles and opted-in abandoned effect retirement coordinated
18+
with effect claims, with atomic recovery/status mailbox notifications and
19+
per-effect extending heartbeat timeouts and heartbeats throughout long-running
20+
effect handlers. See [effect recovery](effect-recovery.md)
21+
for the SQL lock protocol, retention boundary, and external idempotency limit.
1722
- Public RBS effect success/failure envelopes and error records, checked against
1823
the runtime constructors and a packaged consumer with strict Steep diagnostics
1924
- Actor-to-actor asynchronous outbox delivery. Effects and broadcasts use

‎examples/at_least_once/boot.rb‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ def self.call(database_path)
2323
require "solid_objects/database_adapter"
2424
%w[
2525
record process instance message ready_message claimed_message
26-
reminder effect broadcast dead_letter
26+
reminder effect effect_recovery broadcast dead_letter
2727
].each { |model| require File.join(ROOT, "app/models/solid_objects", model) }
2828

2929
SolidObjects.configuration.authorize_message = ->(**) { true }
@@ -40,8 +40,10 @@ def self.migrate
4040
require File.join(ROOT, "db/migrate/20260805000000_create_solid_objects_tables")
4141
require File.join(ROOT, "db/migrate/20260806000000_add_state_revision_to_solid_objects_instances")
4242
require File.join(ROOT, "db/migrate/20260813000000_rename_message_dispatch_columns")
43+
require File.join(ROOT, "db/migrate/20260915000000_add_solid_objects_effect_recoveries")
4344
CreateSolidObjectsTables.new.migrate(:up)
4445
AddStateRevisionToSolidObjectsInstances.new.migrate(:up)
4546
RenameMessageDispatchColumns.new.migrate(:up)
47+
AddSolidObjectsEffectRecoveries.new.migrate(:up)
4648
end
4749
end

0 commit comments

Comments
 (0)