mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Seed a forked run from a source's records, checkpoints and snapshots
`fabro_petri::fork` is the seam rewind, fork and retry are built on (plan item F5.1): `fork` calls Petri's `host::fork_from` over the server's run store to seed the new run's records up to a checkpoint's position, writes the source's checkpoint records for every kept attempt under the new run at their positions, seeds the new run's snapshot repository per workspace with those checkpoints' refs alone (fetched from the source's repository under its run scratch), and records the new run branch (`fabro/run/<new id>` from the position's commit) with its Git identity. `check` refuses a position Petri would refuse (an unknown execution, or one inside a child invocation) before anything is written; `stage_labels` reads the projector's fold state so a timeline can label checkpoints and resolve `node@visit` targets. The resume then restores the fresh workspace itself: `scope_acquired` now brings a host workspace to its durable snapshot too (verified after a restart, restored from the seeded repository for a fork), through `recovery::bring_host_to`, which restores a directory with no history instead of resetting it. The projection folds `forked_from` from the fork's `run.started`. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
ac81543482
commit
0cd0df43ae
8 changed files with 645 additions and 19 deletions
|
|
@ -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<bool, CheckpointError> {
|
||||
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<dyn ExecEnv>,
|
||||
sha: &str,
|
||||
) -> Result<bool, CheckpointError> {
|
||||
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<bool, CheckpointError> {
|
||||
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}}"),
|
||||
|
|
|
|||
534
lib/components/fabro-petri/src/fork.rs
Normal file
534
lib/components/fabro-petri/src/fork.rs
Normal file
|
|
@ -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/<new id>`) 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<dyn RunStore>,
|
||||
/// The platform records of both runs.
|
||||
pub records: Arc<dyn PlatformRecords>,
|
||||
/// 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<KeptCheckpoint>,
|
||||
/// 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<KeptCheckpoint>,
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
}
|
||||
|
||||
/// 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<Forked, ForkError> {
|
||||
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<String, Vec<CheckpointKey>> = 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<CheckpointKey> {
|
||||
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<S: AsRef<str>>(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<fabro_types::ForkOrigin> {
|
||||
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<Option<fabro_types::ForkOrigin>, 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<StageLabels, ProjectError> {
|
||||
let fold_json: Option<String> =
|
||||
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<Vec<StoredPlatformRecord>, PlatformRecordError> {
|
||||
records
|
||||
.read_kind(&run_id, PlatformRecordKind::Checkpoint)
|
||||
.await
|
||||
}
|
||||
|
|
@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -381,10 +381,23 @@ impl RunView {
|
|||
|
||||
fn fold_coordinator(&mut self, record: &CoordinatorEvent, event: &RunEvent, at: DateTime<Utc>) {
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -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<Recovery, RecoveryError
|
|||
let mut recovered = Vec::new();
|
||||
for (workspace, target) in targets {
|
||||
let action = if request.host_workspaces {
|
||||
bring_to(&workspaces, &workspace, &target).await?
|
||||
bring_host_to(&workspaces, &workspace, &target).await?
|
||||
} else {
|
||||
WorkspaceAction::Deferred
|
||||
};
|
||||
|
|
@ -470,8 +470,11 @@ async fn newest(
|
|||
Ok(chosen.clone())
|
||||
}
|
||||
|
||||
/// Verify, reset or restore the host workspace onto its target.
|
||||
async fn bring_to(
|
||||
/// Verify, reset or restore the host workspace onto its target: a
|
||||
/// workspace that still holds the commit is verified or reset in place; a
|
||||
/// gone one, or a fresh directory with no history (a fork's first
|
||||
/// acquisition), is restored from the snapshot repository.
|
||||
pub async fn bring_host_to(
|
||||
workspaces: &RunWorkspaces,
|
||||
workspace: &str,
|
||||
target: &RestoreTarget,
|
||||
|
|
@ -480,7 +483,11 @@ async fn bring_to(
|
|||
workspace: workspace.to_string(),
|
||||
source,
|
||||
};
|
||||
if workspaces.workspace_exists(workspace).await {
|
||||
if workspaces
|
||||
.has_commit(workspace, &target.sha)
|
||||
.await
|
||||
.map_err(failed)?
|
||||
{
|
||||
if workspaces
|
||||
.matches(workspace, &target.sha)
|
||||
.await
|
||||
|
|
|
|||
|
|
@ -148,9 +148,9 @@ pub use run_intent::{
|
|||
TargetValidationError, ValidatedGitRunTarget, ValidatedRunTarget,
|
||||
};
|
||||
pub use run_projection::{
|
||||
CheckpointRecord, PendingInterviewRecord, RunArtifact, RunProjection, StageContextWindow,
|
||||
StageContextWindowUnavailableReason, StageInferenceProjection, StageModelUsage,
|
||||
StageProjection, StageToolBatchProjection, first_event_seq,
|
||||
CheckpointRecord, ForkOrigin, PendingInterviewRecord, RunArtifact, RunProjection,
|
||||
StageContextWindow, StageContextWindowUnavailableReason, StageInferenceProjection,
|
||||
StageModelUsage, StageProjection, StageToolBatchProjection, first_event_seq,
|
||||
};
|
||||
pub use run_sandbox::{
|
||||
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,
|
||||
|
|
|
|||
|
|
@ -47,6 +47,11 @@ pub struct RunProjection {
|
|||
pub superseded_by: Option<RunId>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub retried_from: Option<RunId>,
|
||||
/// 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<ForkOrigin>,
|
||||
/// 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<Utc>,
|
||||
}
|
||||
|
||||
/// 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(),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue