Restore the timeline, fork, rewind and retry operations over checkpoints

The operations the cutover removed come back in their Petri shape, with
no engine of their own: the timeline is the run's checkpoint records
labelled by stage (`RunTimeline::build`), and a target (`@ordinal`, a
node, or `node@visit`) resolves to one of them. A fork's run row is the
source's spec under a new id with `fork_source_ref` naming the source and
the checkpoint's commit (`persist_forked_run`), and `retried_from` when
it is a retry. A rewind needs a terminal source and records
`run.superseded` on it; a retry needs a terminal source and reruns the
last checkpointed stage when the run failed on it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 23:17:25 -04:00
parent 0cd0df43ae
commit f18f02206b
No known key found for this signature in database
5 changed files with 672 additions and 0 deletions

View file

@ -0,0 +1,255 @@
//! Forking a run at a checkpoint: the Fabro run the fork becomes.
//!
//! A fork is a new run whose records Petri seeds from the source's up to a
//! checkpoint's position (`fabro_petri::fork`, over the timeline here). What
//! Fabro itself makes of it is a run row like any other: the `run.created`
//! record carrying the source's spec (its admission, settings and target)
//! under the new id, with `fork_source_ref` naming the source and the
//! checkpoint's commit, and `retried_from` when the fork is a retry; then
//! the `submitted` lifecycle transition. The run is then started in resume
//! mode, as a run left in flight is.
use std::path::PathBuf;
use fabro_store::Database;
use fabro_store::platform_records::{
PlatformRecord, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
};
use fabro_types::{ForkSourceRef, RunId, RunProjection, RunProvenance, RunStatus};
use tokio::fs;
use super::ensure_not_archived;
use super::timeline::{TimelineEntry, TimelinePosition};
use crate::error::Error;
/// The checkpoint a fork was resolved to, as the API reports it.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ResolvedForkTarget {
pub checkpoint_ordinal: usize,
pub node_id: String,
pub visit: usize,
pub position: TimelinePosition,
pub checkpoint_sha: String,
}
impl ResolvedForkTarget {
/// The entry as a fork target, refused when it has no commit.
pub fn of(entry: &TimelineEntry) -> Result<Self, Error> {
let checkpoint_sha = entry.run_commit_sha.clone().ok_or_else(|| {
Error::Validation(format!(
"checkpoint @{} has no git_commit_sha; cannot fork",
entry.ordinal
))
})?;
Ok(Self {
checkpoint_ordinal: entry.ordinal,
node_id: entry.node_name.clone(),
visit: usize::try_from(entry.visit).unwrap_or(1),
position: entry.position,
checkpoint_sha,
})
}
#[must_use]
pub fn response_target(&self) -> String {
format!("@{}", self.checkpoint_ordinal)
}
}
/// The new run a fork creates.
#[derive(Debug)]
pub struct ForkedRunInput<'a> {
pub source: &'a RunProjection,
pub new_run_id: RunId,
/// The new run's scratch directory.
pub run_dir: PathBuf,
pub checkpoint_sha: String,
/// Who forked, when the fork records its own provenance; `None` keeps
/// the source's.
pub provenance: Option<RunProvenance>,
pub web_url: Option<String>,
/// The source, when the fork is a retry of it.
pub retried_from: Option<RunId>,
}
/// A run can be forked unless it is archived.
pub fn ensure_forkable(source: &RunProjection, run_id: &RunId) -> Result<(), Error> {
ensure_not_archived(source.archived_at.is_some(), run_id)
}
/// A run must be terminal to be rewound or retried: its records are
/// complete, and nothing is writing them.
pub fn ensure_terminal(source: &RunProjection, run_id: &RunId, verb: &str) -> Result<(), Error> {
let current = source.status;
if current.is_terminal() {
Ok(())
} else {
Err(Error::Precondition(format!(
"run {run_id} must be terminal (succeeded, failed, or dead) to {verb}; current status \
is {current}"
)))
}
}
/// The `run.created` record of the fork: the source's spec under the new
/// id, naming where it came from.
#[must_use]
pub fn forked_run_record(input: &ForkedRunInput<'_>) -> RunCreatedRecord {
let mut spec = input.source.spec.clone();
spec.run_id = input.new_run_id;
spec.fork_source_ref = Some(ForkSourceRef {
source_run_id: input.source.spec.run_id,
checkpoint_sha: input.checkpoint_sha.clone(),
});
if let Some(provenance) = &input.provenance {
spec.provenance = provenance.clone();
}
RunCreatedRecord {
spec,
title: Some(input.source.title().into_owned()),
parent_id: input.source.parent_id,
retried_from: input.retried_from,
web_url: input.web_url.clone(),
}
}
/// Create the fork's run: its scratch directory, then its first records
/// (`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| {
Error::Io(format!(
"creating run directory {}: {err}",
input.run_dir.display()
))
})?;
let created = PlatformRecord::RunCreated(forked_run_record(input));
let submitted = PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(RunLifecycleKind::Submitted).with_status(RunStatus::Submitted),
);
let summaries = store.run_summary_store();
let platform_records = summaries.platform_records();
for record in [created, submitted] {
platform_records
.append(&input.new_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);
Ok(())
}
#[cfg(test)]
mod tests {
use chrono::Utc;
use fabro_types::{FailureReason, Graph, PetriAdmission, RunSpec, WorkflowSettings, fixtures};
use super::*;
fn source(status: RunStatus) -> RunProjection {
let mut projection = RunProjection::new(
"Source title".to_string(),
RunSpec {
run_id: fixtures::RUN_1,
settings: WorkflowSettings::default(),
graph: Graph::new("source"),
graph_source: Some("digraph source { start -> exit }".to_string()),
workflow_slug: Some("source".to_string()),
workflow_version_id: None,
target: None,
automation: None,
source_directory: None,
labels: std::collections::HashMap::new(),
provenance: fabro_types::test_support::test_run_provenance(),
definition_blob: None,
spec_blob: None,
git: None,
fork_source_ref: None,
admission: PetriAdmission::default(),
},
Utc::now(),
);
projection.status = status;
projection.parent_id = Some(fixtures::RUN_2);
projection
}
fn entry(ordinal: usize, sha: Option<&str>) -> TimelineEntry {
TimelineEntry {
ordinal,
checkpoint_seq: 3,
position: TimelinePosition {
execution: 0,
firing: 2,
attempt: 1,
},
stage_id: Some("build@1".to_string()),
node_name: "build".to_string(),
visit: 1,
workspace: None,
run_commit_sha: sha.map(ToOwned::to_owned),
diff_summary: None,
}
}
#[test]
fn the_record_carries_the_source_spec_under_the_new_id_and_names_the_source() {
let source = source(RunStatus::Succeeded {
reason: fabro_types::SuccessReason::Completed,
});
let record = forked_run_record(&ForkedRunInput {
source: &source,
new_run_id: fixtures::RUN_3,
run_dir: PathBuf::from("/tmp/unused"),
checkpoint_sha: "abc".to_string(),
provenance: None,
web_url: Some("http://localhost/runs/x".to_string()),
retried_from: Some(fixtures::RUN_1),
});
assert_eq!(record.spec.run_id, fixtures::RUN_3);
assert_eq!(
record.spec.fork_source_ref,
Some(ForkSourceRef {
source_run_id: fixtures::RUN_1,
checkpoint_sha: "abc".to_string(),
})
);
assert_eq!(record.spec.graph.name, "source");
assert_eq!(record.title.as_deref(), Some("Source title"));
assert_eq!(record.parent_id, Some(fixtures::RUN_2));
assert_eq!(record.retried_from, Some(fixtures::RUN_1));
assert_eq!(record.web_url.as_deref(), Some("http://localhost/runs/x"));
}
#[test]
fn a_target_needs_a_commit() {
let resolved = ResolvedForkTarget::of(&entry(2, Some("abc"))).unwrap();
assert_eq!(resolved.response_target(), "@2");
assert_eq!(resolved.checkpoint_sha, "abc");
assert!(matches!(
ResolvedForkTarget::of(&entry(2, None)),
Err(Error::Validation(message)) if message.contains("no git_commit_sha")
));
}
#[test]
fn a_rewind_or_retry_needs_a_terminal_source() {
let running = source(RunStatus::Running);
assert!(matches!(
ensure_terminal(&running, &fixtures::RUN_1, "rewind"),
Err(Error::Precondition(message)) if message.contains("must be terminal")
));
for status in [
RunStatus::Dead,
RunStatus::Failed {
reason: FailureReason::Cancelled,
},
RunStatus::Succeeded {
reason: fabro_types::SuccessReason::Completed,
},
] {
ensure_terminal(&source(status), &fixtures::RUN_1, "retry").unwrap();
}
ensure_forkable(&running, &fixtures::RUN_1).unwrap();
}
}

