Preserve committed outcomes during host failure cleanup

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-07 15:53:08 -04:00
parent 21d4db6766
commit 59f9088b04
10 changed files with 227 additions and 76 deletions

View file

@ -1121,6 +1121,7 @@ mod tests {
let server = MockServer::start_async().await;
let mut state = terminal_run_state_response(run_id);
state["status"] = serde_json::json!({"kind": "running"});
state["conclusion"] = serde_json::Value::Null;
let state: server_client::RunProjection = serde_json::from_value(state).unwrap();
let executed = serde_json::json!({
"run_id": run_id, "stream_seq": 1, "kind": "petri", "id": "executed",
@ -1146,6 +1147,14 @@ mod tests {
.json_body(serde_json::to_value(&state).unwrap());
})
.await;
server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/questions"));
then.status(200)
.json_body(serde_json::json!({ "data": [], "meta": { "has_more": false } }));
})
.await;
let waiting = server
.mock_async(|when, then| {
when.method("GET")
@ -1157,7 +1166,7 @@ mod tests {
.await;
let client = server_client::Client::new_no_proxy(&server.base_url()).unwrap();
let task = tokio::spawn(async move {
attach_petri_run_with_client(
Box::pin(attach_petri_run_with_client(
&client,
&run_id,
&state,
@ -1169,7 +1178,7 @@ mod tests {
json_output: true,
},
Printer::Default,
)
))
.await
.unwrap()
});
@ -1184,7 +1193,6 @@ mod tests {
!task.is_finished(),
"successful execution does not end attach while publication is pending"
);
waiting.delete_async().await;
let finished = serde_json::json!({
"run_id": run_id, "stream_seq": 2, "kind": "petri", "id": "finished",
"recorded_at": 2000,
@ -1201,6 +1209,7 @@ mod tests {
.body(format!("data: {finished}\n\n"));
})
.await;
waiting.delete_async().await;
assert_eq!(
tokio::time::timeout(Duration::from_secs(5), task)
.await

View file

@ -1307,10 +1307,10 @@ mod tests {
assert_eq!(exit_code_of(&finished), Some(1));
let mut state = PrettyState::default();
let line = format_pretty(&finished, &Styles::new(false), &mut state).unwrap();
insta::assert_snapshot!(line, @r###"
12:43:08 ✗ FAILED 0s
12:43:08 the push was rejected
"###);
insta::assert_snapshot!(line, @"
04:43:08 ✗ FAILED 0ms
04:43:08 the push was rejected
");
}
#[test]

View file

@ -976,6 +976,7 @@ mod tests {
expected: RunStatus,
message: Option<&str>,
) {
state.petri_projector.settle(run_id).await;
let mut values = Vec::new();
for suffix in ["", "/state"] {
let response = app
@ -999,7 +1000,9 @@ mod tests {
}
assert_eq!(
values[0]["lifecycle"]["status"],
serde_json::to_value(expected).unwrap()
serde_json::to_value(expected).unwrap(),
"public state: {:#?}",
values[1]
);
assert_eq!(values[1]["status"], serde_json::to_value(expected).unwrap());
assert_eq!(state.test_managed_run_status(&run_id), Some(expected));
@ -1018,6 +1021,51 @@ mod tests {
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_rejected_platform_failure_does_not_settle_the_run() {
let runtime = Arc::new(HeldWorkerRuntime::default());
let (state, app, run_id, token) = held_worker_run(&runtime).await;
run_to_running_as_worker(&app, run_id, &token).await;
let pool = state.stores.run_summaries.pool();
sqlx::query(
"CREATE TRIGGER reject_platform_record BEFORE INSERT ON platform_records \
BEGIN SELECT RAISE(FAIL, 'scripted append failure'); END",
)
.execute(&pool)
.await
.unwrap();
let response = append_lifecycle_as_worker(
&app,
run_id,
&token,
RunLifecycleKind::Failed,
RunStatus::Failed {
reason: FailureReason::LaunchFailed,
},
)
.await;
fabro_test::assert_axum_status(
response,
StatusCode::INTERNAL_SERVER_ERROR,
"rejected platform terminal append",
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
crate::server::persist_run_failure(
&state,
run_id,
FailureReason::LaunchFailed,
"worker launch failed".into(),
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
sqlx::query("DROP TRIGGER reject_platform_record")
.execute(&pool)
.await
.unwrap();
runtime.end_worker();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_publication_result_agrees_across_worker_api_projection_and_cleanup() {
use fabro_petri::petri::LogId;
@ -1109,10 +1157,36 @@ mod tests {
}
assert_eq!(concluded_runs(&app).await, 1);
assert_public_result(&state, &app, run_id, expected, rejection).await;
let (rebuilt, _, _) =
fabro_petri::test_support::rebuild(&state.db_pool, &state.db_pool, run_id)
// A host error arriving after cleanup must neither append a
// competing terminal record nor replace the useful failure.
let platform_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM platform_records WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_one(&state.stores.run_summaries.pool())
.await
.unwrap();
crate::server::persist_run_failure(
&state,
run_id,
FailureReason::Terminated,
"worker wait failed during teardown".into(),
)
.await;
assert_public_result(&state, &app, run_id, expected, rejection).await;
let after_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM platform_records WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_one(&state.stores.run_summaries.pool())
.await
.unwrap();
assert_eq!(after_count, platform_count);
let (rebuilt, _, _) = fabro_petri::test_support::rebuild(
&state.db_pool,
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap();
let rebuilt = rebuilt.unwrap();
assert_eq!(rebuilt.status, expected);
assert_eq!(

View file

@ -3515,13 +3515,70 @@ async fn fail_run_before_execution(
reason: FailureReason,
message: String,
) {
if let Err(err) =
run_records::lifecycle(state, run_id, run_records::failed(reason, message.clone())).await
{
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
}
persist_run_failure(state, run_id, reason, message).await;
}
fail_managed_run(state, run_id, reason, message);
/// Record a host failure only while no terminal result is committed. This
/// also handles a worker wait/launch error racing its durable Petri finish.
pub(crate) async fn persist_run_failure(
state: &Arc<AppState>,
run_id: RunId,
reason: FailureReason,
message: String,
) {
match run_records::projection(state, run_id).await {
Ok(Some(projection)) if projection.status.is_terminal() => {
let failure = projection
.conclusion
.as_ref()
.and_then(|conclusion| conclusion.failure.as_ref())
.map(|failure| failure.detail.message.clone());
settle_managed_run_at_finish(state, run_id, projection.status, failure);
}
Ok(Some(_)) => {
match run_records::lifecycle(
state,
run_id,
run_records::failed(reason, message.clone()),
)
.await
{
Ok(_) => match run_records::projection(state, run_id).await {
Ok(Some(committed)) if committed.status.is_terminal() => {
let failure = committed
.conclusion
.as_ref()
.and_then(|conclusion| conclusion.failure.as_ref())
.map(|failure| failure.detail.message.clone());
settle_managed_run_at_finish(state, run_id, committed.status, failure);
}
Ok(_) => {
error!(run_id = %run_id, "Stored host failure has no terminal projection");
}
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to read the committed host failure");
}
},
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
}
}
}
Ok(None) => {
error!(run_id = %run_id, "Run missing when recording a host failure");
}
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to read the committed run result");
}
}
// Resource cleanup is operational; it cannot substitute for a committed
// outcome if storage is unavailable.
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
clear_live_run_state(managed_run);
}
drop(runs);
cleanup_worker_control_bus_for_run(state, run_id);
state.scheduler_notify.notify_one();
}
@ -3777,26 +3834,13 @@ async fn fail_worker_launch(state: &Arc<AppState>, run_id: RunId, err: anyhow::E
None
}
};
let launch_message = format!("Failed to spawn worker: {err}");
let (error, reason) = failure_honoring_pending_cancel(pending_control, || {
(
WorkflowError::engine_with_anyhow("Failed to spawn worker", err),
FailureReason::LaunchFailed,
)
});
let message = if reason == FailureReason::Cancelled {
"Run cancelled before worker launch completed".to_string()
} else {
launch_message
};
let _ = run_records::lifecycle(
state,
run_id,
run_records::failed(reason, error.to_string()),
)
.await;
fail_managed_run(state, run_id, reason, message);
state.scheduler_notify.notify_one();
persist_run_failure(state, run_id, reason, error.to_string()).await;
}
/// A worker that exited without recording the run's end left it failed,
@ -4257,15 +4301,17 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
Err(err) => {
tracing::error!(run_id = %run_id, error = %err, "Failed while waiting on worker");
let message = format!("Worker wait failed: {err}");
let superseded = {
let runs = state.runs.lock().expect("runs lock poisoned");
runs.get(&run_id)
.is_some_and(|run| run.worker_ref.as_ref() != Some(&worker_ref))
};
if superseded {
return;
}
state.worker_runtime.force_stop(&worker_ref).await;
let _ = run_records::lifecycle(
&state,
run_id,
run_records::failed(FailureReason::Terminated, message.clone()),
)
.await;
fail_managed_run(&state, run_id, FailureReason::Terminated, message);
state.scheduler_notify.notify_one();
state.petri_runs.worker_exited(run_id);
persist_run_failure(&state, run_id, FailureReason::Terminated, message).await;
return;
}
};

View file

@ -496,8 +496,8 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
run_id: run_id.to_string(),
run_dir: run_dir.join("petri"),
execution,
// The projector's signal follows each durable append; the managed
// run's settle at Petri's finish precedes it.
// The coordinator finish is stored before managed status settles;
// the projector also reads only durable records.
store: Arc::new(SettlingStore {
inner: state
.petri_projector
@ -550,7 +550,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
}
};
match run_records::lifecycle(&state, run_id, record).await {
Ok(()) => finish(&state, run_id, status, error),
Ok(_) => finish(&state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
release_live_state(&state, run_id);
@ -815,7 +815,7 @@ async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &s
error!(run_id = %run_id, error = message, "Petri run cannot start");
let (status, error, record) = failed(FailureReason::WorkflowError, message.to_string());
match run_records::lifecycle(state, run_id, record).await {
Ok(()) => finish(state, run_id, status, error),
Ok(_) => finish(state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
release_live_state(state, run_id);
@ -825,9 +825,9 @@ async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &s
/// Settle the managed run at its terminal record and release its
/// scheduler slot. A run that Petri finished settled already, at the
/// `run.finished` record ([`SettlingStore`]); this refines its status and
/// error and ends its live state. A run deleted since is gone from the map
/// and stays gone.
/// `run.finished` record ([`SettlingStore`]); this preserves that status,
/// fills any missing failure detail and ends its live state. A run deleted
/// since is gone from the map and stays gone.
fn finish(state: &Arc<AppState>, run_id: RunId, status: RunStatus, error: Option<String>) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {

View file

@ -150,7 +150,7 @@ pub enum RunStatus {
#[derive(Clone, Debug)]
pub struct RunOutcome {
pub status: RunStatus,
/// The root invocation's failure message, when it failed.
/// Required-finalization failure detail, or the root execution's failure.
pub failure: Option<String>,
/// Whether the record is whole: the run recorded its finish and every
/// log replays byte for byte.
@ -409,8 +409,9 @@ pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result<RunOutcome
/// How Fabro reports what [`run`] returned. A cancelled run is a failure
/// with the cancelled reason, as the legacy executor reports one; a failed
/// store interrupted the run; every other shortfall is a workflow error
/// whose message says what the record, or the host, said.
/// store interrupted the run; a publication rejection keeps its
/// publish_failed reason. Other shortfalls are workflow errors whose message
/// says what the record, or the host, said.
#[must_use]
pub fn conclusion(result: &Result<RunOutcome, RunError>) -> Conclusion {
match result {

View file

@ -53,16 +53,18 @@
//!
//! # Operation identities
//!
//! Every external effect here is keyed on `(run key, execution, DecisionId,
//! effect kind)` from the hook context and deduplicated on retry: the
//! checkpoint's key is the attempt's decision in its execution, effect
//! Checkpoint and artifact effects are keyed on `(run key, execution,
//! DecisionId, effect kind)` from the hook context and deduplicated on retry:
//! the checkpoint's key is the attempt's decision in its execution, effect
//! `checkpoint`; an artifact's is the same decision, effect `artifact`, with
//! the file's path and content digest as the identity within it. A
//! re-dispatched attempt whose commit already landed reuses it when the
//! workspace still sits on it unchanged (see [`RunWorkspaces::commit`]); a
//! reissued routing decision finds the record, or the commit by its
//! trailers, and writes nothing twice; a file already collected under the
//! same path and digest is not collected again.
//! same path and digest is not collected again. Publication keeps the
//! publisher's reconciliation policy; required finalization adds no independent
//! effect ledger or guarantee of deduplication across every external crash.
//!
//! # Where the workspace is
//!

View file

@ -10,6 +10,7 @@ use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::FinalizationFailure;
use petri_runtime::{RunOptions, Runtime};
use petri_store::{Access, LogId, MemoryRunStore, Record, RunKey, RunStore as _};
use tokio::fs;
use tokio::sync::{Notify, Semaphore};
use crate::providers::{self, SandboxProviderConfig};
@ -76,7 +77,7 @@ pub async fn test_run_records(
) -> TestRunRecords {
let root = tempfile::tempdir().expect("the fixture has an isolated directory");
let workflow = root.path().join("workflow.fabro");
tokio::fs::write(
fs::write(
&workflow,
r#"digraph Finalization {
graph [goal="Check required publication"]
@ -88,7 +89,7 @@ pub async fn test_run_records(
)
.await
.expect("the fixture workflow writes");
tokio::fs::write(
fs::write(
root.path().join("workflow.toml"),
"_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n",
)
@ -129,7 +130,7 @@ pub async fn test_run_records(
)];
for execution in inspection.executions {
let id = LogId::Execution(execution.execution);
records.push((id.clone(), logs.read(&id).await.expect("execution reads")));
records.push((id, logs.read(&id).await.expect("execution reads")));
}
let mut blobs = Vec::new();
for graph in inspection.graphs {

View file

@ -44,8 +44,9 @@ use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind};
use object_store::local::LocalFileSystem;
use petri_execution::inspect::{self, RunInspection};
use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _};
use tokio::fs;
use tokio::process::Command;
use tokio::sync::{Notify, Semaphore};
use tokio::{fs, time};
use tokio_util::sync::CancellationToken;
mod support;
@ -1609,8 +1610,8 @@ async fn a_run_diff_failure_cannot_silently_skip_publication() {
struct GatedPublisher {
inner: Arc<RecordingPublisher>,
entered: tokio::sync::Notify,
release: tokio::sync::Semaphore,
entered: Notify,
release: Semaphore,
}
#[async_trait::async_trait]
@ -1621,7 +1622,11 @@ impl RunPublisher for GatedPublisher {
async fn publish(&self, publication: &Publication) -> Result<(), String> {
self.entered.notify_one();
self.release.acquire().await.unwrap().forget();
self.release
.acquire()
.await
.expect("publication gate stays open")
.forget();
self.inner.publish(publication).await
}
}
@ -1631,8 +1636,8 @@ async fn required_publication_blocks_the_terminal_result_and_cleanup() {
for rejection in [None, Some("the push was rejected")] {
let publisher = Arc::new(GatedPublisher {
inner: RecordingPublisher::new(rejection),
entered: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
entered: Notify::new(),
release: Semaphore::new(0),
});
let mut harness = Harness::new().await;
harness.publisher = Some(publisher.clone());
@ -1649,7 +1654,7 @@ async fn required_publication_blocks_the_terminal_result_and_cleanup() {
)
.await
});
tokio::time::timeout(
time::timeout(
std::time::Duration::from_secs(15),
publisher.entered.notified(),
)
@ -1674,10 +1679,6 @@ async fn required_publication_blocks_the_terminal_result_and_cleanup() {
.all(|record| record.record["body"]["event"] != "scope.released"),
"scope cleanup waits for publication"
);
assert!(
pending.invocations[0].result.is_some(),
"execution already ended"
);
assert!(
harness.workspace_path(&harness.workspace().await).exists(),
"workspace is available to publication"

View file

@ -27,7 +27,8 @@ use fabro_petri::providers::SandboxProviderConfig;
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::{SqliteRunStore, providers, test_support as petri_support};
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
CheckpointRecord, PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunDiffRecord,
RunLifecycleKind, RunLifecycleRecord,
};
use fabro_store::{BlobStore, test_support};
use fabro_types::{
@ -42,7 +43,7 @@ use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::RunStatus as PetriRunStatus;
use petri_store::{RunKey, RunStore};
use tokio::fs;
use tokio::time::sleep;
use tokio::time::{self, sleep};
const COMMAND_WORKFLOW: &str = r#"digraph Command {
graph [goal="Run one command"]
@ -1387,14 +1388,30 @@ async fn required_finalization_projects_only_the_committed_overall_result() {
.await
.unwrap()
});
tokio::time::timeout(Duration::from_secs(15), finalizer.entered.notified())
time::timeout(Duration::from_secs(15), finalizer.entered.notified())
.await
.unwrap();
projector.settle(scenario.run_id).await;
let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
// The driver can still flush its execution journal after entering
// finalization. Wait for the exit-stage evidence, rather than racing
// that writer while comparing the live view with a full rebuild.
let pending = time::timeout(Duration::from_secs(5), async {
loop {
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
if pending.iter_stages().any(|(id, stage)| {
id.node_id() == "exit" && stage.state == StageState::Succeeded
}) {
break pending;
}
time::sleep(Duration::from_millis(5)).await;
}
})
.await
.expect("execution evidence flushes while publication is held");
assert_eq!(pending.status, RunStatus::Running);
assert!(pending.conclusion.is_none());
assert!(!task.is_finished());
@ -1408,7 +1425,7 @@ async fn required_finalization_projects_only_the_committed_overall_result() {
let patch_blob = BlobHash::new(b"final patch");
let platform = PlatformRecordStore::new(scenario.pool.clone());
for record in [
PlatformRecord::Checkpoint(fabro_store::platform_records::CheckpointRecord {
PlatformRecord::Checkpoint(CheckpointRecord {
execution: 0,
firing: 0,
attempt: Some(1),
@ -1418,7 +1435,7 @@ async fn required_finalization_projects_only_the_committed_overall_result() {
patch_blob: Some(patch_blob),
operation: None,
}),
PlatformRecord::RunDiff(fabro_store::platform_records::RunDiffRecord {
PlatformRecord::RunDiff(RunDiffRecord {
base_sha: None,
head_sha: Some(head_sha.to_string()),
diff_summary: Some(summary),