diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index df82a04e6..85a186f0d 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -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 { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 36914d942..7fa93e192 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -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, - 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, + pub git: RunGitSettings, /// The `[run.artifacts] include` patterns: which files of a stage's /// workspace are collected after the stage. - pub artifacts: Vec, - /// 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, /// 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, + pub test_gates: Option, } 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, 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>, ) -> 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(), diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index 76d56d7ea..ee62ec941 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -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, - pub records: Arc, - 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, + pub records: Arc, + 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> = 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 Result Result { + 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, 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 { + let mut candidates: BTreeMap> = 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, 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 { + 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 { @@ -375,16 +463,14 @@ fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option { }) } -/// The last `StepFinished` of an execution's log. -fn last_finish(execution: &ExecutionInspection) -> Option { +/// The last `StepFinished` of an execution's log, as the key of its +/// snapshot. +fn last_finish(execution: &ExecutionInspection) -> Option { 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 { - 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 diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 1ba189a95..5849b4633 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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, - 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, + 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, - records: Arc::clone(&self.records) as Arc, - 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, + records: Arc::clone(&self.records) as Arc, + git: RunGitSettings::default(), }) .await .expect("recovery decides")