Plan recovery through a Planner and share the run's Git settings with the hooks

RecoveryRequest::for_run and HooksSpec::for_run derived the author, the
identity source, the checkpoint settings and host_workspaces from the
run namespace with the same expressions. RunGitSettings, in checkpoint.rs
beside RunWorkspaces, is that derivation once; both specs carry it.

recovery::plan nested the "recorded, else found and reconciled" lookup
two loops deep and tracked a found flag over a tuple list. A Planner
holds the records, the workspace lookup and the recorded checkpoints,
and its methods read in order: targets, snapshot_of, reconcile_record,
newest. Candidate names the (key, sha) pair, and the Target struct that
duplicated its key's execution is gone; last_finish yields the
CheckpointKey itself, whose Display the failure reason now uses.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-20 15:15:11 -04:00
parent 4216364e92
commit 12483e081e
No known key found for this signature in database
4 changed files with 288 additions and 239 deletions

View file

@ -43,8 +43,8 @@ use std::time::Duration;
use fabro_checkpoint::author::GitAuthor;
use fabro_checkpoint::trailer::{self, Trailer};
use fabro_store::platform_records::{DecisionRef, OperationKey};
use fabro_types::DiffSummary;
use fabro_types::settings::run::RunCheckpointSettings;
use fabro_types::settings::run::{RunCheckpointSettings, RunNamespace};
use fabro_types::{DiffSummary, GitIdentitySource, SandboxProviderKind};
use petri_runtime::executor::{EnvError, ExecEnv, OutputMode, ProcessSpec, Sig};
use petri_runtime::ir::LogStream;
use tokio::process::Command;
@ -92,6 +92,54 @@ pub const EXCLUDE_DIRS: &[&str] = &[
".pytest_cache",
];
/// The settings a run's Git work runs under, as its namespace gives them:
/// who authors the checkpoint commits and where that identity came from,
/// the checkpoint settings, and whether the sandbox provider keeps the
/// workspaces on this host. The hooks and recovery both start from it.
#[derive(Clone, Debug)]
pub struct RunGitSettings {
pub author: GitAuthor,
pub identity_source: GitIdentitySource,
pub checkpoint: RunCheckpointSettings,
/// Whether the run's workspaces are on this host (the local sandbox
/// provider). A run elsewhere snapshots inside its sandboxes.
pub host_workspaces: bool,
}
impl From<&RunNamespace> for RunGitSettings {
fn from(settings: &RunNamespace) -> Self {
let author = settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default();
let identity_source = if author.is_default() {
GitIdentitySource::Default
} else {
GitIdentitySource::Explicit
};
Self {
author,
identity_source,
checkpoint: settings.checkpoint.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
}
}
}
impl Default for RunGitSettings {
/// Fabro's default author and checkpoint settings, on this host.
fn default() -> Self {
Self {
author: GitAuthor::default(),
identity_source: GitIdentitySource::Default,
checkpoint: RunCheckpointSettings::default(),
host_workspaces: true,
}
}
}
/// The identity of one snapshot: the attempt whose files it holds.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct CheckpointKey {

View file

