Retry a finished run from the start

Retry forked the source run at its last checkpoint and reran the failed
stage. It now creates a new run from the source's saved spec and starts
the workflow from the beginning in a fresh workspace, as retry did
before the Petri cutover. The new run records `retried_from` and no
`fork_source_ref`.

Retry no longer needs a checkpoint, a retained workspace, or a published
run branch, so it works for any terminal run that is not archived. To
continue from where a run stopped, fork it at a checkpoint.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-01 11:14:45 -04:00
parent 5879ebf099
commit 3c62d851fe
11 changed files with 108 additions and 103 deletions

View file

@ -2386,12 +2386,11 @@ paths:
tags: [Runs]
summary: Retry Run
description: >
Creates a new run from the terminal source run's last checkpoint and
starts it. When the source failed on a stage, that stage runs again on
the files of the stage before it; otherwise the new run continues from
the last checkpoint as it stands. The new run records `retried_from`
and `fork_source_ref`; the source run is left unchanged. Active and
archived runs are not retryable.
Creates a new run from the terminal source run's saved specification
and starts the workflow from the beginning in a fresh workspace.
No Git checkpoint is required. The new run records `retried_from`
without a `fork_source_ref`; the source run is left unchanged.
Active and archived runs are not retryable.
parameters:
- $ref: "#/components/parameters/RunId"
responses:
@ -2402,7 +2401,7 @@ paths:
schema:
$ref: "#/components/schemas/Run"
"400":
description: The source has no checkpoint to retry from
description: Invalid retry request
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"

View file

@ -90,7 +90,7 @@ fabro [OPTIONS] [COMMAND]
| `fabro provider` | Provider operations |
| `fabro repo` | Repository commands |
| `fabro resume` | Resume an interrupted workflow run |
| `fabro retry` | Retry a finished workflow run from its last checkpoint in a new run |
| `fabro retry` | Retry a finished workflow from the start in a new run |
| `fabro rewind` | Rewind a workflow run to an earlier checkpoint, replacing it |
| `fabro rm` | Remove one or more workflow runs |
| `fabro run` | Register a workflow version, create a run, and start it |
@ -1026,7 +1026,7 @@ fabro resume [OPTIONS] <RUN>
### `fabro retry`
Retry a finished workflow run from its last checkpoint in a new run
Retry a finished workflow from the start in a new run
```bash
fabro retry [OPTIONS] <RUN_ID>

View file

@ -1321,7 +1321,7 @@ pub(crate) enum RunCommands {
Logs(LogsArgs),
/// Resume an interrupted workflow run
Resume(ResumeArgs),
/// Retry a finished workflow run from its last checkpoint in a new run
/// Retry a finished workflow from the start in a new run
Retry(RetryArgs),
/// Fork a workflow run from an earlier checkpoint into a new run
Fork(ForkArgs),

View file

@ -5,8 +5,7 @@ use crate::args::RetryArgs;
use crate::command_context::CommandContext;
use crate::shared::print_json_pretty;
/// Retry a finished run from its last checkpoint in a new run, which starts
/// at once.
/// Retry a finished workflow from the start in a new run.
pub(crate) async fn run(args: &RetryArgs, base_ctx: &CommandContext) -> Result<()> {
let printer = base_ctx.printer();
let ctx = base_ctx.with_target(&args.server)?;

View file

@ -19,7 +19,7 @@ fn help() {
events View the event log of a workflow run
logs View the raw worker tracing log of a workflow run
resume Resume an interrupted workflow run
retry Retry a finished workflow run from its last checkpoint in a new run
retry Retry a finished workflow from the start in a new run
fork Fork a workflow run from an earlier checkpoint into a new run
rewind Rewind a workflow run to an earlier checkpoint, replacing it
timeline Show the checkpoint timeline of a workflow run

View file

@ -251,11 +251,10 @@ async fn a_fork_at_the_first_stage_continues_with_the_rest_on_its_files() {
server.shutdown();
}
/// A retry of a run whose last stage failed reruns that stage: the failure
/// was transient, so the retry passes it on the files of the stage before
/// and finishes the run.
/// A retry starts the workflow over in a new run: every stage runs again
/// in a fresh workspace, so a transient failure passes the second time.
#[tokio::test(flavor = "multi_thread")]
async fn a_retry_reruns_the_failed_stage_and_succeeds_when_the_failure_was_transient() {
async fn a_retry_starts_over_and_succeeds_when_the_failure_was_transient() {
let context = test_context!();
let server = RunningServer::start().await;
let marker = context.temp_dir.join("flaky.marker");
@ -280,22 +279,21 @@ async fn a_retry_reruns_the_failed_stage_and_succeeds_when_the_failure_was_trans
assert_eq!(nodes(&retry_timeline), [
"start", "one", "flaky", "three", "exit"
]);
assert_eq!(retry_timeline["forked_from"]["source_run_id"], source);
assert_eq!(retry_timeline["forked_from"]["rerun_last"], true);
assert!(retry_timeline["forked_from"].is_null());
let retry_workspace = workspace(&server, &retry);
assert_eq!(read(&retry_workspace, "one.txt"), "one\n");
assert_eq!(read(&retry_workspace, "flaky.txt"), "flaky\n");
assert_eq!(read(&retry_workspace, "three.txt"), "three\n");
assert_eq!(commit_subjects(&retry_workspace), [
format!("fabro({source}): start (success)"),
format!("fabro({source}): one (success)"),
format!("fabro({retry}): start (success)"),
format!("fabro({retry}): one (success)"),
format!("fabro({retry}): flaky (success)"),
format!("fabro({retry}): three (success)"),
format!("fabro({retry}): exit (success)"),
]);
let state = run_json(&server, &format!("runs/{retry}/state")).await;
assert_eq!(state["retried_from"], source);
assert_eq!(state["spec"]["fork_source_ref"]["source_run_id"], source);
assert!(state["spec"]["fork_source_ref"].is_null());
// The source is untouched, and a second retry is refused for the
// running or archived cases alone: it is terminal, so it may retry again.

View file

@ -9,9 +9,9 @@
//! (`fabro_petri::fork`), and queues it in resume mode, so its worker
//! acquires a fresh workspace, restores the checkpoint's commit into it and
//! continues from the position. A rewind is a fork of a terminal run that
//! archives the source and records `run.superseded` on it; a retry is a
//! fork of a terminal run at its last checkpoint, with the failed stage run
//! again when the run failed on one.
//! archives the source and records `run.superseded` on it. A retry is not a
//! fork: it is a new run from the source's saved spec, started from the
//! beginning, so it needs no checkpoint or workspace from the source.
use std::sync::Arc;
@ -53,7 +53,6 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
enum ForkKind {
Fork,
Rewind,
Retry,
}
/// A fork made and queued.
@ -86,14 +85,13 @@ async fn run_timeline(
async fn fork_run(
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
State(state): State<Arc<AppState>>,
headers: HeaderMap,
body: Option<Json<api::ForkRequest>>,
) -> Response {
let target = match parse_fork_target(body.and_then(|Json(body)| body.target)) {
Ok(target) => target,
Err(err) => return err.into_response(),
};
match fork_at(state.as_ref(), id, actor, &headers, ForkKind::Fork, target).await {
match fork_at(state.as_ref(), id, actor, ForkKind::Fork, target).await {
Ok(outcome) => (
StatusCode::OK,
Json(api::ForkResponse {
@ -114,23 +112,13 @@ async fn fork_run(
async fn rewind_run(
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
State(state): State<Arc<AppState>>,
headers: HeaderMap,
body: Option<Json<api::RewindRequest>>,
) -> Response {
let target = match parse_fork_target(body.and_then(|Json(body)| body.target)) {
Ok(target) => target,
Err(err) => return err.into_response(),
};
let outcome = match fork_at(
state.as_ref(),
id,
actor.clone(),
&headers,
ForkKind::Rewind,
target,
)
.await
{
let outcome = match fork_at(state.as_ref(), id, actor.clone(), ForkKind::Rewind, target).await {
Ok(outcome) => outcome,
Err(err) => return err.into_response(),
};
@ -175,8 +163,29 @@ async fn retry_run(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
) -> Response {
match fork_at(state.as_ref(), id, actor, &headers, ForkKind::Retry, None).await {
Ok(outcome) => run_response(state.as_ref(), outcome.new_run_id, StatusCode::CREATED).await,
let result = async {
let source = run_records::require_projection(state.as_ref(), id).await?;
operations::ensure_retryable(&source, &id).map_err(workflow_operation_error)?;
let new_run_id = RunId::new();
let storage = Storage::new(state.server_storage_dir());
let run_dir = storage.run_scratch(&new_run_id).root().to_path_buf();
Box::pin(operations::persist_retried_run(
state.store_ref().as_ref(),
&source,
new_run_id,
&run_dir,
run_provenance(&headers, &actor),
state.run_web_url(&new_run_id),
))
.await
.map_err(workflow_operation_error)?;
let projection = run_records::require_projection(state.as_ref(), new_run_id).await?;
queue_run(state.as_ref(), new_run_id, &projection, false, actor).await?;
Ok::<_, ApiError>(new_run_id)
}
.await;
match result {
Ok(id) => run_response(state.as_ref(), id, StatusCode::CREATED).await,
Err(err) => err.into_response(),
}
}
@ -252,7 +261,6 @@ async fn fork_at(
state: &AppState,
id: RunId,
actor: Principal,
headers: &HeaderMap,
kind: ForkKind,
target: Option<ForkTarget>,
) -> Result<ForkOutcome, ApiError> {
@ -260,17 +268,13 @@ async fn fork_at(
match kind {
ForkKind::Fork => operations::ensure_forkable(&source, &id),
ForkKind::Rewind => operations::ensure_rewindable(&source, &id),
ForkKind::Retry => operations::ensure_retryable(&source, &id),
}
.map_err(workflow_operation_error)?;
let timeline = timeline(state, id).await?;
let entry = match kind {
ForkKind::Retry => timeline.latest(),
ForkKind::Fork | ForkKind::Rewind => timeline.resolve_or_latest(target.as_ref()),
}
.map_err(workflow_operation_error)?;
let entry = timeline
.resolve_or_latest(target.as_ref())
.map_err(workflow_operation_error)?;
let resolved = ResolvedForkTarget::of(entry).map_err(workflow_operation_error)?;
let rerun_last = kind == ForkKind::Retry && operations::reruns_last(source.status);
// A position Petri would refuse is refused before the new run exists.
let position = petri_fork::position(resolved.position.execution, resolved.position.firing);
@ -285,15 +289,14 @@ async fn fork_at(
let storage = Storage::new(state.server_storage_dir());
let source_run_dir = storage.run_scratch(&id).root().to_path_buf();
let run_dir = storage.run_scratch(&new_run_id).root().to_path_buf();
let provenance = (kind == ForkKind::Retry).then(|| run_provenance(headers, &actor));
operations::persist_forked_run(state.store_ref().as_ref(), &operations::ForkedRunInput {
source: &source,
new_run_id,
run_dir: run_dir.clone(),
checkpoint_sha: resolved.checkpoint_sha.clone(),
provenance,
provenance: None,
web_url: state.run_web_url(&new_run_id),
retried_from: (kind == ForkKind::Retry).then_some(id),
retried_from: None,
})
.await
.map_err(workflow_operation_error)?;
@ -308,7 +311,7 @@ async fn fork_at(
&state.stores.run_summaries,
))),
position,
rerun_last,
rerun_last: false,
settings: source.spec.settings.run.clone(),
})
.await;
@ -337,7 +340,7 @@ async fn fork_at(
source_run_id: id,
new_run_id,
target: resolved,
rerun_last,
rerun_last: false,
})
}

View file

@ -117,13 +117,28 @@ pub fn forked_run_record(input: &ForkedRunInput<'_>) -> RunCreatedRecord {
/// (`run.created` and the `submitted` transition), which wake its
/// projector.
pub async fn persist_forked_run(store: &Database, input: &ForkedRunInput<'_>) -> Result<(), Error> {
fs::create_dir_all(&input.run_dir).await.map_err(|err| {
persist_new_run(
store,
input.new_run_id,
&input.run_dir,
forked_run_record(input),
)
.await
}
pub(super) async fn persist_new_run(
store: &Database,
run_id: RunId,
run_dir: &std::path::Path,
created: RunCreatedRecord,
) -> Result<(), Error> {
fs::create_dir_all(run_dir).await.map_err(|err| {
Error::Io(format!(
"creating run directory {}: {err}",
input.run_dir.display()
run_dir.display()
))
})?;
let created = PlatformRecord::RunCreated(forked_run_record(input));
let created = PlatformRecord::RunCreated(created);
let submitted = PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(RunLifecycleKind::Submitted).with_status(RunStatus::Submitted),
);
@ -131,11 +146,11 @@ pub async fn persist_forked_run(store: &Database, input: &ForkedRunInput<'_>) ->
let platform_records = summaries.platform_records();
for record in [created, submitted] {
platform_records
.append(&input.new_run_id, &record, None)
.append(&run_id, &record, None)
.await
.map_err(|err| Error::engine_with_source("run store operation failed", err))?;
}
summaries.notify_platform_record(input.new_run_id);
summaries.notify_platform_record(run_id);
Ok(())
}

View file

@ -14,7 +14,7 @@ pub use fork::{
ForkedRunInput, ResolvedForkTarget, ensure_forkable, ensure_terminal, forked_run_record,
persist_forked_run,
};
pub use retry::{ensure_retryable, reruns_last};
pub use retry::{ensure_retryable, persist_retried_run};
pub use rewind::{ensure_rewindable, superseded_record};
pub use timeline::{
ForkTarget, RunTimeline, StageLabel, StageLabels, TimelineEntry, TimelinePosition,

View file

@ -1,47 +1,38 @@
//! Retrying a run: a fork from its last checkpoint.
//!
//! A retry forks a terminal run at its last checkpoint. When the run failed
//! on a stage (its last durable finish is the failed stage's), the position's
//! firing runs again, so the retry reruns the failed stage on the files of
//! the stage before it; a run that succeeded, was cancelled or died forks at
//! the last position as it stands, and continues from there.
//! Retrying starts a new execution from the original spec. It does not
//! require Git checkpoints or reconstruct the previous workspace.
use fabro_types::{FailureReason, RunId, RunProjection, RunStatus};
use std::path::Path;
use super::fork::{ensure_forkable, ensure_terminal};
use fabro_store::platform_records::RunCreatedRecord;
use fabro_store::{Database, RunProjection};
use fabro_types::{RunId, RunProvenance};
use super::fork;
use crate::error::Error;
/// A run can be retried when it is terminal and not archived.
pub fn ensure_retryable(source: &RunProjection, run_id: &RunId) -> Result<(), Error> {
ensure_forkable(source, run_id)?;
ensure_terminal(source, run_id, "retry")
fork::ensure_forkable(source, run_id)?;
fork::ensure_terminal(source, run_id, "retry")
}
/// Whether the retry reruns the last checkpointed stage: it does when the
/// run failed on its own terms, since that stage's finish is the failure.
#[must_use]
pub fn reruns_last(status: RunStatus) -> bool {
matches!(
status,
RunStatus::Failed { reason } if reason != FailureReason::Cancelled
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_failed_run_reruns_its_last_stage_and_the_others_continue() {
assert!(reruns_last(RunStatus::Failed {
reason: FailureReason::WorkflowError,
}));
assert!(!reruns_last(RunStatus::Failed {
reason: FailureReason::Cancelled,
}));
assert!(!reruns_last(RunStatus::Dead));
assert!(!reruns_last(RunStatus::Succeeded {
reason: fabro_types::SuccessReason::Completed,
}));
}
pub async fn persist_retried_run(
store: &Database,
source: &RunProjection,
run_id: RunId,
run_dir: &Path,
provenance: RunProvenance,
web_url: Option<String>,
) -> Result<(), Error> {
let mut spec = source.spec.clone();
spec.run_id = run_id;
spec.fork_source_ref = None;
spec.provenance = provenance;
let created = RunCreatedRecord {
spec,
title: Some(source.title().into_owned()),
parent_id: source.parent_id,
retried_from: Some(source.spec.run_id),
web_url,
};
fork::persist_new_run(store, run_id, run_dir, created).await
}

View file

@ -1164,7 +1164,7 @@ export const RunsApiAxiosParamCreator = function (configuration?: Configuration)
};
},
/**
* Creates a new run from the terminal source run\'s last checkpoint and starts it. When the source failed on a stage, that stage runs again on the files of the stage before it; otherwise the new run continues from the last checkpoint as it stands. The new run records `retried_from` and `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* Creates a new run from the terminal source run\'s saved specification and starts the workflow from the beginning in a fresh workspace. No Git checkpoint is required. The new run records `retried_from` without a `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* @summary Retry Run
* @param {string} id Unique run identifier (ULID).
* @param {*} [options] Override http request option.
@ -1925,7 +1925,7 @@ export const RunsApiFp = function(configuration?: Configuration) {
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
},
/**
* Creates a new run from the terminal source run\'s last checkpoint and starts it. When the source failed on a stage, that stage runs again on the files of the stage before it; otherwise the new run continues from the last checkpoint as it stands. The new run records `retried_from` and `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* Creates a new run from the terminal source run\'s saved specification and starts the workflow from the beginning in a fresh workspace. No Git checkpoint is required. The new run records `retried_from` without a `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* @summary Retry Run
* @param {string} id Unique run identifier (ULID).
* @param {*} [options] Override http request option.
@ -2331,7 +2331,7 @@ export const RunsApiFactory = function (configuration?: Configuration, basePath?
return localVarFp.retrieveRunGraphSource(id, options).then((request) => request(axios, basePath));
},
/**
* Creates a new run from the terminal source run\'s last checkpoint and starts it. When the source failed on a stage, that stage runs again on the files of the stage before it; otherwise the new run continues from the last checkpoint as it stands. The new run records `retried_from` and `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* Creates a new run from the terminal source run\'s saved specification and starts the workflow from the beginning in a fresh workspace. No Git checkpoint is required. The new run records `retried_from` without a `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* @summary Retry Run
* @param {string} id Unique run identifier (ULID).
* @param {*} [options] Override http request option.
@ -2730,7 +2730,7 @@ export class RunsApi extends BaseAPI {
}
/**
* Creates a new run from the terminal source run\'s last checkpoint and starts it. When the source failed on a stage, that stage runs again on the files of the stage before it; otherwise the new run continues from the last checkpoint as it stands. The new run records `retried_from` and `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* Creates a new run from the terminal source run\'s saved specification and starts the workflow from the beginning in a fresh workspace. No Git checkpoint is required. The new run records `retried_from` without a `fork_source_ref`; the source run is left unchanged. Active and archived runs are not retryable.
* @summary Retry Run
* @param {string} id Unique run identifier (ULID).
* @param {*} [options] Override http request option.