View file

@ -1,5 +1,9 @@
mod create;
mod fork;
mod retry;
mod rewind;
mod source;
mod timeline;
mod validate;
pub use create::{
@ -8,7 +12,16 @@ pub use create::{
make_run_dir, materialize_admitted_run, persist_create_run,
};
use fabro_types::RunId;
pub use fork::{
ForkedRunInput, ResolvedForkTarget, ensure_forkable, ensure_terminal, forked_run_record,
persist_forked_run,
};
pub use retry::{ensure_retryable, reruns_last};
pub use rewind::{ensure_rewindable, superseded_record};
pub use source::WorkflowInput;
pub use timeline::{
ForkTarget, RunTimeline, StageLabel, StageLabels, TimelineEntry, TimelinePosition,
};
pub use validate::{ValidateInput, validate};
pub use crate::error::Error;

View file

@ -0,0 +1,47 @@
//! 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.
use fabro_types::{FailureReason, RunId, RunProjection, RunStatus};
use super::fork::{ensure_forkable, ensure_terminal};
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")
}
/// 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,
}));
}
}

View file

@ -0,0 +1,30 @@
//! Rewinding a run: a fork that replaces its source.
//!
//! A rewind forks a terminal run at a checkpoint (`fork`), then archives the
//! source and records `run.superseded` on it, naming the new run and the
//! checkpoint it continues from. The archive is the server's own archive
//! operation; what this module holds is the precondition and the record.
use fabro_store::platform_records::{PlatformRecord, RunSupersededRecord};
use fabro_types::{RunId, RunProjection};
use super::fork::{ResolvedForkTarget, ensure_forkable, ensure_terminal};
use crate::error::Error;
/// A run can be rewound when it is terminal and not archived.
pub fn ensure_rewindable(source: &RunProjection, run_id: &RunId) -> Result<(), Error> {
ensure_forkable(source, run_id)?;
ensure_terminal(source, run_id, "rewind")
}
/// The record a rewound source carries: which run replaced it, from which
/// checkpoint.
#[must_use]
pub fn superseded_record(new_run_id: RunId, target: &ResolvedForkTarget) -> PlatformRecord {
PlatformRecord::RunSuperseded(RunSupersededRecord {
new_run_id,
target_checkpoint_ordinal: target.checkpoint_ordinal,
target_node_id: target.node_id.clone(),
target_visit: target.visit,
})
}

