mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
Commit checkpoint failures through required finalization
A failed checkpoint cancels the run, so Petri recorded it as cancelled and the worker overrode its own outcome in memory. Runs that checkpoint now declare required finalization, and finalize_run rejects with checkpoint_failed before publishing. The projection and the engine outcome report a cancelled finish carrying that failure as a workflow failure with the checkpoint's message, and the in-memory override is gone. A run whose checkpoint failed is never published. Retry a host failure that storage rejects with a short backoff, so a brief storage fault does not leave the run active and holding its scheduler slot. A failure that never commits still leaves the run's status alone. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
parent
b4aea6128d
commit
04309977e5
9 changed files with 195 additions and 45 deletions
|
|
@ -1,7 +1,11 @@
|
|||
# Required run finalization
|
||||
|
||||
Petri owns the durable run result. Fabro declares required finalization when
|
||||
its hooks have a publisher and implements it in `FabroHooks::finalize_run`.
|
||||
its hooks checkpoint or have a publisher, and implements it in
|
||||
`FabroHooks::finalize_run`. A failed checkpoint rejects with
|
||||
`checkpoint_failed` and skips publication. The checkpoint failure cancelled the
|
||||
run, so Petri may commit `cancelled`; Fabro reports that finish as a workflow
|
||||
failure with the checkpoint's message.
|
||||
For successful workflow execution this prepares the final diff, then uses the
|
||||
existing publisher to retry the push and reconcile or create the pull request.
|
||||
A preparation or publication error rejects with `publish_failed` and a rendered
|
||||
|
|
|
|||
|
|
@ -1109,7 +1109,7 @@ fn attach_json_errors_without_prompting_for_human_input() {
|
|||
"middleware_chain": [
|
||||
"circuit-breaker"
|
||||
],
|
||||
"required_finalization": false
|
||||
"required_finalization": true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1077,6 +1077,49 @@ mod tests {
|
|||
runtime.end_worker();
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn a_host_failure_is_retried_until_storage_accepts_it() {
|
||||
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 failure = tokio::spawn({
|
||||
let state = Arc::clone(&state);
|
||||
async move {
|
||||
crate::server::persist_run_failure(
|
||||
&state,
|
||||
run_id,
|
||||
FailureReason::LaunchFailed,
|
||||
"worker launch failed".into(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
});
|
||||
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
|
||||
sqlx::query("DROP TRIGGER reject_platform_record")
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
failure.await.unwrap();
|
||||
let failed = RunStatus::Failed {
|
||||
reason: FailureReason::LaunchFailed,
|
||||
};
|
||||
let stored = crate::server::run_records::projection(&state, run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("the run is stored");
|
||||
assert_eq!(stored.status, failed, "the retried failure is durable");
|
||||
assert_eq!(state.test_managed_run_status(&run_id), Some(failed));
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -3512,29 +3512,55 @@ async fn reject_run_if_sandbox_provider_disabled(
|
|||
/// Record a host failure only while no terminal result is committed. This
|
||||
/// also handles a worker wait/launch error racing its durable Petri finish.
|
||||
/// The managed run settles on whichever terminal result the store committed.
|
||||
/// When nothing could be committed its live state is still released, but its
|
||||
/// status is left alone: the API never reports an outcome storage lacks.
|
||||
/// A failed commit is retried, so a brief storage fault does not leave the
|
||||
/// run active. When nothing could be committed its live state is still
|
||||
/// released, but its status is left alone: the API never reports an outcome
|
||||
/// storage lacks, and the restart reconciliation fails the run.
|
||||
pub(crate) async fn persist_run_failure(
|
||||
state: &Arc<AppState>,
|
||||
run_id: RunId,
|
||||
reason: FailureReason,
|
||||
message: String,
|
||||
) {
|
||||
match commit_host_failure(state, run_id, reason, message).await {
|
||||
Ok(committed) => {
|
||||
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);
|
||||
let mut retry_delays = HOST_FAILURE_RETRY_DELAYS.iter();
|
||||
loop {
|
||||
match commit_host_failure(state, run_id, reason, message.clone()).await {
|
||||
Ok(committed) => {
|
||||
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);
|
||||
break;
|
||||
}
|
||||
Err(err) => {
|
||||
let Some(delay) = retry_delays.next() else {
|
||||
error!(run_id = %run_id, error = %err, "Failed to record a host failure");
|
||||
break;
|
||||
};
|
||||
warn!(
|
||||
run_id = %run_id,
|
||||
error = %err,
|
||||
retry_in_ms = delay.as_millis(),
|
||||
"Failed to record a host failure; retrying"
|
||||
);
|
||||
sleep(*delay).await;
|
||||
}
|
||||
}
|
||||
Err(err) => error!(run_id = %run_id, error = %err, "Failed to record a host failure"),
|
||||
}
|
||||
release_managed_run(state, run_id);
|
||||
state.scheduler_notify.notify_one();
|
||||
}
|
||||
|
||||
/// How long [`persist_run_failure`] waits before each retry of a failed
|
||||
/// commit.
|
||||
const HOST_FAILURE_RETRY_DELAYS: [Duration; 3] = [
|
||||
Duration::from_millis(250),
|
||||
Duration::from_secs(1),
|
||||
Duration::from_secs(4),
|
||||
];
|
||||
|
||||
/// Append the host failure unless the run already ended, and return the
|
||||
/// terminal projection the store committed: the failure, or the finish that
|
||||
/// won the race.
|
||||
|
|
|
|||
|
|
@ -354,17 +354,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
|||
return Err(RunError::StoreFailed(message));
|
||||
}
|
||||
let inspection = inspect(request.store.as_ref(), &key).await?;
|
||||
let mut outcome = outcome(inspection, result.err())?;
|
||||
// A failed checkpoint cancelled the run; what Fabro reports is the
|
||||
// checkpoint failure, not a cancellation.
|
||||
if let Some(failure) = fabro_hooks
|
||||
.as_ref()
|
||||
.and_then(|hooks| hooks.checkpoint_failure())
|
||||
{
|
||||
outcome.status = RunStatus::Failed;
|
||||
outcome.failure = Some(failure);
|
||||
}
|
||||
Ok(outcome)
|
||||
outcome(inspection, result.err())
|
||||
}
|
||||
|
||||
/// When Petri keeps a run's workspaces after their scope is released.
|
||||
|
|
@ -535,9 +525,16 @@ fn outcome(
|
|||
inspection: RunInspection,
|
||||
host_error: Option<HostError>,
|
||||
) -> Result<RunOutcome, RunError> {
|
||||
// A failed checkpoint cancelled the run; its finish records the
|
||||
// checkpoint failure, which is what Fabro reports.
|
||||
let checkpoint_failed = inspection
|
||||
.finalization_failure
|
||||
.as_ref()
|
||||
.is_some_and(projection::is_checkpoint_failure);
|
||||
let status = match inspection.status.as_deref() {
|
||||
Some("success") => RunStatus::Success,
|
||||
Some("failed") => RunStatus::Failed,
|
||||
Some("cancelled") if checkpoint_failed => RunStatus::Failed,
|
||||
Some("cancelled") => RunStatus::Cancelled,
|
||||
_ => {
|
||||
let mut reasons = inspection.incomplete.clone();
|
||||
|
|
|
|||
|
|
@ -31,14 +31,15 @@
|
|||
//! `artifact.collected` record, unless the same file with the same content
|
||||
//! was already collected earlier in the run. A failed write is a recorded
|
||||
//! problem on the transition, never a blocked route.
|
||||
//! - `finalize_run`: the run's diff, its run branch against its base commit, as
|
||||
//! the `run.diff` platform record with the patch as a blob; for a successful
|
||||
//! run, its publication ([`RunPublisher`]: the platform pushes the run branch
|
||||
//! and opens a pull request), whose failure fails the run before its terminal
|
||||
//! record.
|
||||
//! - `run_finished`: best-effort diff preparation for nonpublishing runs, then
|
||||
//! the forwarded point, so the local service runs `run_complete` and
|
||||
//! `run_failed` with the sandbox in place.
|
||||
//! - `finalize_run`, required when the run checkpoints or publishes: the run's
|
||||
//! diff, its run branch against its base commit, as the `run.diff` platform
|
||||
//! record with the patch as a blob; a failed checkpoint, which fails the run
|
||||
//! and skips publication; for a successful run, its publication
|
||||
//! ([`RunPublisher`]: the platform pushes the run branch and opens a pull
|
||||
//! request), whose failure fails the run before its terminal record.
|
||||
//! - `run_finished`: best-effort diff preparation for runs without required
|
||||
//! finalization, then the forwarded point, so the local service runs
|
||||
//! `run_complete` and `run_failed` with the sandbox in place.
|
||||
//! - `scope_acquired`: a fresh run's Git target checked out into the workspace
|
||||
//! from inside the scope ([`crate::source`]); a resumed run uses its
|
||||
//! surviving workspace, while an explicit fork fetches the source run's
|
||||
|
|
@ -541,8 +542,8 @@ impl FabroHooks {
|
|||
}
|
||||
}
|
||||
|
||||
/// The checkpoint failure that ended the run, when one did: what the
|
||||
/// engine reports the run failed with.
|
||||
/// The checkpoint failure that ended the run, when one did: required
|
||||
/// finalization commits it as the run's failure.
|
||||
#[must_use]
|
||||
pub fn checkpoint_failure(&self) -> Option<String> {
|
||||
sync::lock(&self.failure).clone()
|
||||
|
|
@ -1607,7 +1608,9 @@ impl ExecutionHooks for FabroHooks {
|
|||
}
|
||||
|
||||
fn requires_run_finalization(&self) -> bool {
|
||||
self.publisher.is_some() || self.inner.requires_run_finalization()
|
||||
self.publisher.is_some()
|
||||
|| self.checkpoint_enabled
|
||||
|| self.inner.requires_run_finalization()
|
||||
}
|
||||
|
||||
async fn finalize_run(
|
||||
|
|
@ -1615,11 +1618,19 @@ impl ExecutionHooks for FabroHooks {
|
|||
context: &HookContext,
|
||||
finished: RunFinished,
|
||||
) -> Result<(), FinalizationFailure> {
|
||||
let publication = match self.record_run_diff().await {
|
||||
let diff = self.record_run_diff().await.map_err(|error| {
|
||||
let message = error.render();
|
||||
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
|
||||
message
|
||||
});
|
||||
// A failed checkpoint fails the run whatever its execution status,
|
||||
// and its work is never published.
|
||||
if let Some(message) = self.checkpoint_failure() {
|
||||
return Err(projection::checkpoint_failure(message));
|
||||
}
|
||||
let publication = match diff {
|
||||
Ok(publication) => publication,
|
||||
Err(error) => {
|
||||
let message = error.render();
|
||||
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
|
||||
Err(message) => {
|
||||
if self.publisher.is_some() && finished.status == RunStatus::Success {
|
||||
return Err(projection::publish_failure(message));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -226,7 +226,9 @@ impl RunView {
|
|||
}
|
||||
|
||||
/// The status Fabro gives a run at Petri's finish, by the status the finish
|
||||
/// records (`success`, `cancelled`, or a failure).
|
||||
/// records (`success`, `cancelled`, or a failure). A failed checkpoint
|
||||
/// cancels the run, but the run failed: its finish says so with the
|
||||
/// checkpoint's finalization failure.
|
||||
pub(super) fn finished_status(
|
||||
status: &str,
|
||||
finalization_failure: Option<&FinalizationFailure>,
|
||||
|
|
@ -235,9 +237,11 @@ pub(super) fn finished_status(
|
|||
"success" => RunStatus::Succeeded {
|
||||
reason: SuccessReason::Completed,
|
||||
},
|
||||
"cancelled" => RunStatus::Failed {
|
||||
reason: FailureReason::Cancelled,
|
||||
},
|
||||
"cancelled" if !finalization_failure.is_some_and(super::is_checkpoint_failure) => {
|
||||
RunStatus::Failed {
|
||||
reason: FailureReason::Cancelled,
|
||||
}
|
||||
}
|
||||
_ => RunStatus::Failed {
|
||||
reason: if finalization_failure.is_some_and(super::is_publish_failure) {
|
||||
FailureReason::PublishFailed
|
||||
|
|
|
|||
|
|
@ -51,6 +51,8 @@ use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
|
|||
use serde_json::Value;
|
||||
use tracing::debug;
|
||||
|
||||
use crate::checkpoint::CHECKPOINT_FAILED_CLASS;
|
||||
|
||||
/// One item the projector hands the fold, with its delivery sequence.
|
||||
pub enum Item<'a> {
|
||||
Petri(&'a RunEvent),
|
||||
|
|
@ -384,6 +386,20 @@ pub fn is_publish_failure(failure: &FinalizationFailure) -> bool {
|
|||
failure.code == <&'static str>::from(FailureReason::PublishFailed)
|
||||
}
|
||||
|
||||
/// The required-finalization failure Fabro's hooks record when a checkpoint
|
||||
/// failed during the run. The failed checkpoint cancelled the run, so Petri
|
||||
/// may record it as cancelled; Fabro reports it as a workflow failure.
|
||||
#[must_use]
|
||||
pub fn checkpoint_failure(message: impl Into<String>) -> FinalizationFailure {
|
||||
FinalizationFailure::new(CHECKPOINT_FAILED_CLASS, message)
|
||||
}
|
||||
|
||||
/// Whether a required-finalization failure is a failed checkpoint.
|
||||
#[must_use]
|
||||
pub fn is_checkpoint_failure(failure: &FinalizationFailure) -> bool {
|
||||
failure.code == CHECKPOINT_FAILED_CLASS
|
||||
}
|
||||
|
||||
/// The committed overall status and required-finalization failure message.
|
||||
/// Execution failure details remain in the invocation records and projection.
|
||||
#[must_use]
|
||||
|
|
@ -453,6 +469,23 @@ mod tests {
|
|||
reason: fabro_types::FailureReason::WorkflowError,
|
||||
})
|
||||
);
|
||||
assert_eq!(
|
||||
finished_run_result(&coordinator_record(&serde_json::json!({
|
||||
"event": "run.finished",
|
||||
"status": "cancelled",
|
||||
"finalization_failure": {
|
||||
"code": "checkpoint_failed",
|
||||
"message": "checkpoint commit of `wreck` failed",
|
||||
},
|
||||
}))),
|
||||
Some((
|
||||
RunStatus::Failed {
|
||||
reason: fabro_types::FailureReason::WorkflowError,
|
||||
},
|
||||
Some("checkpoint commit of `wreck` failed".to_string()),
|
||||
)),
|
||||
"a failed checkpoint's cancellation is the checkpoint's failure"
|
||||
);
|
||||
assert_eq!(
|
||||
finished_run_result(&coordinator_record(&serde_json::json!({
|
||||
"event": "run.paused",
|
||||
|
|
|
|||
|
|
@ -689,8 +689,9 @@ async fn a_failed_stage_is_committed_and_its_route_sees_the_files() {
|
|||
}
|
||||
|
||||
/// A checkpoint commit that fails is fatal: the stage's outcome is recorded
|
||||
/// as `checkpoint_failed`, no route is taken, the run ends failed with the
|
||||
/// checkpoint's error, and a restart reports it failed without resuming.
|
||||
/// as `checkpoint_failed`, no route is taken, the run's committed finish
|
||||
/// fails it with the checkpoint's error, and a restart reports it failed
|
||||
/// without resuming.
|
||||
#[tokio::test]
|
||||
async fn a_failed_checkpoint_ends_the_run_with_no_route() {
|
||||
let harness = Harness::new().await;
|
||||
|
|
@ -713,6 +714,16 @@ async fn a_failed_checkpoint_ends_the_run_with_no_route() {
|
|||
);
|
||||
|
||||
let inspection = harness.inspection().await;
|
||||
let committed = inspection
|
||||
.finalization_failure
|
||||
.clone()
|
||||
.expect("the finish commits the checkpoint failure");
|
||||
assert_eq!(committed.code, CHECKPOINT_FAILED_CLASS);
|
||||
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
|
||||
.await
|
||||
.expect("the finished run reads back");
|
||||
assert_eq!(stored.status, RunStatus::Failed, "{stored:?}");
|
||||
assert_eq!(stored.failure, outcome.failure);
|
||||
let attempts: Vec<_> = inspection
|
||||
.executions
|
||||
.iter()
|
||||
|
|
@ -1359,6 +1370,27 @@ async fn a_failed_run_is_not_published() {
|
|||
assert!(publisher.published.lock().unwrap().is_empty());
|
||||
}
|
||||
|
||||
/// A run whose checkpoint failed is never published, and its committed
|
||||
/// failure is the checkpoint's, not a publication's.
|
||||
#[tokio::test]
|
||||
async fn a_run_whose_checkpoint_failed_is_not_published() {
|
||||
let publisher = RecordingPublisher::new(None);
|
||||
let (harness, outcome) = published_run(
|
||||
"script=\"rm -rf .git && echo garbage > .git && echo wrecked > out.txt\"",
|
||||
&publisher,
|
||||
)
|
||||
.await;
|
||||
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
|
||||
assert!(!outcome.publish_failed);
|
||||
assert!(publisher.published.lock().unwrap().is_empty());
|
||||
let failure = harness
|
||||
.inspection()
|
||||
.await
|
||||
.finalization_failure
|
||||
.expect("the finish commits the checkpoint failure");
|
||||
assert_eq!(failure.code, CHECKPOINT_FAILED_CLASS);
|
||||
}
|
||||
|
||||
struct OriginPublisher {
|
||||
origin: String,
|
||||
pushed: std::sync::Mutex<Vec<String>>,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue