Skip to content

Commit 15d2237

Browse files
committed
fix: Heartbeat while effect handlers run
1 parent 2fdf5d0 commit 15d2237

8 files changed

Lines changed: 175 additions & 1 deletion

File tree

‎CHANGELOG.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@
22

33
## 0.15.0 - 2026-09-15
44

5+
- Maintain effect-owner heartbeats during long-running handlers, completion, and
6+
failure handling, so healthy external I/O cannot trigger abandoned recovery.
57
- Return a stable effect handle from every `emit`. Wrappers must return it;
68
operations relying on an implicit `nil` result should return `nil` explicitly.
79
- Add abandoned effect recovery with `on_recovery`, optional `on_status`, staged

‎docs/effect-recovery.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,12 @@ floored at that runtime threshold; changing configuration changes the effective
118118
floor even for existing effects. Database lookup errors surface as errors,
119119
never as missing/stale observations.
120120

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+
121127
## Compatibility and installation
122128

123129
Upgrade all effect workers and process cleanup roles before emitting effects

‎docs/roadmap.md‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,8 @@
1616
originally staged arguments for callback correlation
1717
- Stable `emit` handles and opted-in abandoned effect retirement coordinated
1818
with effect claims, with atomic recovery/status mailbox notifications and
19-
per-effect extending heartbeat timeouts. See [effect recovery](effect-recovery.md)
19+
per-effect extending heartbeat timeouts and heartbeats throughout long-running
20+
effect handlers. See [effect recovery](effect-recovery.md)
2021
for the SQL lock protocol, retention boundary, and external idempotency limit.
2122
- Public RBS effect success/failure envelopes and error records, checked against
2223
the runtime constructors and a packaged consumer with strict Steep diagnostics

‎lib/solid_objects.rb‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@
5757
require "solid_objects/effect_registry"
5858
require "solid_objects/effect_payload"
5959
require "solid_objects/effect_recovery_coordinator"
60+
require "solid_objects/process_heartbeat"
6061
require "solid_objects/commit_action_registry"
6162
require "solid_objects/lease"
6263
require "solid_objects/lease_renewer"

‎lib/solid_objects/effect_executor.rb‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,12 +45,16 @@ def run_once
4545
effect = claim_next
4646
return false unless effect
4747

48+
heartbeat = ProcessHeartbeat.new(process_registry:)
49+
heartbeat.start
4850
result = deliver(effect)
4951
complete(effect, result)
5052
true
5153
rescue => error
5254
fail_effect(effect, error) if effect
5355
false
56+
ensure
57+
heartbeat&.stop
5458
end
5559

5660
# @rbs () -> void
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
# rbs_inline: enabled
2+
3+
module SolidObjects
4+
class ProcessHeartbeat
5+
# @rbs @process_registry: ProcessRegistry
6+
# @rbs @mutex: Thread::Mutex
7+
# @rbs @condition: Thread::ConditionVariable
8+
# @rbs @stopped: bool
9+
# @rbs @thread: Thread?
10+
11+
# @rbs (process_registry: ProcessRegistry) -> void
12+
def initialize(process_registry:)
13+
@process_registry = process_registry
14+
@mutex = Thread::Mutex.new
15+
@condition = Thread::ConditionVariable.new
16+
@stopped = false
17+
@thread = nil
18+
end
19+
20+
# @rbs () -> void
21+
def start
22+
@thread = Thread.new do
23+
Thread.current.report_on_exception = false
24+
loop do
25+
break if wait_for_interval
26+
27+
Record.connection_pool.with_connection { process_registry.heartbeat }
28+
end
29+
end
30+
end
31+
32+
# @rbs () -> void
33+
def stop
34+
mutex.synchronize do
35+
@stopped = true
36+
condition.broadcast
37+
end
38+
@thread&.join
39+
end
40+
41+
private
42+
43+
attr_reader :process_registry, :mutex, :condition
44+
45+
# @rbs () -> bool
46+
def wait_for_interval
47+
mutex.synchronize do
48+
condition.wait(mutex, SolidObjects.configuration.process_heartbeat_interval) unless @stopped
49+
@stopped
50+
end
51+
end
52+
end
53+
end
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
# Generated from lib/solid_objects/process_heartbeat.rb with RBS::Inline
2+
3+
module SolidObjects
4+
class ProcessHeartbeat
5+
@process_registry: ProcessRegistry
6+
7+
@mutex: Thread::Mutex
8+
9+
@condition: Thread::ConditionVariable
10+
11+
@stopped: bool
12+
13+
@thread: Thread?
14+
15+
# @rbs (process_registry: ProcessRegistry) -> void
16+
def initialize: (process_registry: ProcessRegistry) -> void
17+
18+
# @rbs () -> void
19+
def start: () -> void
20+
21+
# @rbs () -> void
22+
def stop: () -> void
23+
24+
private
25+
26+
attr_reader process_registry: untyped
27+
28+
attr_reader mutex: untyped
29+
30+
attr_reader condition: untyped
31+
32+
# @rbs () -> bool
33+
def wait_for_interval: () -> bool
34+
end
35+
end