View file

@ -0,0 +1,327 @@
//! The checkpoint timeline of a run: every checkpoint Fabro recorded, at
//! its Petri position, with the commit it made and the stage it belongs
//! to, and the targets a fork names one of them by.
use std::collections::BTreeMap;
use std::str::FromStr;
use fabro_store::platform_records::{CheckpointRecord, DecisionRef};
use fabro_store::{PlatformRecord, StoredPlatformRecord};
use fabro_types::DiffSummary;
use crate::error::Error;
/// How a caller names a checkpoint: by ordinal (`@2`), by the latest visit
/// of a node (`build`), or by one visit of it (`build@1`).
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ForkTarget {
Ordinal(usize),
LatestVisit(String),
SpecificVisit(String, usize),
}
impl FromStr for ForkTarget {
type Err = Error;
fn from_str(s: &str) -> Result<Self, Error> {
if let Some(rest) = s.strip_prefix('@') {
let n: usize = rest
.parse()
.map_err(|_| Error::Validation(format!("invalid ordinal: @{rest}")))?;
if n == 0 {
return Err(Error::Validation("ordinal must be >= 1".to_string()));
}
return Ok(Self::Ordinal(n));
}
if let Some((name, visit)) = s.rsplit_once('@') {
if !name.is_empty() && !visit.is_empty() {
if let Ok(visit) = visit.parse::<usize>() {
if visit == 0 {
return Err(Error::Validation("visit number must be >= 1".to_string()));
}
return Ok(Self::SpecificVisit(name.to_string(), visit));
}
}
}
if s.trim().is_empty() {
return Err(Error::Validation("a target names a checkpoint".to_string()));
}
Ok(Self::LatestVisit(s.to_string()))
}
}
/// A stage as the run's projection labels it, by its Petri position: what
/// the timeline shows beside a checkpoint and what a target such as
/// `build@2` resolves through.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct StageLabel {
/// The stage id (`node@visit`) when the projection shows the stage.
pub stage_id: Option<String>,
pub node_name: String,
pub visit: u32,
}
/// The stages of a run by `(execution, firing)`.
pub type StageLabels = BTreeMap<(u64, u64), StageLabel>;
/// The Petri position of a checkpoint: the attempt whose files it holds.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct TimelinePosition {
pub execution: u64,
pub firing: u64,
pub attempt: u32,
}
/// One checkpoint of the run.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TimelineEntry {
/// 1-based, in the order the checkpoints were recorded.
pub ordinal: usize,
/// The checkpoint record's seq among the run's platform records.
pub checkpoint_seq: u64,
pub position: TimelinePosition,
/// The stage id (`node@visit`) when the projection shows the stage.
pub stage_id: Option<String>,
pub node_name: String,
pub visit: u32,
/// The Petri workspace id the commit was made in.
pub workspace: Option<String>,
pub run_commit_sha: Option<String>,
pub diff_summary: Option<DiffSummary>,
}
/// The run's checkpoints in order.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct RunTimeline {
pub entries: Vec<TimelineEntry>,
}
impl RunTimeline {
/// The timeline of `checkpoints` (the run's `checkpoint` platform
/// records, in seq order), labelled through `labels`.
#[must_use]
pub fn build(checkpoints: &[StoredPlatformRecord], labels: &StageLabels) -> Self {
let mut entries = Vec::new();
for stored in checkpoints {
let PlatformRecord::Checkpoint(record) = &stored.record else {
continue;
};
let position = position_of(record);
let label = labels.get(&(position.execution, position.firing));
entries.push(TimelineEntry {
ordinal: entries.len() + 1,
checkpoint_seq: stored.seq,
position,
stage_id: label.and_then(|label| label.stage_id.clone()),
node_name: label
.map(|label| label.node_name.clone())
.unwrap_or_default(),
visit: label.map_or(1, |label| label.visit),
workspace: record.workspace.clone(),
run_commit_sha: record.git_commit_sha.clone(),
diff_summary: record.diff_summary,
});
}
Self { entries }
}
/// The latest checkpoint, the default target.
pub fn latest(&self) -> Result<&TimelineEntry, Error> {
self.entries
.last()
.ok_or_else(|| Error::Validation("the run has no checkpoint to fork at".to_string()))
}
/// The checkpoint `target` names.
pub fn resolve(&self, target: &ForkTarget) -> Result<&TimelineEntry, Error> {
match target {
ForkTarget::Ordinal(n) => self
.entries
.iter()
.find(|entry| entry.ordinal == *n)
.ok_or_else(|| {
Error::Validation(format!(
"ordinal @{n} out of range (max @{})",
self.entries.len()
))
}),
ForkTarget::LatestVisit(name) => self
.entries
.iter()
.rev()
.find(|entry| entry.node_name == *name)
.ok_or_else(|| Error::Validation(format!("no checkpoint found for node '{name}'"))),
ForkTarget::SpecificVisit(name, visit) => self
.entries
.iter()
.find(|entry| {
entry.node_name == *name && usize::try_from(entry.visit) == Ok(*visit)
})
.ok_or_else(|| {
Error::Validation(format!("no visit {visit} found for node '{name}'"))
}),
}
}
/// The checkpoint `target` names, or the latest one.
pub fn resolve_or_latest(&self, target: Option<&ForkTarget>) -> Result<&TimelineEntry, Error> {
match target {
Some(target) => self.resolve(target),
None => self.latest(),
}
}
}
/// The position a checkpoint record names: its operation identity's
/// attempt, else the attempt it recorded, else the first.
fn position_of(record: &CheckpointRecord) -> TimelinePosition {
let attempt = match record
.operation
.as_ref()
.map(|operation| &operation.decision)
{
Some(DecisionRef::AttemptStart { attempt, .. } | DecisionRef::Route { attempt, .. }) => {
*attempt
}
Some(DecisionRef::ExecutionStart) | None => record.attempt.unwrap_or(1),
};
TimelinePosition {
execution: record.execution,
firing: record.firing,
attempt,
}
}
#[cfg(test)]
mod tests {
use fabro_store::platform_records::OperationKey;
use super::*;
fn checkpoint(
seq: u64,
execution: u64,
firing: u64,
sha: Option<&str>,
) -> StoredPlatformRecord {
StoredPlatformRecord {
seq,
recorded_at: seq * 1_000,
record: PlatformRecord::Checkpoint(CheckpointRecord {
execution,
firing,
attempt: Some(1),
workspace: Some("invocation-0-scope-0".to_string()),
git_commit_sha: sha.map(ToOwned::to_owned),
diff_summary: None,
patch_blob: None,
operation: Some(OperationKey {
execution,
decision: DecisionRef::AttemptStart { firing, attempt: 1 },
effect: "checkpoint".to_string(),
}),
}),
position: None,
}
}
fn label(node: &str, visit: u32) -> StageLabel {
StageLabel {
stage_id: Some(format!("{node}@{visit}")),
node_name: node.to_string(),
visit,
}
}
fn timeline() -> RunTimeline {
let labels: StageLabels = [
((0, 1), label("start", 1)),
((0, 2), label("build", 1)),
((0, 3), label("build", 2)),
]
.into_iter()
.collect();
RunTimeline::build(
&[
checkpoint(7, 0, 1, Some("aaa")),
checkpoint(9, 0, 2, Some("bbb")),
checkpoint(11, 0, 3, Some("ccc")),
],
&labels,
)
}
#[test]
fn a_target_parses_as_an_ordinal_a_node_or_a_visit() {
assert_eq!("@4".parse::<ForkTarget>().unwrap(), ForkTarget::Ordinal(4));
assert_eq!(
"step2".parse::<ForkTarget>().unwrap(),
ForkTarget::LatestVisit("step2".to_string())
);
assert_eq!(
"build@2".parse::<ForkTarget>().unwrap(),
ForkTarget::SpecificVisit("build".to_string(), 2)
);
assert!("@0".parse::<ForkTarget>().is_err());
assert!("@x".parse::<ForkTarget>().is_err());
}
#[test]
fn the_timeline_orders_checkpoints_and_labels_them() {
let timeline = timeline();
let ordinals: Vec<_> = timeline
.entries
.iter()
.map(|entry| (entry.ordinal, entry.node_name.as_str(), entry.visit))
.collect();
assert_eq!(ordinals, [
(1, "start", 1),
(2, "build", 1),
(3, "build", 2)
]);
assert_eq!(timeline.entries[1].checkpoint_seq, 9);
assert_eq!(timeline.entries[1].position, TimelinePosition {
execution: 0,
firing: 2,
attempt: 1,
});
assert_eq!(timeline.entries[2].stage_id.as_deref(), Some("build@2"));
}
#[test]
fn a_target_resolves_to_its_entry() {
let timeline = timeline();
assert_eq!(
timeline.resolve(&ForkTarget::Ordinal(2)).unwrap().ordinal,
2
);
assert_eq!(
timeline
.resolve(&ForkTarget::LatestVisit("build".to_string()))
.unwrap()
.ordinal,
3
);
assert_eq!(
timeline
.resolve(&ForkTarget::SpecificVisit("build".to_string(), 1))
.unwrap()
.ordinal,
2
);
assert_eq!(timeline.resolve_or_latest(None).unwrap().ordinal, 3);
assert!(matches!(
timeline.resolve(&ForkTarget::Ordinal(4)),
Err(Error::Validation(message)) if message.contains("out of range")
));
assert!(matches!(
timeline.resolve(&ForkTarget::LatestVisit("test".to_string())),
Err(Error::Validation(message)) if message.contains("no checkpoint found")
));
}
#[test]
fn an_empty_timeline_has_no_latest() {
assert!(RunTimeline::default().latest().is_err());
}
}