feat(reborn): add concrete TurnRunner worker composition - #3457
serrrfirat wants to merge 5 commits into
Conversation
There was a problem hiding this comment.
Code Review
This pull request implements the TurnRunnerWorker and LoopExitApplier components to manage turn run execution, including task claiming, lease heartbeating, and exit validation. It also introduces a BlockedProcess status and a require_final_checkpoint policy. Feedback suggests optimizing the worker loop to drain the work queue and configuring the heartbeat timer to skip missed ticks to prevent request bursts.
| if let Err(err) = self.try_claim_and_run().await { | ||
| warn!( | ||
| runner_id = ?self.runner_id, | ||
| error = %err, | ||
| "claim-and-run cycle failed" | ||
| ); | ||
| } |
There was a problem hiding this comment.
The worker currently claims and executes only one run per wake signal or poll interval. If multiple runs are queued, the worker will wait for the next poll interval (or another wake signal) before processing the next run. It is more efficient to drain the queue by looping until no more runs are available to claim. This requires updating try_claim_and_run to return a boolean indicating whether a run was claimed.
while !cancel.is_cancelled() {
match self.try_claim_and_run().await {
Ok(true) => continue,
Ok(false) => break,
Err(err) => {
warn!(
runner_id = ?self.runner_id,
error = %err,
"claim-and-run cycle failed"
);
break;
}
}
}References
- Errors should be surfaced via tracing::warn! to ensure reliability in background tasks without failing the entire operation.
| } | ||
|
|
||
| /// Attempt one claim-and-run cycle. | ||
| async fn try_claim_and_run(&self) -> Result<(), TurnRunnerError> { |
There was a problem hiding this comment.
Update the signature of try_claim_and_run to return Result<bool, TurnRunnerError> to support draining the work queue in the main loop.
async fn try_claim_and_run(&self) -> Result<bool, TurnRunnerError> {References
- Create specific error variants for different failure modes (e.g., TurnRunnerError) to provide semantically correct and clear error messages.
|
|
||
| let Some(claimed) = claimed else { | ||
| debug!(runner_id = ?self.runner_id, "no runs available to claim"); | ||
| return Ok(()); |
| ); | ||
|
|
||
| self.execute_claimed_run(claimed).await; | ||
| Ok(()) |
| interval: Duration, | ||
| cancel: CancellationToken, | ||
| ) -> Result<(), TurnError> { | ||
| let mut tick = tokio::time::interval(interval); |
There was a problem hiding this comment.
The heartbeat loop uses the default Burst behavior for tokio::time::interval. If the worker or runtime is under heavy load and ticks are missed, this can lead to a burst of heartbeat requests being sent once the task resumes. For heartbeats, it is generally safer to use MissedTickBehavior::Skip to maintain a steady rate without catching up on missed intervals.
let mut tick = tokio::time::interval(interval);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);References
- Adjusting timing behavior (like MissedTickBehavior) is a pragmatic tradeoff for resource management in non-critical paths like heartbeats.
Summary
TurnRunnerWorkercomposition on a clean branch fromorigin/reborn-integrationTurnRunnerIdand fresh per-claimTurnLeaseTokenLoopExitapplicationLoopExitApplierand recovery mapping for missing drivers, host failures, driver errors/panics, heartbeat lease loss, cancellation, and exit-application failuresRebornLoopDriverHostFactoryinto the runnerHostFactoryseamCloses #3404
Supersedes #3446 with a clean branch directly from
reborn-integration.Tests
cargo fmt --all -- --checkCARGO_TARGET_DIR=/tmp/ironclaw-target-3404-clean cargo test -p ironclaw_rebornCARGO_TARGET_DIR=/tmp/ironclaw-target-3404-clean cargo clippy -p ironclaw_reborn --all-targets --all-features -- -D warnings