‎test/integration/effect_recovery_test.rb‎

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,78 @@ def checked(effect_id:, outcome:, arguments: nil, result: nil)
6060
assert_empty actor.send(:drain_effect_intents)
6161
end
6262

63+
test "a running effect keeps its process heartbeat fresh" do
64+
SolidObjects.configuration.process_heartbeat_interval = 0.02
65+
SolidObjects.configuration.process_alive_threshold = 0.1
66+
ExportActor.ref("long-handler").async.start_recoverable_export
67+
worker = SolidObjects::Worker.new
68+
worker.run_until_idle
69+
entered = Queue.new
70+
release = Queue.new
71+
heartbeats = Queue.new
72+
executing = nil
73+
registry = SolidObjects::ProcessRegistry.new
74+
registry.define_singleton_method(:heartbeat) do
75+
updated = super()
76+
heartbeats << Thread.current if updated && executing
77+
updated
78+
end
79+
SolidObjects.register_effect(:build_report) do
80+
executing = Thread.current
81+
entered << true
82+
release.pop
83+
{ "artifact_key" => "report.pdf" }
84+
end
85+
executor = SolidObjects::EffectExecutor.new(process_registry: registry)
86+
effect_thread = Thread.new do
87+
SolidObjects::Record.connection_pool.with_connection { executor.run_once }
88+
end
89+
Timeout.timeout(5) { entered.pop }
90+
heartbeat_thread = Timeout.timeout(2) { heartbeats.pop }
91+
Timeout.timeout(2) { 7.times { heartbeats.pop } }
92+
SolidObjects::EffectRecoveryCoordinator.new.recover_available
93+
assert_equal "processing", SolidObjects::Effect.find_by!(name: "build_report").status
94+
assert_empty SolidObjects::Message.where(operation: "retired")
95+
release << true
96+
assert effect_thread.value
97+
refute heartbeat_thread.alive?
98+
assert_equal "completed", SolidObjects::Effect.find_by!(name: "build_report").status
99+
ensure
100+
release&.push(true)
101+
effect_thread&.join(5)
102+
executor&.stop
103+
worker&.stop
104+
end
105+
106+
test "a failed handler stops its heartbeat before scheduling a retry" do
107+
SolidObjects.configuration.process_heartbeat_interval = 0.01
108+
ExportActor.ref("failing-handler").async.start_recoverable_export
109+
worker = SolidObjects::Worker.new
110+
worker.run_until_idle
111+
heartbeats = Queue.new
112+
registry = SolidObjects::ProcessRegistry.new
113+
registry.define_singleton_method(:heartbeat) do
114+
updated = super()
115+
heartbeats << Thread.current if updated
116+
updated
117+
end
118+
heartbeat_thread = nil
119+
SolidObjects.register_effect(:build_report) do
120+
heartbeats.pop(true) until heartbeats.empty?
121+
heartbeat_thread = Timeout.timeout(2) { heartbeats.pop }
122+
raise "remote request failed"
123+
end
124+
executor = SolidObjects::EffectExecutor.new(process_registry: registry)
125+
refute executor.run_once
126+
refute_nil heartbeat_thread
127+
refute heartbeat_thread.alive?
128+
assert_equal "pending", SolidObjects::Effect.find_by!(name: "build_report").status
129+
assert_empty SolidObjects::Message.where(operation: "retired")
130+
ensure
131+
executor&.stop
132+
worker&.stop
133+
end
134+
63135
test "a longer recovery timeout survives ordinary process cleanup" do
64136
ExportActor.ref("long-grace").async.start_with_timeout(timeout: 120)
65137
worker = SolidObjects::Worker.new

0 commit comments

Comments
 (0)