diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index 9eaa57b2d..be4ef37cc 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -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, diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index b496fc27c..16f1df26e 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -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, - 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, 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, diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 197d54b2d..176fd5232 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -549,13 +549,7 @@ pub(crate) async fn execute(state: Arc, 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, 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, + run_id: RunId, + record: RunLifecycleRecord, + status: RunStatus, + error: Option, +) { 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, 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 diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index b4c49c8bb..2e9baacc0 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -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) diff --git a/lib/components/fabro-petri/src/fork.rs b/lib/components/fabro-petri/src/fork.rs index 9f6a8827d..3b96c023f 100644 --- a/lib/components/fabro-petri/src/fork.rs +++ b/lib/components/fabro-petri/src/fork.rs @@ -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 { + 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 { - 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 { .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 diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 9cedffee4..687a18d61 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -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"); } diff --git a/lib/components/fabro-petri/src/projection/coordinator.rs b/lib/components/fabro-petri/src/projection/coordinator.rs index 888f77c35..4af5fe1b7 100644 --- a/lib/components/fabro-petri/src/projection/coordinator.rs +++ b/lib/components/fabro-petri/src/projection/coordinator.rs @@ -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 diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index d93a449b7..39f06d948 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -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 { 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 { - finished_run_result(record).map(|(status, _)| status) +pub fn publish_failure(message: impl Into) -> 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"}), diff --git a/lib/components/fabro-petri/src/test_support/finalization.rs b/lib/components/fabro-petri/src/test_support/finalization.rs index b396c9342..b4ee72814 100644 --- a/lib/components/fabro-petri/src/test_support/finalization.rs +++ b/lib/components/fabro-petri/src/test_support/finalization.rs @@ -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())) }) } }