@ -76,15 +76,12 @@ use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use fabro_checkpoint::author::GitAuthor;
use fabro_store::platform_records::{
ArtifactCollectedRecord, CheckpointRecord, GitIdentityRecord, RunBranchRecord, RunDiffRecord,
};
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition};
use fabro_types::settings::run::{RunCheckpointSettings, RunNamespace};
use fabro_types::{
BlobHash, DiffSummary, GitIdentity, GitIdentitySource, RunId, SandboxProviderKind,
};
use fabro_types::settings::run::RunNamespace;
use fabro_types::{BlobHash, DiffSummary, GitIdentity, RunId};
use fabro_util::error::collect_chain;
use fabro_util::sync;
use fabro_util::workspace_glob::{WorkspaceGlobError, WorkspaceGlobSet};
@ -103,8 +100,8 @@ use tracing::{debug, info, warn};
use crate::blobs::Blobs;
use crate::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, EXCLUDE_DIRS, RunWorkspaces, Site,
Snapshot, WorkspaceDiff,
CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, EXCLUDE_DIRS, RunGitSettings,
RunWorkspaces, Site, Snapshot, WorkspaceDiff,
};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::recovery::{self, Plan, RecoveryError, RestoreTarget};
@ -213,50 +210,27 @@ impl HookError {
}
/// What Fabro's hooks need beside the run: where the platform records go,
/// who authors the commits, the checkpoint settings, and which files are
/// the run's artifacts.
/// the run's Git settings, and which files are the run's artifacts.
pub struct HooksSpec {
pub records: Arc<dyn PlatformRecords>,
pub author: GitAuthor,
/// Where the author identity came from: the run's settings, or Fabro's
/// default.
pub identity_source: GitIdentitySource,
pub checkpoint: RunCheckpointSettings,
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
/// The `[run.artifacts] include` patterns: which files of a stage's
/// workspace are collected after the stage.
pub artifacts: Vec<String>,
/// Whether the run's workspaces are on this host (the local sandbox
/// provider). A run elsewhere snapshots inside its sandboxes.
pub host_workspaces: bool,
pub artifacts: Vec<String>,
/// A test's gate directory: a checkpoint point named by a `.hold` file
/// there waits for its `.release` file. `None` outside tests.
pub test_gates: Option<PathBuf>,
pub test_gates: Option<PathBuf>,
}
impl HooksSpec {
/// The spec a run's settings give: its Git author, its checkpoint
/// settings, its artifact patterns, and whether its sandbox provider
/// keeps workspaces on this host.
/// The spec a run's settings give: its Git settings and its artifact
/// patterns.
#[must_use]
pub fn for_run(records: Arc<dyn PlatformRecords>, settings: &RunNamespace) -> Self {
let author = settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default();
let identity_source = if author.is_default() {
GitIdentitySource::Default
} else {
GitIdentitySource::Explicit
};
Self {
records,
author,
identity_source,
checkpoint: settings.checkpoint.clone(),
git: RunGitSettings::from(settings),
artifacts: settings.artifacts.include.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
test_gates: None,
}
}
@ -432,12 +406,16 @@ impl FabroHooks {
blobs: Option<Arc<dyn Blobs>>,
) -> Self {
let identity = GitIdentity {
name: spec.author.name.clone(),
email: spec.author.email.clone(),
source: spec.identity_source,
name: spec.git.author.name.clone(),
email: spec.git.author.email.clone(),
source: spec.git.identity_source,
};
let workspaces =
RunWorkspaces::new(run_dir, run_id.to_string(), spec.author, &spec.checkpoint);
let workspaces = RunWorkspaces::new(
run_dir,
run_id.to_string(),
spec.git.author,
&spec.git.checkpoint,
);
Self {
inner,
run_id,
@ -446,7 +424,7 @@ impl FabroHooks {
workspaces,
lookup: WorkspaceLookup::new(Arc::clone(&store), run_key),
identity,
host_workspaces: spec.host_workspaces,
host_workspaces: spec.git.host_workspaces,
test_gates: spec.test_gates,
handle: OnceLock::new(),
checkpoints: CheckpointLedger::default(),

View file

@ -37,11 +37,10 @@ use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use fabro_checkpoint::author::GitAuthor;
use fabro_store::platform_records::CheckpointRecord;
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition};
use fabro_types::settings::run::{RunCheckpointSettings, RunNamespace};
use fabro_types::{RunId, SandboxProviderKind};
use fabro_types::RunId;
use fabro_types::settings::run::RunNamespace;
use petri_execution::host::{self, HostError};
use petri_execution::inspect::{self, ExecutionInspection, InspectError};
use petri_execution::{Access, InvocationId, RunKey, RunStore};
@ -49,28 +48,24 @@ use petri_store::StoreError;
use tracing::info;
use crate::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces, Site,
CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunGitSettings, RunWorkspaces, Site,
};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::workspace::{WorkspaceLookup, WorkspaceLookupError};
/// What recovery needs: the run, where its workspaces are, its records.
/// What recovery needs: the run, where its workspaces are, its records,
/// and its Git settings.
pub struct RecoveryRequest {
pub run_id: RunId,
pub run_id: RunId,
/// The run directory Petri ran under (the run's `petri` scratch).
pub run_dir: PathBuf,
pub store: Arc<dyn RunStore>,
pub records: Arc<dyn PlatformRecords>,
pub author: GitAuthor,
pub checkpoint: RunCheckpointSettings,
/// Whether the run's workspaces are on this host.
pub host_workspaces: bool,
pub run_dir: PathBuf,
pub store: Arc<dyn RunStore>,
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
}
impl RecoveryRequest {
/// The request a run's settings give: its Git author, its checkpoint
/// settings, and whether its sandbox provider keeps workspaces on this
/// host.
/// The request a run's settings give.
#[must_use]
pub fn for_run(
run_id: RunId,
@ -84,14 +79,7 @@ impl RecoveryRequest {
run_dir,
store,
records,
author: settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default(),
checkpoint: settings.checkpoint.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
git: RunGitSettings::from(settings),
}
}
}
@ -173,12 +161,6 @@ pub enum RecoveryError {
},
}
/// One live execution's last durable finish and the snapshot it names.
struct Target {
execution: u64,
key: CheckpointKey,
}
/// Decide how the run continues: the snapshot every live workspace must sit
/// on, from the records and the snapshot repository, with a lost record
/// reconciled from the repository. Nothing is touched.
@ -221,80 +203,14 @@ pub async fn plan(
if let Some(failed) = checkpoint_failure(&inspection.executions) {
return Ok(Plan::Failed { reason: failed });
}
let lookup = WorkspaceLookup::new(Arc::clone(&store), key);
let recorded = recorded_checkpoints(records, run_id).await?;
// The snapshot each live execution's workspace must sit on. A live
// execution is one whose log records no exit: `inspect_run` reports it
// as incomplete.
let mut candidates: BTreeMap<String, Vec<(Target, String)>> = BTreeMap::new();
for execution in inspection
.executions
.iter()
.filter(|execution| execution.status == "incomplete")
{
let Some(target) = last_finish(execution) else {
continue;
};
let owned = lookup
.of_invocation(execution.invocation)
.await
.map_err(RecoveryError::Lookup)?;
if owned.is_empty() {
continue;
}
let mut found = false;
for workspace in owned {
let sha = match recorded.get(&target.key) {
Some((recorded_workspace, sha))
if recorded_workspace.as_deref().is_none_or(|w| w == workspace) =>
{
Some(sha.clone())
}
_ => {
let sha = workspaces
.find(&workspace, target.key)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.clone(),
source,
})?;
if let Some(sha) = &sha {
reconcile_record(records, run_id, target.key, &workspace, sha).await?;
}
sha
}
};
if let Some(sha) = sha {
found = true;
candidates
.entry(workspace)
.or_default()
.push((Target { ..target }, sha));
}
}
if !found {
return Ok(Plan::Failed {
reason: format!(
"no checkpoint snapshot exists for the last durable finish of execution {} \
(firing {} attempt {}); the run cannot resume on stale files",
target.execution, target.key.firing, target.key.attempt
),
});
}
}
let mut targets = BTreeMap::new();
for (workspace, candidates) in candidates {
let sha = newest(workspaces, &workspace, &candidates).await?;
let key = candidates
.iter()
.find(|(_, candidate)| *candidate == sha)
.map_or(candidates[0].0.key, |(target, _)| target.key);
targets.insert(workspace, RestoreTarget { key, sha });
}
Ok(Plan::Resume { targets })
let planner = Planner {
records,
run_id,
workspaces,
lookup: WorkspaceLookup::new(store, key),
recorded: recorded_checkpoints(records, run_id).await?,
};
planner.targets(&inspection.executions).await
}
/// Decide how the run continues, and bring its host workspaces to their
@ -303,8 +219,8 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
let workspaces = RunWorkspaces::new(
request.run_dir.clone(),
request.run_id.to_string(),
request.author.clone(),
&request.checkpoint,
request.git.author.clone(),
&request.git.checkpoint,
);
let targets = match plan(
Arc::clone(&request.store),
@ -321,7 +237,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 {
let action = if request.git.host_workspaces {
bring_to(
&workspaces,
&workspaces.host(&workspace),
@ -350,6 +266,178 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
})
}
/// The snapshot one live execution's last durable finish names on a
/// workspace.
struct Candidate {
key: CheckpointKey,
sha: String,
}
/// The decision over one run's records and snapshot repository.
struct Planner<'a> {
records: &'a dyn PlatformRecords,
run_id: &'a RunId,
workspaces: &'a RunWorkspaces,
lookup: WorkspaceLookup,
/// The run's checkpoint records by key: the workspace they name and
/// the commit.
recorded: BTreeMap<CheckpointKey, (Option<String>, String)>,
}
impl Planner<'_> {
/// The snapshot each live execution's workspace must sit on. A live
/// execution is one whose log records no exit: `inspect_run` reports
/// it as incomplete. A workspace several live executions share is
/// brought to the newest of their snapshots.
async fn targets(&self, executions: &[ExecutionInspection]) -> Result<Plan, RecoveryError> {
let mut candidates: BTreeMap<String, Vec<Candidate>> = BTreeMap::new();
for execution in executions
.iter()
.filter(|execution| execution.status == "incomplete")
{
let Some(key) = last_finish(execution) else {
continue;
};
let owned = self
.lookup
.of_invocation(execution.invocation)
.await
.map_err(RecoveryError::Lookup)?;
if owned.is_empty() {
continue;
}
let mut found = false;
for workspace in owned {
if let Some(sha) = self.snapshot_of(&workspace, key).await? {
found = true;
candidates
.entry(workspace)
.or_default()
.push(Candidate { key, sha });
}
}
if !found {
return Ok(Plan::Failed {
reason: format!(
"no checkpoint snapshot exists for the last durable finish of {key}; the \
run cannot resume on stale files"
),
});
}
}
let mut targets = BTreeMap::new();
for (workspace, candidates) in candidates {
let target = self.newest(&workspace, &candidates).await?;
targets.insert(workspace, target);
}
Ok(Plan::Resume { targets })
}
/// The snapshot of `key` in `workspace`: the commit its record names,
/// when the record names this workspace or none; else the commit found
/// by its key in the workspace's snapshot repository or history, which
/// is then recorded again for the record the crash lost. `None` when no
/// snapshot exists.
async fn snapshot_of(
&self,
workspace: &str,
key: CheckpointKey,
) -> Result<Option<String>, RecoveryError> {
if let Some((recorded_workspace, sha)) = self.recorded.get(&key) {
if recorded_workspace
.as_deref()
.is_none_or(|recorded| recorded == workspace)
{
return Ok(Some(sha.clone()));
}
}
let found = self
.workspaces
.find(workspace, key)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?;
if let Some(sha) = &found {
self.reconcile_record(key, workspace, sha).await?;
}
Ok(found)
}
/// Write the record a crash lost, from the commit found by its key.
async fn reconcile_record(
&self,
key: CheckpointKey,
workspace: &str,
sha: &str,
) -> Result<(), RecoveryError> {
info!(
run_id = %self.run_id,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
sha,
"checkpoint record reconciled from the run branch"
);
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace.to_string()),
git_commit_sha: Some(sha.to_string()),
diff_summary: None,
patch_blob: None,
operation: Some(key.operation()),
});
self.records
.append(
self.run_id,
&record,
Some(StagePosition {
execution: key.execution,
firing: key.firing,
}),
)
.await
.map_err(RecoveryError::Records)?;
Ok(())
}
/// Of the snapshots live executions name on one workspace, the one
/// every other descends from, else the last named. Two executions that
/// name the same commit share it under the first one's key.
async fn newest(
&self,
workspace: &str,
candidates: &[Candidate],
) -> Result<RestoreTarget, RecoveryError> {
let mut chosen = &candidates[0];
for candidate in &candidates[1..] {
if self
.workspaces
.is_ancestor(workspace, &chosen.sha, &candidate.sha)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?
{
chosen = candidate;
}
}
let chosen = candidates
.iter()
.find(|candidate| candidate.sha == chosen.sha)
.unwrap_or(chosen);
Ok(RestoreTarget {
key: chosen.key,
sha: chosen.sha.clone(),
})
}
}
/// The reason a run with a failed checkpoint is reported failed, when it
/// has one.
fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option<String> {
@ -375,16 +463,14 @@ fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option<String> {
})
}
/// The last `StepFinished` of an execution's log.
fn last_finish(execution: &ExecutionInspection) -> Option<Target> {
/// The last `StepFinished` of an execution's log, as the key of its
/// snapshot.
fn last_finish(execution: &ExecutionInspection) -> Option<CheckpointKey> {
let attempt = execution.engine.as_ref()?.attempts.last()?;
Some(Target {
Some(CheckpointKey {
execution: execution.execution.raw(),
key: CheckpointKey {
execution: execution.execution.raw(),
firing: attempt.firing,
attempt: attempt.attempt,
},
firing: attempt.firing,
attempt: attempt.attempt,
})
}
@ -414,69 +500,6 @@ async fn recorded_checkpoints(
Ok(recorded)
}
/// Write the record a crash lost, from the commit found by its key.
async fn reconcile_record(
records: &dyn PlatformRecords,
run_id: &RunId,
key: CheckpointKey,
workspace: &str,
sha: &str,
) -> Result<(), RecoveryError> {
info!(
run_id = %run_id,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
sha,
"checkpoint record reconciled from the run branch"
);
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace.to_string()),
git_commit_sha: Some(sha.to_string()),
diff_summary: None,
patch_blob: None,
operation: Some(key.operation()),
});
records
.append(
run_id,
&record,
Some(StagePosition {
execution: key.execution,
firing: key.firing,
}),
)
.await
.map_err(RecoveryError::Records)?;
Ok(())
}
/// Of the snapshots live executions name on one workspace, the one every
/// other descends from, else the last named.
async fn newest(
workspaces: &RunWorkspaces,
workspace: &str,
targets: &[(Target, String)],
) -> Result<String, RecoveryError> {
let mut chosen = &targets[0].1;
for (_, sha) in &targets[1..] {
if workspaces
.is_ancestor(workspace, chosen, sha)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?
{
chosen = sha;
}
}
Ok(chosen.clone())
}
/// Verify, reset or restore the workspace at `site` onto its target: a
/// workspace that still holds the commit is verified or reset in place; a
/// gone one, a fresh directory with no history (a fork's first

View file

@ -24,7 +24,9 @@ use fabro_checkpoint::author::GitAuthor;
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::blobs::Blobs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use fabro_petri::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointKey, RunGitSettings, RunWorkspaces,
};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
@ -166,13 +168,13 @@ impl Harness {
fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec {
HooksSpec {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
author: GitAuthor::default(),
identity_source: GitIdentitySource::Default,
checkpoint: RunCheckpointSettings::default(),
artifacts: self.artifacts.clone(),
host_workspaces: *provider == SandboxProviderKind::LOCAL,
test_gates: None,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings {
host_workspaces: *provider == SandboxProviderKind::LOCAL,
..RunGitSettings::default()
},
artifacts: self.artifacts.clone(),
test_gates: None,
}
}
@ -285,13 +287,11 @@ impl Harness {
async fn recover(&self) -> Recovery {
recovery::recover(RecoveryRequest {
run_id: self.run_id,
run_dir: self.run_dir.clone(),
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
author: GitAuthor::default(),
checkpoint: RunCheckpointSettings::default(),
host_workspaces: true,
run_id: self.run_id,
run_dir: self.run_dir.clone(),
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings::default(),
})
.await
.expect("recovery decides")