diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index 4be756797..f6c404911 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -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. diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index a21ce3692..9f0024990 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -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?; diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 2d3201c35..31f4980ca 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, - run_id: RunId, - run_state: &fabro_store::RunProjection, -) -> anyhow::Result { +async fn relaunch(state: &Arc, run_id: RunId) -> anyhow::Result { stop_previous_worker(state, run_id).await?; let held = match state.petri_runs.release_for_restart(run_id).await { Ok(()) => true, diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 0a118a3a9..d548f3a79 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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)) }, diff --git a/lib/foundation/fabro-proc/src/process_lock.rs b/lib/foundation/fabro-proc/src/process_lock.rs index ceaadc8d0..3b1873dc6 100644 --- a/lib/foundation/fabro-proc/src/process_lock.rs +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -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