From 52b504d40b17e60413703a27ecb22617ad92d620 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 20 Sep 2026 14:51:04 -0400 Subject: [PATCH] Take the workspace site as a parameter instead of paired host/sandbox methods RunWorkspaces already ran every git command through a private Site (a host path or a sandbox environment) but exposed each operation twice, as commit/commit_in, matches/matches_in, has_commit/has_commit_in, reset/reset_in, restore/restore_in and workspace_head/workspace_head_in. The pairs propagated into every caller: recovery had bring_host_to and bring_sandbox_to, the hooks had snapshot and snapshot_in_sandbox, and restore_host and restore_sandbox. Site is now the public parameter, each operation exists once, and the callers collapse to one function each. The hooks resolve a scope's workspace and site in one place, site_of. Co-Authored-By: Claude Fable 5.1 --- lib/components/fabro-petri/src/checkpoint.rs | 211 ++++++++----------- lib/components/fabro-petri/src/hooks.rs | 175 +++++---------- lib/components/fabro-petri/src/recovery.rs | 75 ++----- 3 files changed, 159 insertions(+), 302 deletions(-) diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index 48e360d1b..3564c208b 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -340,53 +340,19 @@ impl RunWorkspaces { .unwrap_or(false) } - /// The host site of a workspace. - fn host(&self, workspace: &str) -> Site { + /// The host site of a workspace: where `git` runs for a workspace kept + /// on this host. + #[must_use] + pub fn host(&self, workspace: &str) -> Site { Site::Host(self.workspace_path(workspace)) } /// Commit the workspace's files on the run branch as the snapshot of /// `key`, and publish it. An earlier commit of the same key that the - /// workspace still sits on, unchanged, is reused. + /// workspace still sits on, unchanged, is reused. On a sandbox site + /// `git` runs in the scope through its environment, and the commit + /// reaches the snapshot repository as a bundle. pub async fn commit( - &self, - workspace: &str, - key: CheckpointKey, - node: &str, - status: &str, - ) -> Result { - if !self.workspace_exists(workspace).await { - return Err(CheckpointError::WorkspaceMissing { - workspace: workspace.to_string(), - path: self.workspace_path(workspace), - }); - } - self.commit_at(&self.host(workspace), workspace, key, node, status) - .await - } - - /// [`commit`](Self::commit) for a workspace inside a sandbox: `git` - /// runs in the scope through `env`, and the commit reaches the - /// snapshot repository as a bundle. - pub async fn commit_in( - &self, - env: &Arc, - workspace: &str, - key: CheckpointKey, - node: &str, - status: &str, - ) -> Result { - self.commit_at( - &Site::Sandbox(Arc::clone(env)), - workspace, - key, - node, - status, - ) - .await - } - - async fn commit_at( &self, site: &Site, workspace: &str, @@ -394,6 +360,14 @@ impl RunWorkspaces { node: &str, status: &str, ) -> Result { + if let Site::Host(path) = site { + if !fs::try_exists(path).await.unwrap_or(false) { + return Err(CheckpointError::WorkspaceMissing { + workspace: workspace.to_string(), + path: path.clone(), + }); + } + } let branched = self.ensure_repository(site).await?; if let Some(existing) = self.published_sha(workspace, key).await? { if self.head(site).await?.as_deref() == Some(existing.as_str()) @@ -575,59 +549,20 @@ impl RunWorkspaces { .is_some()) } - /// The workspace's `HEAD`, or `None` when it has no commit. - pub async fn workspace_head(&self, workspace: &str) -> Result, CheckpointError> { - self.head(&self.host(workspace)).await - } - - /// [`workspace_head`](Self::workspace_head) for a workspace inside a - /// sandbox. - pub async fn workspace_head_in( - &self, - env: &Arc, - ) -> Result, CheckpointError> { - self.head(&Site::Sandbox(Arc::clone(env))).await - } - /// Whether the workspace sits on `sha` with nothing changed since. - pub async fn matches(&self, workspace: &str, sha: &str) -> Result { - self.matches_at(&self.host(workspace), sha).await - } - - /// [`matches`](Self::matches) for a workspace inside a sandbox. - pub async fn matches_in( - &self, - env: &Arc, - sha: &str, - ) -> Result { - self.matches_at(&Site::Sandbox(Arc::clone(env)), sha).await - } - - async fn matches_at(&self, site: &Site, sha: &str) -> Result { + pub async fn matches(&self, site: &Site, sha: &str) -> Result { 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); + /// Whether the workspace's repository holds the commit `sha`, so a + /// reset can reach it (without a transfer, in a sandbox); a directory + /// that is gone or is no repository holds none. + pub async fn has_commit(&self, site: &Site, sha: &str) -> Result { + if let Site::Host(path) = site { + if !fs::try_exists(path).await.unwrap_or(false) { + 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( - &self, - env: &Arc, - sha: &str, - ) -> Result { - 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"]) .await? @@ -647,16 +582,7 @@ impl RunWorkspaces { /// Bring the workspace back to `sha`: tracked files reset, untracked /// files removed, the excluded caches left alone. - pub async fn reset(&self, workspace: &str, sha: &str) -> Result<(), CheckpointError> { - self.reset_at(&self.host(workspace), sha).await - } - - /// [`reset`](Self::reset) for a workspace inside a sandbox. - pub async fn reset_in(&self, env: &Arc, sha: &str) -> Result<(), CheckpointError> { - self.reset_at(&Site::Sandbox(Arc::clone(env)), sha).await - } - - async fn reset_at(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> { + pub async fn reset(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> { self.git(site, "reset", &["reset", "-q", "--hard", sha]) .await?; let mut clean = vec!["clean".to_string(), "-fdq".to_string()]; @@ -673,21 +599,38 @@ impl RunWorkspaces { } /// Recreate a gone workspace from the published snapshot `key`, at - /// `sha`, on the run branch. + /// `sha`, on the run branch. On a host site the workspace directory is + /// created and the snapshot fetched from the repository beside it. Into + /// a sandbox the snapshot enters the scope as a bundle of the + /// checkpoint's ref, and the workspace, fresh or stale, is fetched from + /// it and forced onto the run branch at `sha`. pub async fn restore( &self, + site: &Site, workspace: &str, key: CheckpointKey, sha: &str, ) -> Result<(), CheckpointError> { - let path = self.workspace_path(workspace); - fs::create_dir_all(&path) + match site { + Site::Host(path) => self.restore_on_host(path, workspace, key, sha).await, + Site::Sandbox(env) => self.restore_in_sandbox(env, workspace, key, sha).await, + } + } + + async fn restore_on_host( + &self, + path: &Path, + workspace: &str, + key: CheckpointKey, + sha: &str, + ) -> Result<(), CheckpointError> { + fs::create_dir_all(path) .await .map_err(|source| CheckpointError::Io { - path: path.clone(), + path: path.to_path_buf(), source, })?; - let site = Site::Host(path); + let site = Site::Host(path.to_path_buf()); self.git(&site, "init", &["init", "-q"]).await?; let repository = self.snapshot_repository(workspace); let repository = repository.to_string_lossy().into_owned(); @@ -710,10 +653,7 @@ impl RunWorkspaces { self.verify_restored(&site, sha).await } - /// [`restore`](Self::restore) into a sandbox: the snapshot enters the - /// scope as a bundle of the checkpoint's ref, and the workspace, fresh - /// or stale, is fetched from it and forced onto the run branch at `sha`. - pub async fn restore_in( + async fn restore_in_sandbox( &self, env: &Arc, workspace: &str, @@ -766,7 +706,7 @@ impl RunWorkspaces { "FETCH_HEAD", ]) .await?; - self.reset_at(&site, "HEAD").await?; + self.reset(&site, "HEAD").await?; self.verify_restored(&site, sha).await } .await; @@ -1083,7 +1023,8 @@ impl RunWorkspaces { .await } - async fn head(&self, site: &Site) -> Result, CheckpointError> { + /// The workspace's `HEAD`, or `None` when it has no commit. + pub async fn head(&self, site: &Site) -> Result, CheckpointError> { self.git_status(site, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"]) .await } @@ -1360,7 +1301,13 @@ mod tests { }; let first = workspaces - .commit(workspace, key, "build", "success") + .commit( + &workspaces.host(workspace), + workspace, + key, + "build", + "success", + ) .await .expect("the commit"); assert!(!first.reused); @@ -1370,7 +1317,13 @@ mod tests { "the first commit created the run branch in a fresh repository" ); let again = workspaces - .commit(workspace, key, "build", "success") + .commit( + &workspaces.host(workspace), + workspace, + key, + "build", + "success", + ) .await .expect("the second commit"); assert_eq!(again, Snapshot { @@ -1408,7 +1361,7 @@ mod tests { ); assert!( workspaces - .matches(workspace, &first.sha) + .matches(&workspaces.host(workspace), &first.sha) .await .expect("matches") ); @@ -1422,12 +1375,12 @@ mod tests { .expect("an untracked file"); assert!( !workspaces - .matches(workspace, &first.sha) + .matches(&workspaces.host(workspace), &first.sha) .await .expect("matches") ); workspaces - .reset(workspace, &first.sha) + .reset(&workspaces.host(workspace), &first.sha) .await .expect("the reset"); assert_eq!( @@ -1449,7 +1402,7 @@ mod tests { "the snapshot repository still knows the commit" ); workspaces - .restore(workspace, key, &first.sha) + .restore(&workspaces.host(workspace), workspace, key, &first.sha) .await .expect("the restore"); assert_eq!( @@ -1460,7 +1413,7 @@ mod tests { ); assert_eq!( workspaces - .workspace_head(workspace) + .head(&workspaces.host(workspace)) .await .expect("the head"), Some(first.sha) @@ -1483,7 +1436,13 @@ mod tests { attempt: 1, }; let error = workspaces - .commit(workspace, key, "build", "success") + .commit( + &workspaces.host(workspace), + workspace, + key, + "build", + "success", + ) .await .expect_err("the commit fails"); assert!(matches!(error, CheckpointError::Command { .. }), "{error}"); @@ -1537,7 +1496,13 @@ mod tests { attempt: 1, }; let first = workspaces - .commit(workspace, first_key, "start", "success") + .commit( + &workspaces.host(workspace), + workspace, + first_key, + "start", + "success", + ) .await .expect("the first commit"); assert_eq!( @@ -1555,7 +1520,13 @@ mod tests { attempt: 1, }; let second = workspaces - .commit(workspace, second_key, "write", "success") + .commit( + &workspaces.host(workspace), + workspace, + second_key, + "write", + "success", + ) .await .expect("the second commit"); assert_eq!(second.branched, None); diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index a9e503ff6..73dc9d6da 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -103,7 +103,8 @@ use tracing::{debug, info, warn}; use crate::blobs::Blobs; use crate::checkpoint::{ - CHECKPOINT_FAILED_CLASS, CheckpointKey, EXCLUDE_DIRS, RunWorkspaces, Snapshot, WorkspaceDiff, + CHECKPOINT_FAILED_CLASS, CheckpointKey, EXCLUDE_DIRS, RunWorkspaces, Site, Snapshot, + WorkspaceDiff, }; use crate::platform_records::PlatformRecords; use crate::recovery::{self, Plan, RestoreTarget}; @@ -390,6 +391,28 @@ impl FabroHooks { .cloned() } + /// Where `scope`'s workspace is and what to call it: on this host, the + /// directory the records name (`None` when it does not exist yet); in + /// a sandbox, the environment kept at `scope_acquired` (`None` before + /// the scope was acquired). + async fn site_of( + &self, + context: &HookContext, + scope: ScopeId, + ) -> Result, String> { + if !self.host_workspaces { + return Ok(self + .env_of(context, scope) + .map(|(workspace, env)| (workspace, Site::Sandbox(env)))); + } + let workspace = self.workspace_of(context, scope).await?; + if !self.workspaces.workspace_exists(&workspace).await { + return Ok(None); + } + let site = self.workspaces.host(&workspace); + Ok(Some((workspace, site))) + } + /// The checkpoint commit for one attempt's result. `Ok(Some)` is the /// note to record, `Ok(None)` nothing to record, `Err` the fatal /// failure message. @@ -402,15 +425,9 @@ impl FabroHooks { status: &Status, origin: ResultOrigin, ) -> Result, String> { - if !self.host_workspaces { - return self - .snapshot_in_sandbox(context, scope, key, node, status, origin) - .await; - } - let workspace = self.workspace_of(context, scope).await?; - if !self.workspaces.workspace_exists(&workspace).await { + let Some((workspace, site)) = self.site_of(context, scope).await? else { // A skipped node or a driver-made outcome may precede the scope's - // environment; nothing of the stage's is on disk to snapshot. + // environment; nothing of the stage's exists to snapshot. if origin == ResultOrigin::Driver || matches!(status, Status::Skipped) { return Ok(Some(Note::new( CHECKPOINT_NOTE, @@ -418,22 +435,21 @@ impl FabroHooks { "execution": key.execution, "firing": key.firing, "attempt": key.attempt, - "workspace": workspace, - "skipped": "the workspace does not exist yet", + "skipped": "the scope has no workspace yet", }), ))); } return Err(format!( - "the workspace `{workspace}` of scope {scope} does not exist at {}", - self.workspaces.workspace_path(&workspace).display() + "scope {scope} of execution {} has no workspace to snapshot", + context.execution )); - } + }; self.gate("commit", node).await; let serialized = self.workspace_lock(&workspace); let _held = serialized.lock().await; match self .workspaces - .commit(&workspace, key, node, status.tag()) + .commit(&site, &workspace, key, node, status.tag()) .await { Ok(snapshot) => { @@ -444,6 +460,7 @@ impl FabroHooks { firing = key.firing, attempt = key.attempt, reused = snapshot.reused, + site = ?site, "checkpoint committed" ); self.committed(key, &workspace, &snapshot).await; @@ -466,74 +483,6 @@ impl FabroHooks { } } - /// [`snapshot`](Self::snapshot) for a workspace inside the scope's - /// sandbox, through the environment kept at `scope_acquired`. - async fn snapshot_in_sandbox( - &self, - context: &HookContext, - scope: ScopeId, - key: CheckpointKey, - node: &str, - status: &Status, - origin: ResultOrigin, - ) -> Result, String> { - let Some((workspace, env)) = self.env_of(context, scope) else { - // A skipped node or a driver-made outcome may precede the scope's - // environment; nothing of the stage's exists to snapshot. - if origin == ResultOrigin::Driver || matches!(status, Status::Skipped) { - return Ok(Some(Note::new( - CHECKPOINT_NOTE, - json!({ - "execution": key.execution, - "firing": key.firing, - "attempt": key.attempt, - "skipped": "the scope has no environment yet", - }), - ))); - } - return Err(format!( - "scope {scope} of execution {} has no sandbox environment to snapshot in", - context.execution - )); - }; - self.gate("commit", node).await; - let serialized = self.workspace_lock(&workspace); - let _held = serialized.lock().await; - match self - .workspaces - .commit_in(&env, &workspace, key, node, status.tag()) - .await - { - Ok(snapshot) => { - debug!( - run_id = %self.run_id, - node, - execution = key.execution, - firing = key.firing, - attempt = key.attempt, - reused = snapshot.reused, - "checkpoint committed in the sandbox" - ); - self.committed(key, &workspace, &snapshot).await; - Ok(Some(Note::new( - CHECKPOINT_NOTE, - json!({ - "execution": key.execution, - "firing": key.firing, - "attempt": key.attempt, - "workspace": workspace, - "git_commit_sha": snapshot.sha, - "reused": snapshot.reused, - }), - ))) - } - Err(error) => Err(format!( - "the checkpoint commit of `{node}` in the sandbox failed: {}", - collect_chain(&error).join(": ") - )), - } - } - /// Remember a commit this process made, and record the run branch when /// this commit created it. async fn committed(&self, key: CheckpointKey, workspace: &str, snapshot: &Snapshot) { @@ -666,12 +615,12 @@ 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> { + /// Bring a workspace to the snapshot the resumed run's durable state + /// names, once, at its first acquisition. After a restart the server + /// already brought a host workspace there, so this verifies; a fork's + /// fresh workspace is restored here from the snapshot repository the + /// fork seeded; a sandbox workspace is only reachable here. + async fn restore(&self, workspace: &str, site: &Site) -> Result<(), ScopeAcquiredError> { let targets = self.restore_targets().await?; let target = sync::lock(targets).remove(workspace); let Some(target) = target else { @@ -679,11 +628,11 @@ impl FabroHooks { }; let serialized = self.workspace_lock(workspace); let _held = serialized.lock().await; - let action = recovery::bring_host_to(&self.workspaces, workspace, &target) + let action = recovery::bring_to(&self.workspaces, site, workspace, &target) .await .map_err(|error| { ScopeAcquiredError::new(format!( - "the host workspace `{workspace}` could not be brought to its snapshot: {}", + "the workspace `{workspace}` could not be brought to its snapshot: {}", collect_chain(&error).join(": ") )) })?; @@ -692,39 +641,8 @@ impl FabroHooks { 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( - &self, - workspace: &str, - env: &Arc, - ) -> Result<(), ScopeAcquiredError> { - let targets = self.restore_targets().await?; - let target = sync::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_sandbox_to(&self.workspaces, env, workspace, &target) - .await - .map_err(|error| { - ScopeAcquiredError::new(format!( - "the sandbox 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, - "sandbox workspace brought to its durable snapshot" + site = ?site, + "workspace brought to its durable snapshot" ); Ok(()) } @@ -1316,11 +1234,12 @@ impl ExecutionHooks for FabroHooks { if !self.resumed { return Ok(()); } - if self.host_workspaces { - self.restore_host(&workspace).await + let site = if self.host_workspaces { + self.workspaces.host(&workspace) } else { - self.restore_sandbox(&workspace, &acquired.env).await - } + Site::Sandbox(Arc::clone(&acquired.env)) + }; + self.restore(&workspace, &site).await } } diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index 2c068ed86..76d56d7ea 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -31,7 +31,7 @@ //! the worker's run reaches, so its target is deferred, and the worker's //! hooks read the same plan and apply it through the scope's environment //! at `scope_acquired`, before the first attempt runs there -//! ([`bring_sandbox_to`]). +//! ([`bring_to`]). use std::collections::BTreeMap; use std::path::PathBuf; @@ -45,11 +45,12 @@ use fabro_types::{RunId, SandboxProviderKind}; use petri_execution::host::{self, HostError}; use petri_execution::inspect::{self, ExecutionInspection, InspectError}; use petri_execution::{Access, InvocationId, RunKey, RunStore}; -use petri_runtime::executor::ExecEnv; use petri_store::StoreError; use tracing::info; -use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces}; +use crate::checkpoint::{ + CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces, Site, +}; use crate::platform_records::{PlatformRecordError, PlatformRecords}; use crate::workspace::{WorkspaceLookup, WorkspaceLookupError}; @@ -321,7 +322,13 @@ pub async fn recover(request: RecoveryRequest) -> Result Result { @@ -484,64 +493,22 @@ pub async fn bring_host_to( source, }; if workspaces - .has_commit(workspace, &target.sha) + .has_commit(site, &target.sha) .await .map_err(failed)? { if workspaces - .matches(workspace, &target.sha) + .matches(site, &target.sha) .await .map_err(failed)? { return Ok(WorkspaceAction::Verified); } - workspaces - .reset(workspace, &target.sha) - .await - .map_err(failed)?; + workspaces.reset(site, &target.sha).await.map_err(failed)?; return Ok(WorkspaceAction::Reset); } workspaces - .restore(workspace, target.key, &target.sha) - .await - .map_err(failed)?; - Ok(WorkspaceAction::Restored) -} - -/// Verify, reset or restore a sandbox workspace onto its target, through -/// the scope's environment: a retained sandbox that still holds the commit -/// is verified or reset in place; a fresh one, or one whose repository -/// lost the commit, is restored from a bundle of the snapshot. -pub async fn bring_sandbox_to( - workspaces: &RunWorkspaces, - env: &Arc, - workspace: &str, - target: &RestoreTarget, -) -> Result { - let failed = |source| RecoveryError::Workspace { - workspace: workspace.to_string(), - source, - }; - if workspaces - .has_commit_in(env, &target.sha) - .await - .map_err(failed)? - { - if workspaces - .matches_in(env, &target.sha) - .await - .map_err(failed)? - { - return Ok(WorkspaceAction::Verified); - } - workspaces - .reset_in(env, &target.sha) - .await - .map_err(failed)?; - return Ok(WorkspaceAction::Reset); - } - workspaces - .restore_in(env, workspace, target.key, &target.sha) + .restore(site, workspace, target.key, &target.sha) .await .map_err(failed)?; Ok(WorkspaceAction::Restored)