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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-20 14:51:04 -04:00
parent 1f6e53c949
commit 52b504d40b
No known key found for this signature in database
3 changed files with 159 additions and 302 deletions

View file

@ -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<Snapshot, CheckpointError> {
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<dyn ExecEnv>,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
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<Snapshot, CheckpointError> {
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<Option<String>, 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<dyn ExecEnv>,
) -> Result<Option<String>, 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<bool, CheckpointError> {
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<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
self.matches_at(&Site::Sandbox(Arc::clone(env)), sha).await
}
async fn matches_at(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
pub async fn matches(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
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);
/// 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<bool, CheckpointError> {
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<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
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"])
.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<dyn ExecEnv>, 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<dyn ExecEnv>,
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<Option<String>, CheckpointError> {
/// The workspace's `HEAD`, or `None` when it has no commit.
pub async fn head(&self, site: &Site) -> Result<Option<String>, 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);

View file

@ -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<Option<(String, Site)>, 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<Option<Note>, 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<Option<Note>, 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<dyn ExecEnv>,
) -> 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
}
}

View file

@ -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<Recovery, RecoveryError
let mut recovered = Vec::new();
for (workspace, target) in targets {
let action = if request.host_workspaces {
bring_host_to(&workspaces, &workspace, &target).await?
bring_to(
&workspaces,
&workspaces.host(&workspace),
&workspace,
&target,
)
.await?
} else {
WorkspaceAction::Deferred
};
@ -470,12 +477,14 @@ async fn newest(
Ok(chosen.clone())
}
/// Verify, reset or restore the host workspace onto its target: a
/// 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, or a fresh directory with no history (a fork's first
/// acquisition), is restored from the snapshot repository.
pub async fn bring_host_to(
/// gone one, a fresh directory with no history (a fork's first
/// acquisition), or a sandbox whose repository lost the commit, is restored
/// from the snapshot repository (a bundle of the snapshot, into a sandbox).
pub async fn bring_to(
workspaces: &RunWorkspaces,
site: &Site,
workspace: &str,
target: &RestoreTarget,
) -> Result<WorkspaceAction, RecoveryError> {
@ -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<dyn ExecEnv>,
workspace: &str,
target: &RestoreTarget,
) -> Result<WorkspaceAction, RecoveryError> {
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)