Adapt recovery to current resume callers and API errors

This commit is contained in:
Scott Werner 2026-10-06 12:01:01 -04:00
parent eb39e15781
commit f12065e6ae
5 changed files with 18 additions and 13 deletions

View file

@ -104,4 +104,7 @@ same transaction as the projection that consumed the record. A client that
resumes from its last `stream_seq` sees every item exactly once.
A worker cannot continue past a record it failed to append: the store's
error reaches the engine and fails the run.
error reaches the engine and interrupts the run's lifetime. The worker
records no terminal lifecycle transition for that interruption. The server
resumes the run from its durable records, up to three times per server
process; a fourth interruption fails the run.

View file

@ -2793,7 +2793,13 @@ async fn delete_run_internal(
// outlived a server crash, is stopped first.
petri_runs::stop_previous_worker(state, id)
.await
.map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, format!("{err:#}")))?;
.map_err(|err| {
error!(run_id = %id, error = %format!("{err:#}"), "Stopping the run's previous worker failed");
ApiError::new(
StatusCode::INTERNAL_SERVER_ERROR,
"failed to stop the run's previous worker",
)
})?;
state.petri_runs.worker_exited(id);
let delete_outcome = delete_run_sandbox_resource(state, id, force).await?;

View file

@ -580,7 +580,7 @@ pub(crate) async fn reconcile_on_startup(
run_id: RunId,
run_state: &fabro_store::RunProjection,
) -> anyhow::Result<()> {
let mode = match relaunch(state, run_id, run_state).await? {
let mode = match relaunch(state, run_id).await? {
Relaunch::Worker(mode) => mode,
Relaunch::Failed { reason } => {
warn!(
@ -660,7 +660,7 @@ pub(crate) async fn resume_after_interruption(
{
return Err(interrupted);
}
let mode = match relaunch(state, run_id, &run_state).await {
let mode = match relaunch(state, run_id).await {
Ok(Relaunch::Worker(mode)) => mode,
Ok(Relaunch::Failed { reason }) => return Err(reason),
Err(err) => {
@ -713,11 +713,7 @@ enum Relaunch {
/// run, else in start mode: a worker that died before it created the run's
/// record left nothing to continue from, so the run starts from its
/// admitted graphs. The caller registers the run with the scheduler.
async fn relaunch(
state: &Arc<AppState>,
run_id: RunId,
run_state: &fabro_store::RunProjection,
) -> anyhow::Result<Relaunch> {
async fn relaunch(state: &Arc<AppState>, run_id: RunId) -> anyhow::Result<Relaunch> {
stop_previous_worker(state, run_id).await?;
let held = match state.petri_runs.release_for_restart(run_id).await {
Ok(()) => true,

View file

@ -191,7 +191,7 @@ impl Harness {
run_id: self.run_id.to_string(),
run_dir: self.run_dir.clone(),
execution: if resumed {
Execution::Resume
Execution::Resume(admit(workflow, settings))
} else {
Execution::Start(admit(workflow, settings))
},

View file

@ -18,7 +18,7 @@ use std::time::{Duration, Instant};
use tokio::fs::OpenOptions;
use tokio::time;
use crate::signal::{sigkill, sigkill_process_group};
use crate::signal;
/// How often [`stop_lock_holder`] checks whether the lock is free.
const POLL: Duration = Duration::from_millis(20);
@ -94,8 +94,8 @@ pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result<Loc
let Some(pid) = holder(&file)? else {
return Ok(LockHolder::None);
};
sigkill_process_group(pid);
sigkill(pid);
signal::sigkill_process_group(pid);
signal::sigkill(pid);
let deadline = Instant::now() + patience;
loop {
match holder(&file)? {