diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index 9d8ad5cb4..48e360d1b 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -607,6 +607,15 @@ impl RunWorkspaces { Ok(self.head(site).await?.as_deref() == Some(sha) && self.is_clean(site).await?) } + /// Whether the host workspace's repository holds the commit `sha`, so a + /// reset can reach it; a directory that is no repository holds none. + pub async fn has_commit(&self, workspace: &str, sha: &str) -> Result { + if !self.workspace_exists(workspace).await { + return Ok(false); + } + self.has_commit_at(&self.host(workspace), sha).await + } + /// Whether a sandbox workspace's repository holds the commit `sha`, so /// a reset can reach it without a transfer. pub async fn has_commit_in( @@ -614,16 +623,20 @@ impl RunWorkspaces { env: &Arc, sha: &str, ) -> Result { - let site = Site::Sandbox(Arc::clone(env)); + self.has_commit_at(&Site::Sandbox(Arc::clone(env)), sha) + .await + } + + async fn has_commit_at(&self, site: &Site, sha: &str) -> Result { if self - .git_status(&site, "rev-parse", &["rev-parse", "--git-dir"]) + .git_status(site, "rev-parse", &["rev-parse", "--git-dir"]) .await? .is_none() { return Ok(false); } Ok(self - .git_status(&site, "cat-file", &[ + .git_status(site, "cat-file", &[ "cat-file", "-e", &format!("{sha}^{{commit}}"), diff --git a/lib/components/fabro-petri/src/fork.rs b/lib/components/fabro-petri/src/fork.rs new file mode 100644 index 000000000..70f76dd95 --- /dev/null +++ b/lib/components/fabro-petri/src/fork.rs @@ -0,0 +1,534 @@ +//! Forking a Fabro run at a checkpoint: the seam over Petri's +//! `host::fork_from` that rewind, fork and retry are built on (the +//! integration plan's F5.1). +//! +//! Fabro's checkpoint record ties a Petri position `(execution, firing)` to +//! a Git commit. A fork seeds a new run from the source's records up to such +//! a position and leaves it ready to resume, in three steps: +//! +//! 1. Petri's `fork_from` writes the new run's records into the server's store +//! under the new run id: the same graphs, the position execution's engine +//! log cut after the firing's routing (before its first record with +//! `rerun_last`), the finished children the kept firings called, and a +//! `run.started` whose `forked_from` names the source and the position. No +//! sandbox lease is carried over, so the resume acquires the position +//! execution's scopes fresh. +//! 2. The source's checkpoint records for every attempt the fork kept are +//! written again under the new run, at their positions, and the snapshot +//! repository of every workspace they name is seeded with those checkpoints' +//! refs alone, fetched from the source's repository under the source's run +//! scratch. That is what the resume's recovery plan +//! ([`crate::recovery::plan`]) reads: the last durable finish of the +//! position execution names the snapshot the fresh workspace is restored to +//! at `scope_acquired`, on the host and in a sandbox alike. +//! 3. The fork's `run.branch` record names the run branch the restore creates +//! (`fabro/run/`) and the commit it starts from, with the +//! `git.identity` beside it, both at the checkpoint's position, so the hooks +//! record nothing twice and the run's diff is measured from the fork point. +//! +//! The Fabro run row (`run.created`, the lifecycle records) and the launch in +//! resume mode are the caller's: `fabro_workflow::operations` shapes the +//! records and the server launches the worker. + +use std::collections::{BTreeMap, BTreeSet}; +use std::path::PathBuf; +use std::process::Stdio; +use std::sync::Arc; + +use fabro_checkpoint::author::GitAuthor; +use fabro_db::DbPool; +use fabro_store::platform_records::{CheckpointRecord, GitIdentityRecord, RunBranchRecord}; +use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord}; +use fabro_types::settings::run::RunNamespace; +use fabro_types::{GitIdentity, GitIdentitySource, RunId}; +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, + StoreError as CoordinatorStoreError, +}; +use petri_runtime::ir::FiringId; +use petri_runtime::{RunOptions, Runtime}; +use petri_store::StoreError; +use tokio::fs; +use tokio::process::Command; +use tracing::{debug, info}; + +use crate::checkpoint::{CheckpointKey, RunWorkspaces}; +use crate::platform_records::{PlatformRecordError, PlatformRecords}; +use crate::projection::FoldState; +use crate::projector::ProjectError; + +/// One fork to seed. +pub struct ForkRequest { + /// The run whose records are copied. + pub source: RunId, + /// The new run's id: its Petri run key and its own run scratch. + pub fork: RunId, + /// The source's Petri run directory (its scratch's `petri`), where its + /// snapshot repositories are. + pub source_run_dir: PathBuf, + /// The fork's Petri run directory, where its snapshot repositories go. + pub fork_run_dir: PathBuf, + /// The server's run store: the source is read from it, the fork is + /// written into it. + pub store: Arc, + /// The platform records of both runs. + pub records: Arc, + /// The position the source's records are kept up to. + pub position: ForkPosition, + /// Whether the position's firing runs again (a retry of a failed + /// stage) instead of keeping its finish. + pub rerun_last: bool, + /// The run's settings, for its Git author and checkpoint settings. + pub settings: RunNamespace, +} + +/// A seeded fork, not yet resumed. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Forked { + /// The source and position, as the fork's own run declaration records + /// them. + pub origin: ForkOrigin, + /// The checkpoints the fork kept, in the source's record order. + pub checkpoints: Vec, + /// The snapshot the fork's workspace starts on, when the kept records + /// name one for the position execution's last durable finish. + pub start: Option, +} + +/// A source checkpoint the fork carries over. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct KeptCheckpoint { + pub key: CheckpointKey, + pub sha: String, + /// The Petri workspace id the commit was made in. + pub workspace: Option, +} + +/// Why the fork could not be seeded. +#[derive(Debug, thiserror::Error)] +pub enum ForkError { + #[error("the source run's record could not be opened")] + Open(#[source] StoreError), + /// Petri refused the position: an unknown execution, a firing whose + /// finish was not routed, or a position inside a child invocation (a + /// branch of a parallel node). The message says which. + #[error("{0}")] + Refused(String), + #[error("Petri could not seed the fork")] + Seed(#[source] HostError), + #[error("the fork's record could not be inspected")] + Inspect(#[source] InspectError), + #[error("the fork's coordinator log could not be read")] + Log(#[source] CoordinatorStoreError), + #[error("the checkpoint records could not be read or written")] + Records(#[source] PlatformRecordError), + #[error("the snapshot repository for `{workspace}` could not be seeded: {detail}")] + Snapshots { + workspace: String, + detail: String, + }, +} + +/// Refuse a position Petri would refuse, before anything is written for +/// the fork: an execution the source does not have, or one inside a child +/// invocation (a branch of a parallel node), whose caller's firing is live +/// at every position inside it. The messages are Petri's own. +pub async fn check( + store: &dyn RunStore, + source: RunId, + position: ForkPosition, +) -> Result<(), ForkError> { + let logs = store + .open(&RunKey::new(source.to_string()), Access::Read) + .await + .map_err(ForkError::Open)?; + 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 {}", + position.execution + ))); + }; + let invocation = execution.declaration.invocation; + if invocation != InvocationId::ROOT { + return Err(ForkError::Refused(format!( + "execution {} belongs to invocation {invocation}, not the root: a position inside a \ + child invocation cannot be forked", + position.execution + ))); + } + Ok(()) +} + +/// Seed the fork: Petri's records, then the kept checkpoints, their +/// snapshots and the run branch. The new run must not exist in the store +/// yet. +pub async fn fork(request: ForkRequest) -> Result { + let source_key = RunKey::new(request.source.to_string()); + let fork_key = RunKey::new(request.fork.to_string()); + let source_logs = request + .store + .open(&source_key, Access::Read) + .await + .map_err(ForkError::Open)?; + + let mut options = RunOptions::new(&request.fork_run_dir); + options.run_key = Some(fork_key.clone()); + let runtime = Runtime::standard() + .options(options) + .store(Arc::clone(&request.store)); + let forked = host::fork_from(&runtime, &*source_logs, request.position, ForkOptions { + rerun_last: request.rerun_last, + }) + .await + .map_err(|error| match error { + HostError::Fork(refused) => ForkError::Refused(refused.to_string()), + other => ForkError::Seed(other), + })?; + drop(source_logs); + info!( + source = %request.source, + fork = %request.fork, + position = %request.position, + rerun_last = request.rerun_last, + "Petri seeded the fork's records" + ); + + // What the fork kept: every attempt with a durable finish in its + // records, and the position execution's last one. + let fork_logs = request + .store + .open(&fork_key, Access::Read) + .await + .map_err(ForkError::Open)?; + let inspection = inspect::inspect_run(&*fork_logs) + .await + .map_err(ForkError::Inspect)?; + drop(fork_logs); + let mut kept_keys = BTreeSet::new(); + let mut start_key = None; + for execution in &inspection.executions { + let Some(engine) = execution.engine.as_ref() else { + continue; + }; + for attempt in &engine.attempts { + let key = CheckpointKey { + execution: execution.execution.raw(), + firing: attempt.firing, + attempt: attempt.attempt, + }; + kept_keys.insert(key); + if execution.execution == request.position.execution { + start_key = Some(key); + } + } + } + + // The source's checkpoint records for the kept attempts, written again + // under the fork at their positions. + let source_checkpoints = request + .records + .read_kind(&request.source, PlatformRecordKind::Checkpoint) + .await + .map_err(ForkError::Records)?; + let mut checkpoints = Vec::new(); + for stored in source_checkpoints { + let PlatformRecord::Checkpoint(record) = &stored.record else { + continue; + }; + let Some(key) = checkpoint_key(record) else { + continue; + }; + if !kept_keys.contains(&key) { + continue; + } + let Some(sha) = record.git_commit_sha.clone() else { + continue; + }; + let mut copied = record.clone(); + copied.attempt = Some(key.attempt); + copied.operation = Some(key.operation()); + request + .records + .append( + &request.fork, + &PlatformRecord::Checkpoint(copied), + Some(StagePosition { + execution: key.execution, + firing: key.firing, + }), + ) + .await + .map_err(ForkError::Records)?; + checkpoints.push(KeptCheckpoint { + key, + sha, + workspace: record.workspace.clone(), + }); + } + + // The snapshot repositories: one per workspace the kept checkpoints + // name, holding those checkpoints' refs alone. + let author = request + .settings + .git + .author + .as_ref() + .map(GitAuthor::from) + .unwrap_or_default(); + let source_workspaces = RunWorkspaces::new( + request.source_run_dir.clone(), + request.source.to_string(), + author.clone(), + &request.settings.checkpoint, + ); + let fork_workspaces = RunWorkspaces::new( + request.fork_run_dir.clone(), + request.fork.to_string(), + author.clone(), + &request.settings.checkpoint, + ); + let mut by_workspace: BTreeMap> = BTreeMap::new(); + for kept in &checkpoints { + if let Some(workspace) = &kept.workspace { + by_workspace + .entry(workspace.clone()) + .or_default() + .push(kept.key); + } + } + for (workspace, keys) in &by_workspace { + seed_snapshots(&source_workspaces, &fork_workspaces, workspace, keys).await?; + } + + // The run branch the restore creates, from the checkpoint the fork + // starts on, and the identity that authors the fork's commits. + let start = start_key.and_then(|key| checkpoints.iter().find(|kept| kept.key == key).cloned()); + if let Some(start) = &start { + let position = StagePosition { + execution: start.key.execution, + firing: start.key.firing, + }; + let branch = PlatformRecord::RunBranch(RunBranchRecord { + run_branch: Some(fork_workspaces.run_branch()), + base_sha: Some(start.sha.clone()), + workspace: start.workspace.clone(), + }); + request + .records + .append(&request.fork, &branch, Some(position)) + .await + .map_err(ForkError::Records)?; + let identity = PlatformRecord::GitIdentity(GitIdentityRecord { + identity: GitIdentity { + name: author.name.clone(), + email: author.email.clone(), + source: if author.is_default() { + GitIdentitySource::Default + } else { + GitIdentitySource::Explicit + }, + }, + }); + request + .records + .append(&request.fork, &identity, Some(position)) + .await + .map_err(ForkError::Records)?; + info!( + fork = %request.fork, + sha = start.sha, + execution = start.key.execution, + firing = start.key.firing, + "the fork's run branch starts at the position's checkpoint" + ); + } else { + debug!( + fork = %request.fork, + "the fork keeps no checkpoint; its workspace starts empty" + ); + } + + Ok(Forked { + origin: forked.origin, + checkpoints, + start, + }) +} + +/// The key a checkpoint record names: its operation identity, else its +/// position with the attempt it recorded. +fn checkpoint_key(record: &CheckpointRecord) -> Option { + record + .operation + .as_ref() + .and_then(CheckpointKey::from_operation) + .or_else(|| { + Some(CheckpointKey { + execution: record.execution, + firing: record.firing, + attempt: record.attempt?, + }) + }) +} + +/// Create the fork's bare snapshot repository for `workspace` and fetch the +/// kept checkpoints' refs into it from the source's. +async fn seed_snapshots( + source: &RunWorkspaces, + fork: &RunWorkspaces, + workspace: &str, + keys: &[CheckpointKey], +) -> Result<(), ForkError> { + let failed = |detail: String| ForkError::Snapshots { + workspace: workspace.to_string(), + detail, + }; + let source_repository = source.snapshot_repository(workspace); + if !fs::try_exists(&source_repository).await.unwrap_or(false) { + return Err(failed(format!( + "the source run has no snapshot repository at {}", + source_repository.display() + ))); + } + let repository = fork.snapshot_repository(workspace); + fs::create_dir_all(&repository).await.map_err(|error| { + failed(format!( + "{} could not be created: {error}", + repository.display() + )) + })?; + git(&repository, &["init", "-q", "--bare"]) + .await + .map_err(failed)?; + let mut args = vec![ + "fetch".to_string(), + "-q".to_string(), + source_repository.to_string_lossy().into_owned(), + ]; + for key in keys { + let name = key.snapshot_ref(); + args.push(format!("+{name}:{name}")); + } + git(&repository, &args).await.map_err(failed)?; + debug!( + workspace, + refs = keys.len(), + repository = %repository.display(), + "the fork's snapshot repository is seeded" + ); + Ok(()) +} + +/// Run `git` in `repository`; a non-zero exit is the error's detail. +async fn git>(repository: &std::path::Path, args: &[S]) -> Result<(), String> { + let output = Command::new("git") + .args(args.iter().map(AsRef::as_ref)) + .current_dir(repository) + .stdin(Stdio::null()) + .output() + .await + .map_err(|error| format!("git could not run: {error}"))?; + if output.status.success() { + Ok(()) + } else { + Err(format!( + "git {} failed ({}): {}", + args.first().map_or("", AsRef::as_ref), + output.status, + String::from_utf8_lossy(&output.stderr).trim() + )) + } +} + +/// The position a checkpoint's execution and firing name, in Petri's ids. +#[must_use] +pub fn position(execution: u64, firing: u64) -> ForkPosition { + ForkPosition { + execution: ExecutionId::new(execution), + firing: FiringId::new(firing), + } +} + +/// A fork origin as Fabro's projection shows it. +#[must_use] +pub fn origin_view(origin: &ForkOrigin) -> Option { + Some(fabro_types::ForkOrigin { + source_run_id: origin.source.to_string().parse().ok()?, + execution: origin.position.execution.raw(), + firing: origin.position.firing.raw(), + rerun_last: origin.rerun_last, + }) +} + +/// Where a run came from, when it is a fork: the `forked_from` of its run +/// declaration. `None` for a run that is not a fork, or that has no record +/// yet. +pub async fn origin_of( + store: &dyn RunStore, + run_id: RunId, +) -> Result, ForkError> { + let logs = match store + .open(&RunKey::new(run_id.to_string()), Access::Read) + .await + { + Ok(logs) => logs, + Err(StoreError::NotFound { .. }) => return Ok(None), + Err(error) => return Err(ForkError::Open(error)), + }; + let records = petri_execution::read_coordinator_log(&*logs) + .await + .map_err(ForkError::Log)?; + Ok(records.first().and_then(|record| match &record.body { + CoordinatorEvent::RunStarted { + forked_from: Some(origin), + .. + } => origin_view(origin), + _ => None, + })) +} + +/// The stages of a run by `(execution, firing)`, as the projector's fold +/// state names them: what a timeline labels its checkpoints with, and what +/// a fork target such as `build@2` resolves through. `views` is the pool +/// the view tables live in. +pub async fn stage_labels(views: &DbPool, run_id: RunId) -> Result { + let fold_json: Option = + sqlx::query_scalar("SELECT fold_json FROM petri_projection WHERE run_id = ?") + .bind(run_id.to_string()) + .fetch_optional(views) + .await + .map_err(ProjectError::Database)?; + let Some(fold_json) = fold_json else { + return Ok(BTreeMap::new()); + }; + let state: FoldState = serde_json::from_str(&fold_json).map_err(ProjectError::Encode)?; + Ok(state + .stages + .iter() + .filter_map(|(key, stage)| { + let (execution, firing) = key.split_once(':')?; + Some(( + (execution.parse().ok()?, firing.parse().ok()?), + StageLabel { + stage_id: stage.shown.then(|| stage.stage_id.to_string()), + node_name: stage.node_name.clone(), + visit: stage.visit, + }, + )) + }) + .collect()) +} + +/// The checkpoint records of a run, in seq order. +pub async fn checkpoints( + records: &dyn PlatformRecords, + run_id: RunId, +) -> Result, PlatformRecordError> { + records + .read_kind(&run_id, PlatformRecordKind::Checkpoint) + .await +} diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index e1f2288ad..e200f0a61 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -62,12 +62,14 @@ //! `git` inside the scope through it, and move the commit out as a bundle //! into the same snapshot repository the host path pushes to. Artifacts are //! read out through the same environment on every provider. The same -//! point is where a resumed run brings a sandbox workspace to the snapshot -//! its durable state names, before the first attempt runs in it: verified, +//! point is where a resumed run brings a workspace to the snapshot its +//! durable state names, before the first attempt runs in it: verified, //! reset, or, in a fresh sandbox (Petri replaces a lost one on Fabro's //! request), restored from a bundle of the checkpoint. The plan is //! [`recovery::plan`], the one the server applied to host workspaces before -//! it relaunched the worker. +//! it relaunched the worker; a host workspace is verified here, unless the +//! run is a fork whose fresh workspace nothing restored yet +//! ([`crate::fork`]), which is restored from the seeded snapshot repository. use std::collections::{BTreeMap, HashMap, HashSet}; use std::path::{Path, PathBuf}; @@ -659,6 +661,37 @@ impl FabroHooks { .await } + /// Bring a host workspace to the snapshot the resumed run's durable + /// state names, once, at its first acquisition. After a restart the + /// server already brought it there, so this verifies; a fork's fresh + /// workspace is restored here from the snapshot repository the fork + /// seeded. + async fn restore_host(&self, workspace: &str) -> Result<(), ScopeAcquiredError> { + let targets = self.restore_targets().await?; + let target = lock(targets).remove(workspace); + let Some(target) = target else { + return Ok(()); + }; + let serialized = self.workspace_lock(workspace); + let _held = serialized.lock().await; + let action = recovery::bring_host_to(&self.workspaces, workspace, &target) + .await + .map_err(|error| { + ScopeAcquiredError::new(format!( + "the host workspace `{workspace}` could not be brought to its snapshot: {}", + collect_chain(&error).join(": ") + )) + })?; + info!( + run_id = %self.run_id, + workspace, + sha = target.sha, + action = ?action, + "host workspace brought to its durable snapshot" + ); + Ok(()) + } + /// Bring a sandbox workspace to the snapshot the resumed run's durable /// state names, once, at its first acquisition. async fn restore_sandbox( @@ -1279,10 +1312,14 @@ impl ExecutionHooks for FabroHooks { (context.execution, acquired.scope), (workspace.clone(), Arc::clone(&acquired.env)), ); - if self.host_workspaces || !self.resumed { + if !self.resumed { return Ok(()); } - self.restore_sandbox(&workspace, &acquired.env).await + if self.host_workspaces { + self.restore_host(&workspace).await + } else { + self.restore_sandbox(&workspace, &acquired.env).await + } } } diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 0a4d392af..c07732ee2 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -48,7 +48,11 @@ //! - [`host_tools`]: Fabro's run tools on every native agent session of a run, //! through Petri's `HostTools` capability; //! - [`controls`]: the controls Fabro drives on a live run (pause, unpause, -//! steer, cancel), over Petri's control service. +//! steer, cancel), over Petri's control service; +//! - [`fork`]: a run seeded from another's records up to a checkpoint's +//! position, over Petri's `host::fork_from`, with the kept checkpoints, their +//! snapshots and the run branch carried over: what rewind, fork and retry are +//! built on. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. @@ -59,6 +63,7 @@ pub mod check; pub mod checkpoint; pub mod controls; pub mod engine; +pub mod fork; pub mod hooks; pub mod host_tools; pub mod http_store; diff --git a/lib/components/fabro-petri/src/projection.rs b/lib/components/fabro-petri/src/projection.rs index e096adde6..902267cf1 100644 --- a/lib/components/fabro-petri/src/projection.rs +++ b/lib/components/fabro-petri/src/projection.rs @@ -381,10 +381,23 @@ impl RunView { fn fold_coordinator(&mut self, record: &CoordinatorEvent, event: &RunEvent, at: DateTime) { match record { - CoordinatorEvent::RunStarted { root, .. } => { + CoordinatorEvent::RunStarted { + root, forked_from, .. + } => { self.state.root = Some(root.raw()); self.state.started_at = Some(event.recorded_at); if let Some(projection) = self.projection.as_mut() { + // A fork's declaration names its source; a parse failure + // means the source was not a Fabro run, which the + // projection cannot show. + projection.forked_from = forked_from.as_ref().and_then(|origin| { + Some(fabro_types::ForkOrigin { + source_run_id: origin.source.as_str().parse().ok()?, + execution: origin.position.execution.raw(), + firing: origin.position.firing.raw(), + rerun_last: origin.rerun_last, + }) + }); apply_status(projection, RunStatus::Running, at); projection.start = Some(StartRecord { start_time: at, diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index eed0f9ac9..2c068ed86 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -14,8 +14,8 @@ //! when the record was lost to the crash, the commit found by its key in the //! workspace's snapshot repository or history, which is then recorded again; //! - a workspace that survives is verified to sit on that commit, unchanged, or -//! reset to it; a workspace that is gone is restored from the run's snapshot -//! repository into a fresh directory; +//! reset to it; a workspace that is gone, or a fresh one with no history (a +//! fork's first acquisition), is restored from the run's snapshot repository; //! - a durable finish with no snapshot fails the run with a named error rather //! than resume it on stale files. //! @@ -321,7 +321,7 @@ pub async fn recover(request: RecoveryRequest) -> Result, #[serde(default, skip_serializing_if = "Option::is_none")] pub retried_from: Option, + /// Where the run's records came from when it is a fork: the source run + /// and the position its records were kept up to, as Petri's own run + /// declaration names them. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub forked_from: Option, /// The Git author/committer identity the run resolved for its commits. /// Absent until the run's first initialization resolves it. #[serde(default, skip_serializing_if = "Option::is_none")] @@ -80,6 +85,17 @@ pub struct PendingInterviewRecord { pub started_at: DateTime, } +/// The source of a forked run: the run whose records were copied, the +/// position (a firing of one of its root executions) they were kept up to, +/// and whether that firing runs again in the fork. +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct ForkOrigin { + pub source_run_id: RunId, + pub execution: u64, + pub firing: u64, + pub rerun_last: bool, +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct CheckpointRecord { pub seq: u32, @@ -712,6 +728,7 @@ impl RunProjection { pull_request_creation: None, superseded_by: None, retried_from: None, + forked_from: None, git_identity: None, pending_interviews: BTreeMap::new(), artifacts: Vec::new(),