diff --git a/lib/components/fabro-workflow/src/operations/fork.rs b/lib/components/fabro-workflow/src/operations/fork.rs new file mode 100644 index 000000000..000b7a6cd --- /dev/null +++ b/lib/components/fabro-workflow/src/operations/fork.rs @@ -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 { + 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, + pub web_url: Option, + /// The source, when the fork is a retry of it. + pub retried_from: Option, +} + +/// 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(); + } +} diff --git a/lib/components/fabro-workflow/src/operations/mod.rs b/lib/components/fabro-workflow/src/operations/mod.rs index 232ef5e00..6cfe66b24 100644 --- a/lib/components/fabro-workflow/src/operations/mod.rs +++ b/lib/components/fabro-workflow/src/operations/mod.rs @@ -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; diff --git a/lib/components/fabro-workflow/src/operations/retry.rs b/lib/components/fabro-workflow/src/operations/retry.rs new file mode 100644 index 000000000..0b493e591 --- /dev/null +++ b/lib/components/fabro-workflow/src/operations/retry.rs @@ -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, + })); + } +} diff --git a/lib/components/fabro-workflow/src/operations/rewind.rs b/lib/components/fabro-workflow/src/operations/rewind.rs new file mode 100644 index 000000000..4452d9d2c --- /dev/null +++ b/lib/components/fabro-workflow/src/operations/rewind.rs @@ -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, + }) +} diff --git a/lib/components/fabro-workflow/src/operations/timeline.rs b/lib/components/fabro-workflow/src/operations/timeline.rs new file mode 100644 index 000000000..19e00ed51 --- /dev/null +++ b/lib/components/fabro-workflow/src/operations/timeline.rs @@ -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 { + 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::() { + 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, + 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, + pub node_name: String, + pub visit: u32, + /// The Petri workspace id the commit was made in. + pub workspace: Option, + pub run_commit_sha: Option, + pub diff_summary: Option, +} + +/// The run's checkpoints in order. +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub struct RunTimeline { + pub entries: Vec, +} + +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::().unwrap(), ForkTarget::Ordinal(4)); + assert_eq!( + "step2".parse::().unwrap(), + ForkTarget::LatestVisit("step2".to_string()) + ); + assert_eq!( + "build@2".parse::().unwrap(), + ForkTarget::SpecificVisit("build".to_string(), 2) + ); + assert!("@0".parse::().is_err()); + assert!("@x".parse::().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()); + } +}