Simplify required run finalization

Name the publish failure once, settle host failures through one flat
helper, share the append-then-finish path for in-process runs, reuse
the source state already read when checking a fork, and drop the
unused finished-run status shim.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-08 12:42:35 -04:00
parent 599b674ec1
commit af41c7c66c
9 changed files with 105 additions and 104 deletions

View file

@ -1170,12 +1170,13 @@ mod tests {
assert_public_result(&state, &app, run_id, expected, rejection).await;
// 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();
let platform_count = fabro_petri::test_support::stored_platform_records(
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap()
.len();
crate::server::persist_run_failure(
&state,
run_id,
@ -1184,12 +1185,13 @@ mod tests {
)
.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();
let after_count = fabro_petri::test_support::stored_platform_records(
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap()
.len();
assert_eq!(after_count, platform_count);
let (rebuilt, _, _) = fabro_petri::test_support::rebuild(
&state.db_pool,

View file

@ -3505,19 +3505,10 @@ async fn reject_run_if_sandbox_provider_disabled(
return false;
};
tracing::warn!(run_id = %run_id, error = %error, "Sandbox provider disabled by server policy");
fail_run_before_execution(state, run_id, FailureReason::LaunchFailed, error).await;
persist_run_failure(state, run_id, FailureReason::LaunchFailed, error).await;
true
}
async fn fail_run_before_execution(
state: &Arc<AppState>,
run_id: RunId,
reason: FailureReason,
message: String,
) {
persist_run_failure(state, run_id, reason, message).await;
}
/// 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(
@ -3526,50 +3517,8 @@ pub(crate) async fn persist_run_failure(
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");
}
if let Err(err) = commit_host_failure(state, run_id, reason, message).await {
error!(run_id = %run_id, error = %err, "Failed to record a host failure");
}
// Resource cleanup is operational; it cannot substitute for a committed
// outcome if storage is unavailable.
@ -3582,6 +3531,41 @@ pub(crate) async fn persist_run_failure(
state.scheduler_notify.notify_one();
}
/// Append the host failure unless the run already ended, then settle the
/// managed run on whichever terminal result the store committed.
async fn commit_host_failure(
state: &AppState,
run_id: RunId,
reason: FailureReason,
message: String,
) -> anyhow::Result<()> {
let mut committed = run_records::projection(state, run_id)
.await?
.context("the run is missing")?;
if !committed.status.is_terminal() {
run_records::lifecycle(state, run_id, run_records::failed(reason, message)).await?;
// The append waited for the projector, so the stored projection
// already folds it, or the finish that won the race.
committed = state
.stores
.run_summaries
.load_petri_projection(&run_id)
.await?
.context("the run is missing")?;
anyhow::ensure!(
committed.status.is_terminal(),
"the stored host failure has no terminal projection"
);
}
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(())
}
fn managed_run(
dot_source: String,
status: RunStatus,
@ -4241,7 +4225,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
let github_app_private_key = match state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await {
Ok(value) => value,
Err(err) => {
fail_run_before_execution(
persist_run_failure(
&state,
run_id,
FailureReason::WorkflowError,

View file

@ -549,13 +549,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
failed(reason, message)
}
};
match run_records::lifecycle(&state, run_id, record).await {
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);
}
}
commit_and_finish(&state, run_id, record, status, error).await;
// The view trails the terminal record; the aggregate reads the settled
// projection, as the worker path reads the final state at worker exit.
state.petri_projector.settle(run_id).await;
@ -814,10 +808,22 @@ fn failed(
async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &str) {
error!(run_id = %run_id, error = message, "Petri run cannot start");
let (status, error, record) = failed(FailureReason::WorkflowError, message.to_string());
commit_and_finish(state, run_id, record, status, error).await;
}
/// Append the run's terminal record, then finish the run. An append that
/// fails only releases the live state: it commits no terminal result.
async fn commit_and_finish(
state: &Arc<AppState>,
run_id: RunId,
record: RunLifecycleRecord,
status: RunStatus,
error: Option<String>,
) {
match run_records::lifecycle(state, run_id, record).await {
Ok(_) => finish(state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
release_live_state(state, run_id);
}
}
@ -837,10 +843,9 @@ fn finish(state: &Arc<AppState>, run_id: RunId, status: RunStatus, error: Option
managed_run.error = error;
}
}
clear_live_run_state(managed_run);
}
drop(runs);
state.scheduler_notify.notify_one();
release_live_state(state, run_id);
}
/// Release controls after an append failure without claiming a new terminal

View file

@ -84,6 +84,7 @@ use crate::admission::AdmittedGraphs;
use crate::blobs::{Blobs, RunBlobs};
use crate::controls::RunControls;
use crate::hooks::{FabroHooks, HooksSpec};
use crate::projection;
use crate::runtime::RuntimeSpec;
use crate::secrets::SharedSecrets;
@ -556,7 +557,7 @@ fn outcome(
let publish_failed = inspection
.finalization_failure
.as_ref()
.is_some_and(|failure| failure.code == "publish_failed");
.is_some_and(projection::is_publish_failure);
let failure = inspection
.finalization_failure
.map(|failure| failure.message)

View file

@ -39,7 +39,7 @@ use fabro_workflow::operations::{StageLabel, StageLabels};
use petri_execution::host::{self, ForkOptions, ForkOrigin, ForkPosition, HostError};
use petri_execution::inspect::{self, InspectError};
use petri_execution::{
Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunStore,
Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunLogs, RunStore,
StoreError as CoordinatorStoreError,
};
use petri_runtime::RunOptions;
@ -132,7 +132,13 @@ pub async fn check(
.open(&RunKey::new(source.to_string()), Access::Read)
.await
.map_err(ForkError::Open)?;
let state = host::stored_state(&*logs).await.map_err(ForkError::Seed)?;
check_logs(&*logs, &position).await.map(|_| ())
}
/// [`check`] over the source's opened logs. Returns whether the source
/// requires run finalization, which the fork inherits.
async fn check_logs(logs: &dyn RunLogs, position: &ForkPosition) -> Result<bool, ForkError> {
let state = host::stored_state(logs).await.map_err(ForkError::Seed)?;
let Some(execution) = state.executions.get(&position.execution) else {
return Err(ForkError::Refused(format!(
"the source run has no execution {}",
@ -147,7 +153,7 @@ pub async fn check(
position.execution
)));
}
let inspection = inspect::inspect_run(&*logs)
let inspection = inspect::inspect_run(logs)
.await
.map_err(ForkError::Inspect)?;
if inspection
@ -167,7 +173,7 @@ pub async fn check(
"the terminal checkpoint has no remaining work to acquire a sandbox; select an earlier checkpoint or retry the workflow from the start".to_string(),
));
}
Ok(())
Ok(state.required_finalization)
}
/// Declaration-only hooks used while copying a fork's records. The fork
@ -198,7 +204,6 @@ impl ExecutionHooks for ForkFinalizationRequirement {
/// Seed the fork: Petri's records, the kept checkpoints and the run branch. The
/// new run must not exist in the store yet.
pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
check(request.store.as_ref(), request.source, request.position).await?;
let source_key = RunKey::new(request.source.to_string());
let fork_key = RunKey::new(request.fork.to_string());
let source_logs = request
@ -207,10 +212,7 @@ pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
.await
.map_err(ForkError::Open)?;
let required = host::stored_state(&*source_logs)
.await
.map_err(ForkError::Seed)?
.required_finalization;
let required = check_logs(&*source_logs, &request.position).await?;
let mut options = RunOptions::new(&request.fork_run_dir);
options.run_key = Some(fork_key.clone());
// A fork only copies records and acquires no sandbox, so it needs no

View file

@ -120,6 +120,7 @@ use crate::checkpoint::{
};
use crate::fork::{self, ForkError};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::projection;
use crate::recovery::{self, Plan, RecoveryError, RestoreTarget};
use crate::source::RunSource;
use crate::workspace::{self, WorkspaceLookup, WorkspaceLookupError};
@ -1620,7 +1621,7 @@ impl ExecutionHooks for FabroHooks {
let message = error.render();
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
if self.publisher.is_some() && finished.status == RunStatus::Success {
return Err(FinalizationFailure::new("publish_failed", message));
return Err(projection::publish_failure(message));
}
None
}
@ -1628,14 +1629,13 @@ impl ExecutionHooks for FabroHooks {
if let Some(publisher) = &self.publisher {
if finished.status == RunStatus::Success {
let publication = publication.ok_or_else(|| {
FinalizationFailure::new(
"publish_failed",
projection::publish_failure(
"the run has no recorded branch and checkpoint to publish",
)
})?;
publisher.publish(&publication).await.map_err(|message| {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
FinalizationFailure::new("publish_failed", message)
projection::publish_failure(message)
})?;
info!(run_id = %self.run_id, branch = publication.run_branch, sha = publication.head_sha, "run published");
}

View file

@ -239,8 +239,7 @@ pub(super) fn finished_status(
reason: FailureReason::Cancelled,
},
_ => RunStatus::Failed {
reason: if finalization_failure.is_some_and(|failure| failure.code == "publish_failed")
{
reason: if finalization_failure.is_some_and(super::is_publish_failure) {
FailureReason::PublishFailed
} else {
FailureReason::WorkflowError

View file

@ -40,10 +40,12 @@ use chrono::{DateTime, TimeZone as _, Utc};
use fabro_store::StagePosition;
use fabro_store::platform_records::StoredPlatformRecord;
use fabro_types::{
RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection,
FailureReason, RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId,
StageProjection,
};
use petri_execution::events::{NodeRef, RunEvent, Subject};
use petri_execution::{CoordinatorEvent, CoordinatorRecord, ExecutionId};
use petri_runtime::ir::FinalizationFailure;
use petri_store::Record;
use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
use serde_json::Value;
@ -369,13 +371,17 @@ pub fn run_id_of(key: &str) -> Option<RunId> {
key.parse().ok()
}
/// The status the view gives the run at Petri's own finish, when the
/// stored record is the coordinator log's `run.finished`: the view reports
/// the run ended from the moment that record is stored, ahead of Fabro's
/// terminal lifecycle record. `None` for any other record.
/// The required-finalization failure Fabro's hooks record when the run's
/// publication fails; its code is [`FailureReason::PublishFailed`]'s.
#[must_use]
pub fn finished_run_status(record: &Record) -> Option<RunStatus> {
finished_run_result(record).map(|(status, _)| status)
pub fn publish_failure(message: impl Into<String>) -> FinalizationFailure {
FinalizationFailure::new(<&'static str>::from(FailureReason::PublishFailed), message)
}
/// Whether a required-finalization failure is the run's failed publication.
#[must_use]
pub fn is_publish_failure(failure: &FinalizationFailure) -> bool {
failure.code == <&'static str>::from(FailureReason::PublishFailed)
}
/// The committed overall status and required-finalization failure message.
@ -423,10 +429,11 @@ mod tests {
#[test]
fn a_finish_record_names_the_status_the_view_ends_the_run_on() {
let finished = |status: &str| {
finished_run_status(&coordinator_record(&serde_json::json!({
finished_run_result(&coordinator_record(&serde_json::json!({
"event": "run.finished",
"status": status,
})))
.map(|(status, _)| status)
};
assert_eq!(
finished("success"),
@ -447,13 +454,13 @@ mod tests {
})
);
assert_eq!(
finished_run_status(&coordinator_record(&serde_json::json!({
finished_run_result(&coordinator_record(&serde_json::json!({
"event": "run.paused",
}))),
None
);
assert_eq!(
finished_run_status(&Record {
finished_run_result(&Record {
seq: 3,
recorded_at: 1_000,
record: serde_json::json!({"event": "run.finished", "status": "success"}),

View file

@ -13,6 +13,7 @@ use petri_store::{Access, LogId, MemoryRunStore, Record, RunKey, RunStore as _};
use tokio::fs;
use tokio::sync::{Notify, Semaphore};
use crate::projection;
use crate::providers::{self, SandboxProviderConfig};
/// A deterministic finalizer gate for testing the committed boundary.
@ -50,7 +51,7 @@ impl ExecutionHooks for TestFinalizer {
.expect("the gate stays open")
.forget();
self.rejection.as_ref().map_or(Ok(()), |message| {
Err(FinalizationFailure::new("publish_failed", message.clone()))
Err(projection::publish_failure(message.clone()))
})
}
}