Record the run branch, identity, artifacts and diffs from the hooks

The commit that creates a workspace's run branch records `run.branch`
(the base commit, or the first checkpoint in a workspace with no history)
and `git.identity`. Every checkpoint record after the first carries the
stage's diff from its parent commit, with the patch as a text blob. After
the checkpoint record, the transition hook lists the stage's workspace
through the scope's environment, on the host and in a sandbox alike, and
collects every file under `[run.artifacts] include` into the blob table as
an `artifact.collected` record, skipping a file already collected under
the same path and digest. At the run's end the hooks diff the branch's
last checkpoint against its base in the snapshot repository and record
`run.diff`.

`engine::retention` maps the environment's lifecycle settings onto
Petri's workspace retention instead of always keeping every workspace:
`preserve`, `stop_on_terminal = false` and the local provider keep them,
anything else keeps a failed scope's only. The hooks docs no longer name a
redundant link target, so rustdoc passes with warnings denied.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 16:00:11 -04:00
parent dc48197553
commit 67ba595b01
No known key found for this signature in database
10 changed files with 1239 additions and 129 deletions

View file

@ -205,6 +205,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
.environment
.provider
.clone(),
retention: engine::retention(&worker.run_state.spec.settings.run.environment),
cancel: cancel_token.clone(),
controls: controls.clone(),
interviewer: Arc::new(petri_interviewer),

View file

@ -318,6 +318,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
.observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))),
runtime: runtime_spec(&state, &eligible, dry_run),
provider: run_state.spec.settings.run.environment.provider.clone(),
retention: engine::retention(&run_state.spec.settings.run.environment),
cancel,
// The in-process test path drives no pause or steer: the server's
// transports for those name the worker.

View file

@ -96,18 +96,46 @@ Every adapter the integration plan describes lands here.
`petri::EVENT_CONTRACT_VERSION`.
- The platform adapters the plan adds after it: hooks and the run tools.
### What the hooks add to the records
`hooks` writes the platform records Petri cannot: `run.branch` and
`git.identity` when the first checkpoint creates the run branch (the base
commit is the workspace's `HEAD` before the branch, or that first commit in
a workspace with no history), `checkpoint` after every route with the
stage's diff from its parent commit (`diff_summary`, and the patch as a
text blob under `patch_blob`), `artifact.collected` for every file under
`[run.artifacts] include` a stage left in its workspace (the bytes go to
the blob table; a file unchanged since an earlier capture is not recorded
again), and `run.diff` at the run's end (the run branch's last checkpoint
against the base, summary and patch blob). The projection folds them into
`start`, `git_identity`, `checkpoints[].diff`, `StageProjection.diff`,
`artifacts` and `Conclusion.diff`; a patch is carried as its
`blob://sha256/<hex>` reference, which `fabro diff` and `dump` resolve.
A stage's `output` and an agent's `response` are carried the same way when
Petri offloaded them: the reference, never the bytes. A dry run's simulated
prompt or agent stage carries the stub's text as its `response`.
### What the projection leaves default
`VIEWS.md` rows with no source yet, or whose source this crate does not read
yet, keep their default value in the projection: `StageProjection.diff` and
`Conclusion.diff.patch` (the checkpoint's `patch_blob` is not resolved),
`Checkpoint`'s engine-derived maps (`completed_nodes`, `node_retries`,
`context_values`, `node_outcomes`, `next_node_id`), `agent_tools`,
`permission_level`, `script_invocation` and `script_timing`, a stage's
`notes`, `StageCompletion` details for a `parsed.note`, the sandbox instance
(the matrix's two gaps), `Run.ask_fabro`, an interview option's
`description` and `preview`, the pull request `creation` state, and the
run's notices, notifications and pairings (recorded, not shown).
yet, keep their default value in the projection: `Checkpoint`'s
engine-derived maps (`completed_nodes`, `node_retries`, `context_values`,
`node_outcomes`, `next_node_id`), `agent_tools`, `permission_level`,
`script_invocation` and `script_timing`, a stage's `notes`,
`StageCompletion` details for a `parsed.note`, the sandbox instance (the
matrix's two gaps), `Run.ask_fabro`, an interview option's `description`
and `preview`, the pull request `creation` state, and the run's notices,
notifications and pairings (recorded, not shown).
### Retention
`engine::retention` maps the run's environment settings onto Petri's
workspace retention: `preserve = true` or `stop_on_terminal = false` keeps
every workspace (`Retention::Always`), as does the local provider, whose
host workspaces live under the run's scratch directory and go with it;
otherwise a failed scope's workspace is kept for debugging and a successful
one is released (`Retention::OnFailure`, Petri's default).
Every run executes on Petri. The server side is `fabro-server`'s
`server::petri_runs`; the worker side is `fabro-cli`'s
@ -137,6 +165,13 @@ Integration tests live under `tests/`:
(`petri_testkit::run_store::conformance`) against `SqliteRunStore`, plus the
operator release, lease exclusivity, a crash between appends, and blob
interoperation with Fabro's `BlobStore`.
- `hooks.rs` runs command-only bundles through the engine assembly with
Fabro's hooks over the memory store, in-memory platform records and an
in-memory blob table: every finish is committed and recorded, a failed
stage's route sees its files, a failed checkpoint ends the run, the run
branch, identity, artifacts, per-checkpoint diffs and the run diff are
recorded, and the Docker and Daytona variants commit inside their
sandboxes.
- `interview.rs` runs human gates through the engine assembly with the
interview adapter over a control interviewer: a gate answered under the
posted id, two parallel gates each bound to their own answer, an expiry

View file

@ -147,7 +147,7 @@ stages live in the child invocation and list under the fork (see Parallel).
| notes | `StageCompletion.notes` | `step.finished` `outcome` notes; `parsed.note {result_prepared, transition}` | attempt |
| files touched | `stage.completed` `files_touched` | Pebble's fold of envelope `ToolCallCompleted` (see Agent activity) | session |
| stage diff | `StageProjection.diff` | platform record `checkpoint {execution, firing, patch_blob}` | stage |
| artifacts | `RunArtifactEntry {stage_id, node_slug, retry, relative_path, size}`, the stage artifact endpoints | `step.progress.recorded` `artifact {name, uri}`; bytes through the store | attempt |
| artifacts | `RunArtifactEntry {stage_id, node_slug, retry, relative_path, size}`, the stage artifact endpoints | platform record `artifact.collected {execution, firing, attempt, path, blob, bytes, digest}` from the `transition` hook, the bytes in the blob table; a file unchanged since an earlier capture is not recorded again | attempt |
| checkout | `setup.*` lines, `attractor.checkout` | the root `start` stage's `custom attractor.checkout {repository, commit, depth, files}` and its log lines; `[run.prepare]` commands are `run_prepare_N` stages | stage |
| hook decisions | none today | `parsed.note {kind: hook}` (`HookReport`), `custom attractor.hook` (a point a step asks itself), `parsed.hook_activity`; `run.note.recorded` for run-level points | attempt |
| budget pause | none today | `parsed.budget {state, attempt, remaining_ms, pending_questions}` | attempt |
@ -258,7 +258,8 @@ where it belongs to a stage. The proposed `record_json` fields follow.
| Fabro fact | Fields | Platform record | Keyed on |
| --- | --- | --- | --- |
| checkpoint | `checkpoints[] {seq, checkpoint, diff}`, `checkpoint.completed`, `checkpoint.failed` | `checkpoint {execution, firing, git_commit_sha, diff_summary, patch_blob}` from the `transition` hook after the commit (decision 5: a failed commit is `checkpoint_failed` on the `step.finished`, so `checkpoint.failed` needs no record); `Checkpoint`'s `completed_nodes`, `node_retries`, `node_visits`, `node_outcomes`, `context_values`, `next_node_id` and failure signatures are derived from Petri's engine state at that position | stage |
| Git commits | `git.commit {sha}`, `git.push`, `git.fetch`, `git.reset`, `RunCommit`, `RunCommitsMeta {base_sha, head_sha}` | `checkpoint {git_commit_sha}` per stage; `run.branch {run_branch, base_sha}`; push, fetch and reset are `git.push {branch, success, attempts}` only when a view needs them (none does today) | stage, run |
| Git commits | `git.commit {sha}`, `git.push`, `git.fetch`, `git.reset`, `RunCommit`, `RunCommitsMeta {base_sha, head_sha}` | `checkpoint {git_commit_sha}` per stage; `run.branch {run_branch, base_sha, workspace}` from the first checkpoint's branch creation; push, fetch and reset are `git.push {branch, success, attempts}` only when a view needs them (none does today) | stage, run |
| run diff | `Run.diff`, `Conclusion.diff` | platform record `run.diff {base_sha, head_sha, diff_summary, patch_blob}` from the `run_finished` hook: the run branch's last checkpoint against the base | run |
| files changed | `FileDiff`, `RunFilesMeta` | live: the run branch or the sandbox, from the `checkpoint` shas | run |
| pull request | `pull_request`, `pull_request_creation`, `PullRequestDetails`, `CheckRun`, `pull_request.*` | `pull_request.requested {creation_id, model, force}`, `pull_request.created {number, owner, repo, html_url, head_sha, draft}`, `pull_request.linked`, `pull_request.unlinked`, `pull_request.failed {creation_id, error}`; details and checks are live from GitHub | run |
| notifications | Slack lifecycle and interview messages | `notification.sent {route, event, channel, thread, message_id}` | run, question |

View file

@ -43,6 +43,7 @@ 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 petri_runtime::executor::{EnvError, ExecEnv, OutputMode, ProcessSpec, Sig};
use petri_runtime::ir::LogStream;
@ -64,6 +65,9 @@ pub const ATTEMPT_TRAILER: &str = "Fabro-Attempt";
const FOOTER: &str = "\u{2692}\u{fe0f} Generated with [Fabro](https://fabro.sh)";
const REFS_PREFIX: &str = "refs/checkpoints/";
/// Git's empty tree: what a root commit is diffed against.
const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904";
/// Where a bundle waits inside a sandbox on its way in or out: outside the
/// workspace, so no checkpoint ever commits it.
const TRANSFER_DIR: &str = "/tmp/fabro-snapshots";
@ -110,13 +114,20 @@ impl CheckpointKey {
/// decision in its execution, effect kind `checkpoint`.
#[must_use]
pub fn operation(self) -> OperationKey {
self.operation_for(CHECKPOINT_EFFECT)
}
/// The operation identity of another effect performed for the same
/// attempt, under `effect`.
#[must_use]
pub fn operation_for(self, effect: &str) -> OperationKey {
OperationKey {
execution: self.execution,
decision: DecisionRef::AttemptStart {
firing: self.firing,
attempt: self.attempt,
},
effect: CHECKPOINT_EFFECT.to_string(),
effect: effect.to_string(),
}
}
@ -233,12 +244,37 @@ struct GitOutput {
stderr: Vec<u8>,
}
/// A checkpoint commit: the commit, and whether an earlier attempt of the
/// same operation had already made it.
/// A checkpoint commit: the commit, whether an earlier attempt of the same
/// operation had already made it, and, when this commit created the run
/// branch in its workspace, where the branch started.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Snapshot {
pub sha: String,
pub reused: bool,
pub sha: String,
pub reused: bool,
pub branched: Option<BranchPoint>,
}
/// Where a workspace's run branch was created: the commit the workspace
/// stood on, or `None` in a repository that had no commit yet.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct BranchPoint {
pub base_sha: Option<String>,
}
/// The difference between two snapshots: the summary `git diff --numstat`
/// gives and the patch itself.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct WorkspaceDiff {
pub summary: DiffSummary,
pub patch: String,
}
impl WorkspaceDiff {
/// Whether the two snapshots hold the same tree.
#[must_use]
pub fn is_empty(&self) -> bool {
self.patch.trim().is_empty()
}
}
/// One published snapshot of a workspace.
@ -358,14 +394,15 @@ impl RunWorkspaces {
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
self.ensure_repository(site).await?;
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())
&& self.is_clean(site).await?
{
return Ok(Snapshot {
sha: existing,
sha: existing,
reused: true,
branched,
});
}
}
@ -406,7 +443,59 @@ impl RunWorkspaces {
Site::Host(path) => self.publish(workspace, path, key, &sha).await?,
Site::Sandbox(env) => self.publish_from_sandbox(env, workspace, key, &sha).await?,
}
Ok(Snapshot { sha, reused: false })
Ok(Snapshot {
sha,
reused: false,
branched,
})
}
/// The parent of a published commit, or `None` for a root commit.
pub async fn commit_parent(
&self,
workspace: &str,
sha: &str,
) -> Result<Option<String>, CheckpointError> {
let repository = Site::Host(self.ensure_snapshot_repository(workspace).await?);
self.git_status(&repository, "rev-parse", &[
"rev-parse",
"-q",
"--verify",
&format!("{sha}^"),
])
.await
}
/// The diff from `base` (the empty tree when `None`) to `head`, both
/// published in the workspace's snapshot repository.
pub async fn diff(
&self,
workspace: &str,
base: Option<&str>,
head: &str,
) -> Result<WorkspaceDiff, CheckpointError> {
let repository = Site::Host(self.ensure_snapshot_repository(workspace).await?);
let base = base.unwrap_or(EMPTY_TREE);
let numstat = self
.git(&repository, "diff --numstat", &[
"diff",
"--numstat",
"--no-color",
base,
head,
])
.await?;
let patch = self
.git(&repository, "diff", &["diff", "--no-color", base, head])
.await?;
let mut patch = patch;
if !patch.is_empty() {
patch.push('\n');
}
Ok(WorkspaceDiff {
summary: numstat_summary(&numstat),
patch,
})
}
/// The commit of `key`, from the snapshot repository first, else from
@ -721,8 +810,9 @@ impl RunWorkspaces {
}
/// A repository on the run branch, initialised when the workspace has
/// none.
async fn ensure_repository(&self, site: &Site) -> Result<(), CheckpointError> {
/// none. `Some` when the run branch was created here, with the commit
/// the workspace stood on.
async fn ensure_repository(&self, site: &Site) -> Result<Option<BranchPoint>, CheckpointError> {
if self
.git_status(site, "rev-parse", &["rev-parse", "--git-dir"])
.await?
@ -739,11 +829,13 @@ impl RunWorkspaces {
"HEAD",
])
.await?;
if current.as_deref() != Some(branch.as_str()) {
self.git(site, "checkout", &["checkout", "-q", "-B", &branch])
.await?;
if current.as_deref() == Some(branch.as_str()) {
return Ok(None);
}
Ok(())
let base_sha = self.head(site).await?;
self.git(site, "checkout", &["checkout", "-q", "-B", &branch])
.await?;
Ok(Some(BranchPoint { base_sha }))
}
/// The bare snapshot repository of the workspace, created on first use.
@ -1151,6 +1243,24 @@ impl RunWorkspaces {
}
}
/// The summary `git diff --numstat` lines add up to: one line per file,
/// `<additions>\t<deletions>\t<path>`, with `-` for a binary file.
fn numstat_summary(numstat: &str) -> DiffSummary {
let mut summary = DiffSummary::default();
for line in numstat.lines() {
let mut parts = line.splitn(3, '\t');
let (Some(additions), Some(deletions), Some(_path)) =
(parts.next(), parts.next(), parts.next())
else {
continue;
};
summary.files_changed += 1;
summary.additions += additions.parse::<i64>().unwrap_or(0);
summary.deletions += deletions.parse::<i64>().unwrap_or(0);
}
summary
}
/// The tail of git's stderr for an error message: what the run's record
/// carries about the failure, bounded.
fn detail(stderr: &[u8]) -> String {
@ -1241,14 +1351,37 @@ mod tests {
.await
.expect("the commit");
assert!(!first.reused);
assert_eq!(
first.branched,
Some(BranchPoint { base_sha: None }),
"the first commit created the run branch in a fresh repository"
);
let again = workspaces
.commit(workspace, key, "build", "success")
.await
.expect("the second commit");
assert_eq!(again, Snapshot {
sha: first.sha.clone(),
reused: true,
sha: first.sha.clone(),
reused: true,
branched: None,
});
assert_eq!(
workspaces
.commit_parent(workspace, &first.sha)
.await
.expect("the parent lookup"),
None
);
let diff = workspaces
.diff(workspace, None, &first.sha)
.await
.expect("the diff from the empty tree");
assert_eq!(diff.summary, DiffSummary {
files_changed: 1,
additions: 1,
deletions: 0,
});
assert!(diff.patch.contains("+one"), "{}", diff.patch);
assert_eq!(
workspaces.find(workspace, key).await.expect("the lookup"),
Some(first.sha.clone())
@ -1343,6 +1476,119 @@ mod tests {
assert!(matches!(error, CheckpointError::Command { .. }), "{error}");
}
#[tokio::test]
async fn a_second_commit_diffs_from_its_parent_and_a_branch_from_its_base() {
let dir = tempfile::tempdir().expect("a temp dir");
let workspaces = workspaces(dir.path());
let workspace = "invocation-0-scope-0";
let path = workspaces.workspace_path(workspace);
fs::create_dir_all(&path).await.expect("the workspace");
fs::write(path.join("story.txt"), "line 1\n")
.await
.expect("a file");
// A source repository with a commit: the run branch starts from it.
for args in [vec!["init", "-q"], vec!["add", "."], vec![
"-c",
"user.name=t",
"-c",
"user.email=t@example.com",
"commit",
"-q",
"-m",
"initial",
]] {
let status = Command::new("git")
.args(&args)
.current_dir(&path)
.status()
.await
.expect("git runs");
assert!(status.success(), "git {args:?}");
}
let base = String::from_utf8(
Command::new("git")
.args(["rev-parse", "HEAD"])
.current_dir(&path)
.output()
.await
.expect("git runs")
.stdout,
)
.expect("utf-8")
.trim()
.to_string();
let first_key = CheckpointKey {
execution: 0,
firing: 1,
attempt: 1,
};
let first = workspaces
.commit(workspace, first_key, "start", "success")
.await
.expect("the first commit");
assert_eq!(
first.branched,
Some(BranchPoint {
base_sha: Some(base.clone()),
})
);
fs::write(path.join("story.txt"), "line 1\nline 2\n")
.await
.expect("a change");
let second_key = CheckpointKey {
execution: 0,
firing: 2,
attempt: 1,
};
let second = workspaces
.commit(workspace, second_key, "write", "success")
.await
.expect("the second commit");
assert_eq!(second.branched, None);
assert_eq!(
workspaces
.commit_parent(workspace, &second.sha)
.await
.expect("the parent lookup"),
Some(first.sha.clone())
);
let stage = workspaces
.diff(workspace, Some(&first.sha), &second.sha)
.await
.expect("the stage diff");
assert_eq!(stage.summary, DiffSummary {
files_changed: 1,
additions: 1,
deletions: 0,
});
assert!(stage.patch.contains("+line 2"), "{}", stage.patch);
let run = workspaces
.diff(workspace, Some(&base), &second.sha)
.await
.expect("the run diff");
assert_eq!(run.summary, stage.summary);
let unchanged = workspaces
.diff(workspace, Some(&base), &first.sha)
.await
.expect("the empty diff");
assert!(unchanged.is_empty());
assert_eq!(unchanged.summary, DiffSummary::default());
}
#[test]
fn numstat_lines_add_up_and_binary_files_count_as_changed() {
assert_eq!(
numstat_summary("3\t1\ta.txt\n-\t-\timage.png\n"),
DiffSummary {
files_changed: 2,
additions: 3,
deletions: 1,
}
);
assert_eq!(numstat_summary(""), DiffSummary::default());
}
#[test]
fn detail_keeps_the_tail_of_long_output() {
let long = "x".repeat(600);

View file

@ -23,9 +23,10 @@
//! hook service: the checkpoint commit before every durable finish and its
//! platform record after every route, with a failed commit ending the run
//! as a `checkpoint_failed` failure. What the standalone runner's defaults
//! give the run: Petri's local hook service for `[[run.hooks]]`, no host
//! tools, and `Retention::Always` for every workspace, Fabro's default.
//! Cancellation rides the caller's token: when it fires, the root
//! give the run: Petri's local hook service for `[[run.hooks]]` and no host
//! tools. The workspaces' retention comes from the run's environment
//! settings through [`retention`]. Cancellation rides the caller's token:
//! when it fires, the root
//! invocation is cancelled politely and Petri records why. The run's other
//! controls (pause, unpause, steer) are the caller's [`RunControls`]: its
//! pause gate is installed over the run's hooks, it observes the run, and
@ -47,6 +48,7 @@
use std::path::PathBuf;
use std::sync::Arc;
use fabro_types::settings::run::RunEnvironmentSettings;
use fabro_types::{FailureReason, RunId, SandboxProviderKind};
use petri_execution::host::{self, HostError, HostRun};
use petri_execution::inspect::{self, InspectError, RunInspection};
@ -55,7 +57,8 @@ use petri_execution::{
RECEIPT_FILE, RunKey, RunStore,
};
use petri_runtime::driver::lifecycle::ExecutionHooks;
use petri_runtime::executor::{Retention, SecretProvider};
pub use petri_runtime::executor::Retention;
use petri_runtime::executor::SecretProvider;
use petri_runtime::{LostSandbox, RunOptions, SandboxBackend};
use tokio::fs;
use tokio_util::sync::CancellationToken;
@ -93,6 +96,8 @@ pub struct RunRequest {
pub runtime: RuntimeSpec,
/// The sandbox provider Fabro resolved for the run's environment.
pub provider: SandboxProviderKind,
/// When the run's workspaces are kept after their scope is released.
pub retention: Retention,
/// Fires to cancel the run.
pub cancel: CancellationToken,
/// The run's pause, unpause and steer controls, which the caller keeps
@ -173,7 +178,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
let key = RunKey::new(request.run_id.as_str());
let mut options = RunOptions::new(&request.run_dir);
options.run_key = Some(key.clone());
options.retention = Retention::Always;
options.retention = request.retention;
options.sandbox.backend = backend;
// Fabro's hooks restore a sandbox workspace from its snapshots at the
// scope's acquisition, so a lease whose sandbox is gone gets a fresh
@ -190,8 +195,8 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
if let Some(secrets) = request.secrets {
runtime = runtime.secrets(SharedSecrets(secrets));
}
if let Some(blobs) = request.blobs {
runtime = runtime.capability(RunBlobs::output_store(blobs));
if let Some(blobs) = &request.blobs {
runtime = runtime.capability(RunBlobs::output_store(Arc::clone(blobs)));
}
let fabro_hooks = request.hooks.map(|spec| {
let inner = runtime
@ -206,6 +211,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
request.run_dir.clone(),
Arc::clone(&request.store),
resumed,
request.blobs.clone(),
))
});
if let Some(hooks) = &fabro_hooks {
@ -279,6 +285,31 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
Ok(outcome)
}
/// When Petri keeps a run's workspaces after their scope is released, from
/// the run's environment settings:
///
/// - `[environments.<id>.lifecycle] preserve = true` asks for the sandbox to
/// stay after the run, so every workspace is kept (`Retention::Always`).
/// - The local provider keeps every workspace too: a host workspace lives under
/// the run's own scratch directory, which `fabro system prune` removes with
/// the run, and the legacy executor never removed it on its own.
/// - `stop_on_terminal = false` asks for the sandbox to outlive the run, so its
/// workspaces are kept (`Retention::Always`).
/// - Otherwise the sandbox is released with the run and Petri's default
/// applies: a failed scope's workspace is kept for debugging, a successful
/// one is not (`Retention::OnFailure`).
#[must_use]
pub fn retention(environment: &RunEnvironmentSettings) -> Retention {
let keep = environment.lifecycle.preserve
|| !environment.lifecycle.stop_on_terminal
|| environment.provider == SandboxProviderKind::LOCAL;
if keep {
Retention::Always
} else {
Retention::OnFailure
}
}
/// The Fabro run id the run key names. A key that is not one (a test's
/// bare key) still gets hooks, under a fresh id for its platform records.
fn spec_run_id(run_id: &str) -> RunId {
@ -460,8 +491,41 @@ async fn write_receipt(run_dir: &std::path::Path, receipt: &petri_execution::Int
#[cfg(test)]
mod tests {
use fabro_types::settings::run::EnvironmentLifecycleSettings;
use super::*;
fn environment(provider: SandboxProviderKind) -> RunEnvironmentSettings {
let mut environment = RunEnvironmentSettings::from_environment(
"test".to_string(),
fabro_types::settings::run::EnvironmentSettings::default(),
);
environment.provider = provider;
environment
}
#[test]
fn retention_follows_the_environment_lifecycle() {
let mut docker = environment(SandboxProviderKind::DOCKER);
assert_eq!(retention(&docker), Retention::OnFailure);
docker.lifecycle = EnvironmentLifecycleSettings {
preserve: true,
stop_on_terminal: true,
auto_stop: None,
};
assert_eq!(retention(&docker), Retention::Always);
docker.lifecycle = EnvironmentLifecycleSettings {
preserve: false,
stop_on_terminal: false,
auto_stop: None,
};
assert_eq!(retention(&docker), Retention::Always);
assert_eq!(
retention(&environment(SandboxProviderKind::LOCAL)),
Retention::Always
);
}
fn outcome_with(status: RunStatus, failure: Option<&str>, complete: bool) -> RunOutcome {
RunOutcome {
status,

View file

@ -1,6 +1,7 @@
//! Fabro's awaited extension points on a Petri run: the checkpoint commit,
//! its platform record, and the run-level ends, wrapped around Petri's own
//! hook service so `[[run.hooks]]` keep running.
//! its platform record, the artifacts a stage leaves behind, the run's diff,
//! and the run-level ends, wrapped around Petri's own hook service so
//! `[[run.hooks]]` keep running.
//!
//! [`FabroHooks`] implements Petri's `ExecutionHooks` and is installed with
//! `Runtime::hooks` by [`engine::run`](crate::engine::run). It holds the
@ -15,26 +16,41 @@
//! cancelled attempt is not. A failed commit is fatal to the run: the outcome
//! becomes a failure of class `checkpoint_failed`, the run is cancelled
//! through the coordinator handle, and `transition` refuses the firing's
//! routes, so no route is taken.
//! routes, so no route is taken. The commit that creates the run branch also
//! records where it started: the `run.branch` platform record (the branch
//! name and the base commit) and the `git.identity` record (who authors the
//! commits, and where that identity came from).
//! - `transition`: the platform checkpoint record, keyed on the Petri position
//! and the checkpoint's operation identity. A failed write is a recorded
//! problem on the transition, never a blocked route.
//! - `run_finished` and `scope_released`: forwarded, so the local service runs
//! `run_complete`, `run_failed` and `sandbox_cleanup` with the sandbox in
//! place. Fabro's own end-of-run work (the terminal lifecycle event,
//! notifications on it) is the run lifecycle path's, on the worker's and
//! server's side of the engine, and the workspace's retention is Petri's
//! (`Retention::Always`).
//! and the checkpoint's operation identity, with the stage's diff from its
//! parent commit (`diff_summary`, and the patch as a blob); then the stage's
//! artifacts: every file under `[run.artifacts] include` in the stage's
//! workspace goes to the blob table and gets an `artifact.collected` record,
//! unless the same file with the same content was already collected earlier
//! in the run. A failed write is a recorded problem on the transition, never
//! a blocked route.
//! - `run_finished`: the run's diff, its run branch against its base commit, as
//! the `run.diff` platform record with the patch as a blob; then the
//! forwarded point, so the local service runs `run_complete` and `run_failed`
//! with the sandbox in place.
//! - `scope_released`: forwarded, so the local service runs `sandbox_cleanup`
//! with the sandbox in place. Fabro's own end-of-run work (the terminal
//! lifecycle event, notifications on it) is the run lifecycle path's, on the
//! worker's and server's side of the engine, and the workspace's retention is
//! Petri's, mapped from the run's environment settings by
//! [`engine::retention`](crate::engine::retention).
//!
//! # Operation identities
//!
//! Every external effect here is keyed on `(run key, execution, DecisionId,
//! effect kind)` from the hook context and deduplicated on retry: the
//! checkpoint's key is the attempt's decision in its execution, effect
//! `checkpoint`. A re-dispatched attempt whose commit already landed
//! reuses it when the workspace still sits on it unchanged (see
//! [`RunWorkspaces::commit`]); a reissued routing decision finds the
//! record, or the commit by its trailers, and writes nothing twice.
//! `checkpoint`; an artifact's is the same decision, effect `artifact`, with
//! the file's path and content digest as the identity within it. A
//! re-dispatched attempt whose commit already landed reuses it when the
//! workspace still sits on it unchanged (see [`RunWorkspaces::commit`]); a
//! reissued routing decision finds the record, or the commit by its
//! trailers, and writes nothing twice; a file already collected under the
//! same path and digest is not collected again.
//!
//! # Where the workspace is
//!
@ -43,25 +59,31 @@
//! On Docker or Daytona the workspace lives inside the scope's sandbox: the
//! hooks keep the environment Petri hands them at `scope_acquired`, run
//! `git` inside the scope through it, and move the commit out as a bundle
//! into the same snapshot repository the host path pushes to. The same
//! 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,
//! 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`](crate::recovery::plan), the one the server applied
//! to host workspaces before it relaunched the worker.
//! [`recovery::plan`], the one the server applied to host workspaces before
//! it relaunched the worker.
use std::collections::{BTreeMap, HashMap, HashSet};
use std::path::PathBuf;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError};
use std::time::Duration;
use fabro_checkpoint::author::GitAuthor;
use fabro_store::platform_records::CheckpointRecord;
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::{RunId, SandboxProviderKind};
use fabro_types::{
BlobHash, DiffSummary, GitIdentity, GitIdentitySource, RunId, SandboxProviderKind,
};
use fabro_util::error::collect_chain;
use fabro_util::workspace_glob::{WorkspaceGlobError, WorkspaceGlobSet};
use petri_execution::{CancelReason, CoordinatorHandle, InvocationId, RunKey, RunStore};
use petri_runtime::driver::lifecycle::{
AdmitAttempt, AttemptDecision, ExecutionHooks, HookContext, Note, PrepareError, PrepareResult,
@ -75,7 +97,10 @@ use tokio::sync::{Mutex as AsyncMutex, OnceCell};
use tokio::{fs, time};
use tracing::{debug, info, warn};
use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use crate::blobs::Blobs;
use crate::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointKey, EXCLUDE_DIRS, RunWorkspaces, Snapshot, WorkspaceDiff,
};
use crate::platform_records::PlatformRecords;
use crate::recovery::{self, Plan, RestoreTarget};
use crate::workspace::{self, WorkspaceLookup};
@ -83,15 +108,36 @@ use crate::workspace::{self, WorkspaceLookup};
/// The note kind the hooks record on a firing about its checkpoint.
pub const CHECKPOINT_NOTE: &str = "fabro.checkpoint";
/// The effect kind of an artifact collection in its operation identity.
pub const ARTIFACT_EFFECT: &str = "artifact";
/// How often a held checkpoint polls its test gate.
const GATE_POLL: Duration = Duration::from_millis(50);
/// The most files one stage's collection keeps, the legacy executor's
/// budget.
const ARTIFACT_MAX_FILES: usize = 100;
/// The largest file collected, the legacy executor's budget.
const ARTIFACT_MAX_FILE_BYTES: u64 = 10 * 1024 * 1024;
/// The most bytes one stage's collection keeps, the legacy executor's
/// budget.
const ARTIFACT_MAX_TOTAL_BYTES: u64 = 50 * 1024 * 1024;
/// How deep a traversal root is listed.
const ARTIFACT_LIST_DEPTH: usize = 64;
/// What Fabro's hooks need beside the run: where the platform records go,
/// who authors the commits, and the checkpoint settings.
/// who authors the commits, the checkpoint 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,
/// 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,
@ -102,19 +148,27 @@ pub struct HooksSpec {
impl HooksSpec {
/// The spec a run's settings give: its Git author, its checkpoint
/// settings, and whether its sandbox provider keeps workspaces on this
/// host.
/// settings, its artifact patterns, and whether its sandbox provider
/// keeps workspaces on this host.
#[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: settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default(),
author,
identity_source,
checkpoint: settings.checkpoint.clone(),
artifacts: settings.artifacts.include.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
test_gates: None,
}
@ -128,45 +182,68 @@ impl HooksSpec {
}
/// A scope's sandbox environment as the hooks keep it: the workspace id
/// the executor named, and the environment `git` runs in.
/// the executor named, and the environment `git` runs in and files are
/// read through.
type AcquiredEnv = (String, Arc<dyn ExecEnv>);
/// The identity of a collected file: its path and content digest.
type ArtifactIdentity = (String, String);
/// The last checkpoint recorded: its workspace and commit.
type LastCheckpoint = (String, String);
/// Fabro's `ExecutionHooks`, around the hooks the runtime installed.
pub struct FabroHooks {
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
host_workspaces: bool,
test_gates: Option<PathBuf>,
handle: OnceLock<CoordinatorHandle>,
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
/// Where an artifact's bytes and a diff's patch go; `None` records
/// summaries alone.
blobs: Option<Arc<dyn Blobs>>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
identity: GitIdentity,
artifact_globs: Result<WorkspaceGlobSet, WorkspaceGlobError>,
host_workspaces: bool,
test_gates: Option<PathBuf>,
handle: OnceLock<CoordinatorHandle>,
/// The workspace and commit of every checkpoint this process made.
committed: Mutex<HashMap<CheckpointKey, (String, String)>>,
committed: Mutex<HashMap<CheckpointKey, (String, String)>>,
/// Which checkpoints have their platform record, loaded from the store
/// once and kept up to date with every append.
recorded: Mutex<HashSet<CheckpointKey>>,
recorded_loaded: OnceCell<()>,
recorded: Mutex<HashSet<CheckpointKey>>,
recorded_loaded: OnceCell<()>,
/// The workspace and commit of the checkpoint recorded last: the head
/// the run's diff is measured to.
last_checkpoint: Mutex<Option<LastCheckpoint>>,
/// The run branch as recorded, once: read from the store, or written
/// by the commit that created the branch.
branch: OnceCell<RunBranchRecord>,
/// Every artifact collected so far, by path and digest, loaded from the
/// store once and kept up to date with every append.
collected: Mutex<HashSet<ArtifactIdentity>>,
collected_loaded: OnceCell<()>,
/// Inherited workspaces resolved through the run's records.
inherited: Mutex<HashMap<InvocationId, Option<String>>>,
inherited: Mutex<HashMap<InvocationId, Option<String>>>,
/// One lock per workspace: the branches of a parallel node and a nested
/// invocation share their caller's workspace, and Git allows one index
/// operation at a time in it.
workspace_locks: Mutex<HashMap<String, Arc<AsyncMutex<()>>>>,
workspace_locks: Mutex<HashMap<String, Arc<AsyncMutex<()>>>>,
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<Option<String>>,
/// The sandbox environment of every acquired scope, by execution and
/// scope, with the workspace id the executor named: where `git` runs
/// when the workspaces are not on this host. Dropped at release.
envs: Mutex<HashMap<(ExecutionId, ScopeId), AcquiredEnv>>,
failure: Mutex<Option<String>>,
/// The environment of every acquired scope, by execution and scope,
/// with the workspace id the executor named: where `git` runs when the
/// workspaces are not on this host, and where artifacts are read from
/// on every provider. Dropped at release.
envs: Mutex<HashMap<(ExecutionId, ScopeId), AcquiredEnv>>,
/// Whether the run continues from its records: a sandbox workspace is
/// then brought to its snapshot when its scope is first acquired.
resumed: bool,
resumed: bool,
/// The snapshot every live sandbox workspace must sit on before work
/// resumes in it, read once from the records; an entry leaves when it
/// is applied.
restore: OnceCell<Mutex<BTreeMap<String, RestoreTarget>>>,
store: Arc<dyn RunStore>,
restore: OnceCell<Mutex<BTreeMap<String, RestoreTarget>>>,
store: Arc<dyn RunStore>,
}
impl FabroHooks {
@ -174,7 +251,8 @@ impl FabroHooks {
/// run whose records are in `store` under `run_key`, with its
/// workspaces under `run_dir`. `resumed` says the run continues from
/// its records, so a sandbox workspace is brought to its snapshot at
/// its scope's first acquisition.
/// its scope's first acquisition. `blobs` is where artifact bytes and
/// diff patches go.
#[must_use]
pub fn new(
spec: HooksSpec,
@ -184,21 +262,34 @@ impl FabroHooks {
run_dir: PathBuf,
store: Arc<dyn RunStore>,
resumed: bool,
blobs: Option<Arc<dyn Blobs>>,
) -> Self {
let identity = GitIdentity {
name: spec.author.name.clone(),
email: spec.author.email.clone(),
source: spec.identity_source,
};
let workspaces =
RunWorkspaces::new(run_dir, run_id.to_string(), spec.author, &spec.checkpoint);
Self {
inner,
run_id,
records: spec.records,
blobs,
workspaces,
lookup: WorkspaceLookup::new(Arc::clone(&store), run_key),
identity,
artifact_globs: WorkspaceGlobSet::try_new(&spec.artifacts),
host_workspaces: spec.host_workspaces,
test_gates: spec.test_gates,
handle: OnceLock::new(),
committed: Mutex::default(),
recorded: Mutex::default(),
recorded_loaded: OnceCell::new(),
last_checkpoint: Mutex::default(),
branch: OnceCell::new(),
collected: Mutex::default(),
collected_loaded: OnceCell::new(),
inherited: Mutex::default(),
workspace_locks: Mutex::default(),
failure: Mutex::default(),
@ -285,6 +376,12 @@ impl FabroHooks {
Ok(inherited.unwrap_or(isolated))
}
/// The environment of `scope` in the context's execution, as
/// `scope_acquired` kept it, with the workspace id the executor named.
fn env_of(&self, context: &HookContext, scope: ScopeId) -> Option<AcquiredEnv> {
lock(&self.envs).get(&(context.execution, scope)).cloned()
}
/// 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.
@ -341,7 +438,7 @@ impl FabroHooks {
reused = snapshot.reused,
"checkpoint committed"
);
lock(&self.committed).insert(key, (workspace.clone(), snapshot.sha.clone()));
self.committed(key, &workspace, &snapshot).await;
Ok(Some(Note::new(
CHECKPOINT_NOTE,
json!({
@ -372,8 +469,7 @@ impl FabroHooks {
status: &Status,
origin: ResultOrigin,
) -> Result<Option<Note>, String> {
let held = lock(&self.envs).get(&(context.execution, scope)).cloned();
let Some((workspace, env)) = held else {
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) {
@ -410,7 +506,7 @@ impl FabroHooks {
reused = snapshot.reused,
"checkpoint committed in the sandbox"
);
lock(&self.committed).insert(key, (workspace.clone(), snapshot.sha.clone()));
self.committed(key, &workspace, &snapshot).await;
Ok(Some(Note::new(
CHECKPOINT_NOTE,
json!({
@ -430,6 +526,96 @@ impl FabroHooks {
}
}
/// 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) {
lock(&self.committed).insert(key, (workspace.to_string(), snapshot.sha.clone()));
let Some(branched) = &snapshot.branched else {
return;
};
// A branch that starts from nothing (a workspace with no history) is
// measured from its first commit: the checkout the run started on.
let base_sha = branched
.base_sha
.clone()
.unwrap_or_else(|| snapshot.sha.clone());
if let Err(error) = self.record_branch(workspace, base_sha).await {
warn!(run_id = %self.run_id, error = %error, "the run branch was not recorded");
}
}
/// The `run.branch` and `git.identity` records, once per run: the first
/// workspace to create the run branch names where it started. A run
/// that already recorded its branch (a resume, or a nested workspace
/// after the root's) records nothing.
async fn record_branch(&self, workspace: &str, base_sha: String) -> Result<(), String> {
let branch = self
.branch
.get_or_try_init(|| async {
if let Some(stored) = self.stored_branch().await? {
return Ok::<_, String>(stored);
}
let record = RunBranchRecord {
run_branch: Some(self.workspaces.run_branch()),
base_sha: Some(base_sha.clone()),
workspace: Some(workspace.to_string()),
};
self.records
.append(
&self.run_id,
&PlatformRecord::RunBranch(record.clone()),
None,
)
.await
.map_err(|error| {
format!(
"the run branch record could not be written: {}",
collect_chain(&error).join(": ")
)
})?;
let identity = PlatformRecord::GitIdentity(GitIdentityRecord {
identity: self.identity.clone(),
});
self.records
.append(&self.run_id, &identity, None)
.await
.map_err(|error| {
format!(
"the git identity record could not be written: {}",
collect_chain(&error).join(": ")
)
})?;
info!(
run_id = %self.run_id,
workspace,
base_sha,
"run branch recorded"
);
Ok(record)
})
.await?;
debug!(run_id = %self.run_id, base_sha = ?branch.base_sha, "the run branch is recorded");
Ok(())
}
/// The run branch the store already holds, when a record exists.
async fn stored_branch(&self) -> Result<Option<RunBranchRecord>, String> {
let stored = self
.records
.read_kind(&self.run_id, PlatformRecordKind::RunBranch)
.await
.map_err(|error| {
format!(
"the run's branch record could not be read: {}",
collect_chain(&error).join(": ")
)
})?;
Ok(stored.into_iter().find_map(|stored| match stored.record {
PlatformRecord::RunBranch(record) => Some(record),
_ => None,
}))
}
/// The restore plan of a resumed run, read once: what every live
/// sandbox workspace must be brought to at its first acquisition.
async fn restore_targets(
@ -491,7 +677,8 @@ impl FabroHooks {
Ok(())
}
/// The checkpoint's platform record, once per operation identity.
/// The checkpoint's platform record, once per operation identity, with
/// the stage's diff from the commit's parent.
async fn record(
&self,
context: &HookContext,
@ -508,9 +695,7 @@ impl FabroHooks {
let (workspace, sha) = if let Some(committed) = committed {
committed
} else {
let acquired = lock(&self.envs)
.get(&(context.execution, scope))
.map(|(workspace, _)| workspace.clone());
let acquired = self.env_of(context, scope).map(|(workspace, _)| workspace);
let workspace = match acquired {
Some(workspace) => workspace,
None => self.workspace_of(context, scope).await?,
@ -534,15 +719,29 @@ impl FabroHooks {
})?;
(workspace, sha)
};
let (diff_summary, patch_blob) = match self.stage_diff(&workspace, &sha).await {
Ok(diff) => diff,
Err(error) => {
// The record still names the commit; the diff is a view.
warn!(
run_id = %self.run_id,
workspace,
sha,
error = %error,
"the checkpoint's diff was not computed"
);
(None, None)
}
};
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace),
git_commit_sha: Some(sha),
diff_summary: None,
patch_blob: None,
operation: Some(key.operation()),
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace.clone()),
git_commit_sha: Some(sha.clone()),
diff_summary,
patch_blob,
operation: Some(key.operation()),
});
self.records
.append(
@ -561,11 +760,56 @@ impl FabroHooks {
)
})?;
lock(&self.recorded).insert(key);
*lock(&self.last_checkpoint) = Some((workspace, sha));
Ok(())
}
/// A stage's diff: its checkpoint commit against the commit's parent.
/// A root commit (the first snapshot of a workspace with no history)
/// has none. The patch goes to the blob table when the run has one and
/// the diff is not empty.
async fn stage_diff(
&self,
workspace: &str,
sha: &str,
) -> Result<(Option<DiffSummary>, Option<BlobHash>), String> {
let parent = self
.workspaces
.commit_parent(workspace, sha)
.await
.map_err(|error| collect_chain(&error).join(": "))?;
let Some(parent) = parent else {
return Ok((None, None));
};
let diff = self
.workspaces
.diff(workspace, Some(&parent), sha)
.await
.map_err(|error| collect_chain(&error).join(": "))?;
let patch_blob = self.patch_blob(&diff).await?;
Ok((Some(diff.summary), patch_blob))
}
/// The patch of a diff in the blob table, when the diff is not empty
/// and the run has a blob table.
async fn patch_blob(&self, diff: &WorkspaceDiff) -> Result<Option<BlobHash>, String> {
if diff.is_empty() {
return Ok(None);
}
let Some(blobs) = &self.blobs else {
return Ok(None);
};
blobs
.write(diff.patch.as_bytes())
.await
.map(Some)
.map_err(|error| format!("the patch could not be stored: {error:#}"))
}
/// The checkpoints already recorded for the run, read once: what a
/// resume's reissued routing decisions must not record again.
/// resume's reissued routing decisions must not record again, and
/// where the run's diff is measured to when this process made no
/// checkpoint yet.
async fn load_recorded(&self) -> Result<(), String> {
let stored = self
.records
@ -578,6 +822,7 @@ impl FabroHooks {
)
})?;
let mut recorded = lock(&self.recorded);
let mut last = None;
for record in stored {
let PlatformRecord::Checkpoint(checkpoint) = &record.record else {
continue;
@ -589,7 +834,191 @@ impl FabroHooks {
{
recorded.insert(key);
}
if let (Some(workspace), Some(sha)) =
(&checkpoint.workspace, &checkpoint.git_commit_sha)
{
last = Some((workspace.clone(), sha.clone()));
}
}
drop(recorded);
let mut last_checkpoint = lock(&self.last_checkpoint);
if last_checkpoint.is_none() {
*last_checkpoint = last;
}
Ok(())
}
/// The artifacts of a finished attempt: every file of its workspace
/// under the run's patterns, stored once. `Ok` is how many files were
/// collected; `Err` names the first problem that stopped the
/// collection.
async fn collect_artifacts(
&self,
context: &HookContext,
scope: ScopeId,
key: CheckpointKey,
) -> Result<usize, String> {
let globs = match &self.artifact_globs {
Ok(globs) => globs,
Err(error) => return Err(format!("invalid run.artifacts.include pattern: {error}")),
};
if globs.is_empty() {
return Ok(0);
}
let Some((_, env)) = self.env_of(context, scope) else {
// A skipped node or a driver-made outcome may precede the scope's
// environment; there is no workspace to collect from.
return Ok(0);
};
let Some(blobs) = &self.blobs else {
return Err("the run has no blob table to collect artifacts into".to_string());
};
self.collected_loaded
.get_or_try_init(|| self.load_collected())
.await?;
let candidates = list_artifacts(env.as_ref(), globs).await?;
let limit = usize::try_from(ARTIFACT_MAX_FILE_BYTES).unwrap_or(usize::MAX);
let mut collected = 0;
let mut total_bytes = 0_u64;
for (path, size) in select_artifacts(candidates) {
if total_bytes.saturating_add(size) > ARTIFACT_MAX_TOTAL_BYTES {
break;
}
let bytes = match env.read_file_limited(Path::new(&path), limit).await {
Ok(Some(bytes)) => bytes,
Ok(None) => continue,
Err(error) => {
warn!(run_id = %self.run_id, path, error = %error, "an artifact could not be read");
continue;
}
};
let digest = BlobHash::new(&bytes);
let identity = (path.clone(), digest.to_string());
if lock(&self.collected).contains(&identity) {
continue;
}
let blob = blobs
.write(&bytes)
.await
.map_err(|error| format!("the artifact `{path}` could not be stored: {error:#}"))?;
let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord {
execution: key.execution,
firing: key.firing,
attempt: key.attempt,
path: path.clone(),
blob,
bytes: u64::try_from(bytes.len()).unwrap_or(u64::MAX),
digest: digest.to_string(),
operation: Some(key.operation_for(ARTIFACT_EFFECT)),
});
self.records
.append(
&self.run_id,
&record,
Some(StagePosition {
execution: key.execution,
firing: key.firing,
}),
)
.await
.map_err(|error| {
format!(
"the artifact record for `{path}` could not be written: {}",
collect_chain(&error).join(": ")
)
})?;
lock(&self.collected).insert(identity);
total_bytes = total_bytes.saturating_add(size);
collected += 1;
}
Ok(collected)
}
/// The artifacts already collected for the run, read once: a file that
/// is unchanged since it was collected is not collected again.
async fn load_collected(&self) -> Result<(), String> {
let stored = self
.records
.read_kind(&self.run_id, PlatformRecordKind::ArtifactCollected)
.await
.map_err(|error| {
format!(
"the run's artifact records could not be read: {}",
collect_chain(&error).join(": ")
)
})?;
let mut collected = lock(&self.collected);
for record in stored {
if let PlatformRecord::ArtifactCollected(artifact) = record.record {
collected.insert((artifact.path, artifact.digest));
}
}
Ok(())
}
/// The run's diff: the run branch's last checkpoint against the base
/// the branch started from, in the snapshot repository on this host.
/// Nothing is recorded for a run that never created its branch or
/// never checkpointed.
async fn record_run_diff(&self) -> Result<(), String> {
self.recorded_loaded
.get_or_try_init(|| self.load_recorded())
.await?;
let branch = match self.branch.get() {
Some(branch) => branch.clone(),
None => match self.stored_branch().await? {
Some(branch) => branch,
None => {
debug!(run_id = %self.run_id, "no run branch is recorded; no run diff");
return Ok(());
}
},
};
let Some(base_sha) = branch.base_sha.clone() else {
return Ok(());
};
let last = lock(&self.last_checkpoint).clone();
let Some((workspace, head_sha)) = last else {
debug!(run_id = %self.run_id, "no checkpoint is recorded; no run diff");
return Ok(());
};
// The run's diff is measured in the workspace the branch started
// in; a last checkpoint elsewhere (a nested invocation's workspace)
// is not this branch's head.
let workspace = branch.workspace.clone().unwrap_or(workspace);
let diff = self
.workspaces
.diff(&workspace, Some(&base_sha), &head_sha)
.await
.map_err(|error| {
format!(
"the run's diff could not be computed: {}",
collect_chain(&error).join(": ")
)
})?;
let patch_blob = self.patch_blob(&diff).await?;
let record = PlatformRecord::RunDiff(RunDiffRecord {
base_sha: Some(base_sha),
head_sha: Some(head_sha),
diff_summary: Some(diff.summary),
patch_blob,
});
self.records
.append(&self.run_id, &record, None)
.await
.map_err(|error| {
format!(
"the run diff record could not be written: {}",
collect_chain(&error).join(": ")
)
})?;
info!(
run_id = %self.run_id,
files_changed = diff.summary.files_changed,
additions = diff.summary.additions,
deletions = diff.summary.deletions,
"run diff recorded"
);
Ok(())
}
@ -611,6 +1040,70 @@ impl FabroHooks {
}
}
/// Every file under the patterns' traversal roots that matches a pattern,
/// with its size, listed through the scope's environment. Directories
/// never committed are never collected either.
async fn list_artifacts(
env: &dyn ExecEnv,
globs: &WorkspaceGlobSet,
) -> Result<Vec<(String, u64)>, String> {
let mut files = Vec::new();
for root in globs.traversal_roots() {
let listed = env
.list_directory(
Path::new(if root.is_empty() { "." } else { root }),
ARTIFACT_LIST_DEPTH,
)
.await
.map_err(|error| {
format!("the workspace could not be listed below `{root}`: {error}")
})?;
for entry in listed {
if entry.is_dir {
continue;
}
let path = entry.path.trim_start_matches("./").to_string();
let path = if root.is_empty() || path.starts_with(&format!("{root}/")) {
path
} else {
format!("{root}/{path}")
};
if path
.split('/')
.any(|segment| EXCLUDE_DIRS.contains(&segment))
{
continue;
}
if !globs.is_match(&path) {
continue;
}
files.push((path, entry.size.unwrap_or(0)));
}
}
files.sort();
files.dedup();
Ok(files)
}
/// The files within the collection's budget: the legacy executor's rule,
/// smallest first, each under the file limit, at most the count limit.
fn select_artifacts(mut candidates: Vec<(String, u64)>) -> Vec<(String, u64)> {
candidates.retain(|(_, size)| *size <= ARTIFACT_MAX_FILE_BYTES);
candidates.sort_by(|left, right| left.1.cmp(&right.1).then_with(|| left.0.cmp(&right.0)));
let mut total = 0_u64;
let mut selected = Vec::new();
for (path, size) in candidates {
if selected.len() >= ARTIFACT_MAX_FILES
|| total.saturating_add(size) > ARTIFACT_MAX_TOTAL_BYTES
{
break;
}
total = total.saturating_add(size);
selected.push((path, size));
}
selected
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
@ -708,6 +1201,30 @@ impl ExecutionHooks for FabroHooks {
);
problems.push(problem);
}
match self.collect_artifacts(context, scope, key).await {
Ok(0) => {}
Ok(collected) => {
debug!(
run_id = %self.run_id,
node,
execution = key.execution,
firing = key.firing,
collected,
"artifacts collected"
);
}
Err(problem) => {
warn!(
run_id = %self.run_id,
node,
execution = key.execution,
firing = key.firing,
error = %problem,
"artifact collection failed"
);
problems.push(format!("artifact collection failed: {problem}"));
}
}
let mut report = self.inner.transition(context, transition).await?;
report.problems.extend(problems);
Ok(report)
@ -718,8 +1235,11 @@ impl ExecutionHooks for FabroHooks {
run_id = %self.run_id,
status = ?finished.status,
failure = finished.failure.as_deref().unwrap_or(""),
"Petri run finished; running the run-end hooks"
"Petri run finished; recording the run's diff and running the run-end hooks"
);
if let Err(error) = self.record_run_diff().await {
warn!(run_id = %self.run_id, error = %error, "the run's diff was not recorded");
}
self.inner.run_finished(context, finished).await
}
@ -742,17 +1262,32 @@ impl ExecutionHooks for FabroHooks {
acquired: ScopeAcquired,
) -> Result<(), ScopeAcquiredError> {
self.inner.scope_acquired(context, acquired.clone()).await?;
if self.host_workspaces {
return Ok(());
}
let workspace = acquired.workspace.as_str().to_owned();
lock(&self.envs).insert(
(context.execution, acquired.scope),
(workspace.clone(), Arc::clone(&acquired.env)),
);
if !self.resumed {
if self.host_workspaces || !self.resumed {
return Ok(());
}
self.restore_sandbox(&workspace, &acquired.env).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_selection_keeps_the_smallest_files_within_the_budgets() {
let mut candidates: Vec<(String, u64)> = (0..(ARTIFACT_MAX_FILES + 5))
.map(|index| (format!("file{index:03}.txt"), 100))
.collect();
candidates.push(("huge.bin".to_string(), ARTIFACT_MAX_FILE_BYTES + 1));
candidates.push(("tiny.txt".to_string(), 1));
let selected = select_artifacts(candidates);
assert_eq!(selected.len(), ARTIFACT_MAX_FILES);
assert_eq!(selected[0], ("tiny.txt".to_string(), 1));
assert!(selected.iter().all(|(path, _)| path != "huge.bin"));
}
}

View file

@ -1,19 +1,61 @@
//! Petri's test kit, for Fabro crates that check a store implementation
//! against Petri's contract from their own tests, and an in-memory platform
//! record store for tests of the hooks and recovery. Compiled only with the
//! `test-support` feature, which a dev-dependency turns on.
//! against Petri's contract from their own tests, an in-memory platform
//! record store and an in-memory blob table for tests of the hooks and
//! recovery. Compiled only with the `test-support` feature, which a
//! dev-dependency turns on.
use std::collections::HashMap;
use std::sync::{Mutex, MutexGuard, PoisonError};
use async_trait::async_trait;
use bytes::Bytes;
use fabro_store::platform_records::now_ms;
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
use fabro_types::RunId;
use fabro_types::{BlobHash, RunId};
pub use petri_testkit::run_store;
use crate::blobs::Blobs;
use crate::platform_records::{PlatformRecordError, PlatformRecords};
/// A blob table in memory.
#[derive(Debug, Default)]
pub struct MemoryBlobs {
rows: Mutex<HashMap<BlobHash, Vec<u8>>>,
}
impl MemoryBlobs {
#[must_use]
pub fn new() -> Self {
Self::default()
}
/// How many blobs the table holds.
#[must_use]
pub fn len(&self) -> usize {
lock(&self.rows).len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
#[async_trait]
impl Blobs for MemoryBlobs {
async fn write(&self, bytes: &[u8]) -> anyhow::Result<BlobHash> {
let hash = BlobHash::new(bytes);
lock(&self.rows).insert(hash, bytes.to_vec());
Ok(hash)
}
async fn read(&self, hash: &BlobHash) -> anyhow::Result<Option<Bytes>> {
Ok(lock(&self.rows)
.get(hash)
.map(|bytes| Bytes::copy_from_slice(bytes)))
}
}
/// Platform records kept in memory, per run, in seq order.
#[derive(Debug, Default)]
pub struct MemoryPlatformRecords {

View file

@ -22,18 +22,19 @@ use std::sync::Arc;
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::controls::RunControls;
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::engine::{self, Execution, Retention, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::platform_records::PlatformRecords;
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::test_support::MemoryPlatformRecords;
use fabro_petri::test_support::{MemoryBlobs, MemoryPlatformRecords};
use fabro_store::{PlatformRecord, PlatformRecordKind};
use fabro_types::settings::run::RunCheckpointSettings;
use fabro_types::{RunId, SandboxProviderKind};
use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind};
use petri_execution::inspect::{self, RunInspection};
use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _};
use tokio::fs;
@ -139,22 +140,27 @@ fn docker_plugin() -> Option<PathBuf> {
/// One run's pieces: the store, its platform records, where it ran.
struct Harness {
run_id: RunId,
run_dir: PathBuf,
store: Arc<MemoryRunStore>,
records: Arc<MemoryPlatformRecords>,
_root: tempfile::TempDir,
run_id: RunId,
run_dir: PathBuf,
store: Arc<MemoryRunStore>,
records: Arc<MemoryPlatformRecords>,
blobs: Arc<MemoryBlobs>,
/// The `[run.artifacts] include` patterns the hooks collect under.
artifacts: Vec<String>,
_root: tempfile::TempDir,
}
impl Harness {
fn new() -> Self {
let root = tempfile::tempdir().expect("a temp dir");
Self {
run_id: RunId::new(),
run_dir: root.path().join("run"),
store: Arc::new(MemoryRunStore::new()),
records: Arc::new(MemoryPlatformRecords::new()),
_root: root,
run_id: RunId::new(),
run_dir: root.path().join("run"),
store: Arc::new(MemoryRunStore::new()),
records: Arc::new(MemoryPlatformRecords::new()),
blobs: Arc::new(MemoryBlobs::new()),
artifacts: Vec::new(),
_root: root,
}
}
@ -162,7 +168,9 @@ impl Harness {
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,
}
@ -191,12 +199,13 @@ impl Harness {
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
runtime: RuntimeSpec::default(),
provider,
retention: Retention::Always,
cancel: CancellationToken::new(),
controls: RunControls::new(),
interviewer,
observers,
secrets: None,
blobs: None,
blobs: Some(Arc::clone(&self.blobs) as Arc<dyn Blobs>),
hooks: Some(hooks),
};
engine::run(request).await.expect("the run executes")
@ -421,6 +430,180 @@ async fn every_finish_is_committed_and_recorded() {
assert_eq!(run_end, "run_complete\nsandbox_cleanup\n");
}
/// The files under `[run.artifacts] include` are collected once per
/// content into the blob table, the run branch and the author identity are
/// recorded when the branch is created, every checkpoint after the first
/// carries its diff from its parent, and the run's diff is recorded at the
/// end.
#[tokio::test]
async fn artifacts_the_branch_and_the_diffs_are_recorded() {
if host_plugin().is_none() {
return;
}
let mut harness = Harness::new();
harness.artifacts = vec!["assets/**".to_string()];
let workflow = workflow(
" write [shape=parallelogram, script=\"mkdir -p assets && printf one > \
assets/report.txt && echo line > story.txt\"]\n keep [shape=parallelogram, \
script=\"test -f assets/report.txt\"]\n change [shape=parallelogram, script=\"printf \
two > assets/report.txt\"]",
" start -> write -> keep -> change -> exit",
);
let outcome = harness.run(&workflow, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let records = harness.records.records(&harness.run_id);
let artifacts: Vec<_> = records
.iter()
.filter_map(|stored| match &stored.record {
PlatformRecord::ArtifactCollected(record) => Some(record.clone()),
_ => None,
})
.collect();
assert_eq!(artifacts.len(), 2, "one capture per content: {artifacts:?}");
assert!(
artifacts
.iter()
.all(|artifact| artifact.path == "assets/report.txt"),
"{artifacts:?}"
);
assert_eq!(artifacts[0].bytes, 3);
assert_eq!(artifacts[0].digest, artifacts[0].blob.to_string());
assert_ne!(artifacts[0].digest, artifacts[1].digest);
let bytes = harness
.blobs
.read(&artifacts[1].blob)
.await
.expect("the blob reads")
.expect("the blob exists");
assert_eq!(bytes.as_ref(), b"two");
// The first capture belongs to `write`, the second to `change`; `keep`
// saw the file unchanged and recorded nothing.
let checkpoint_firings: Vec<(String, u64)> = checkpoint_nodes(&harness).await;
let firing = |node: &str| {
checkpoint_firings
.iter()
.find(|(name, _)| name == node)
.map(|(_, firing)| *firing)
.expect("the node checkpointed")
};
assert_eq!(artifacts[0].firing, firing("write"));
assert_eq!(artifacts[1].firing, firing("change"));
let branches: Vec<_> = records
.iter()
.filter_map(|stored| match &stored.record {
PlatformRecord::RunBranch(record) => Some(record.clone()),
_ => None,
})
.collect();
assert_eq!(branches.len(), 1, "{branches:?}");
let workspace = harness.workspace().await;
assert_eq!(
branches[0].run_branch.as_deref(),
Some(format!("fabro/run/{}", harness.run_id).as_str())
);
assert_eq!(branches[0].workspace.as_deref(), Some(workspace.as_str()));
let checkpoints = harness.checkpoints();
assert_eq!(
branches[0].base_sha.as_deref(),
Some(checkpoints[0].1.as_str()),
"a branch in a fresh repository starts from its first checkpoint"
);
let identities: Vec<_> = records
.iter()
.filter_map(|stored| match &stored.record {
PlatformRecord::GitIdentity(record) => Some(record.identity.clone()),
_ => None,
})
.collect();
assert_eq!(identities.len(), 1, "{identities:?}");
assert_eq!(identities[0].source, GitIdentitySource::Default);
assert_eq!(identities[0].name, GitAuthor::default().name);
// The checkpoints carry their diffs: `start` is the root commit and has
// none; `write` adds two files; `keep` changes nothing; `change` edits
// one file.
let diffs: Vec<_> = records
.iter()
.filter_map(|stored| match &stored.record {
PlatformRecord::Checkpoint(record) => {
Some((record.diff_summary, record.patch_blob.is_some()))
}
_ => None,
})
.collect();
assert_eq!(diffs.len(), 5, "{diffs:?}");
assert_eq!(diffs[0], (None, false));
let write = diffs[1].0.expect("the write diff");
assert_eq!(
(write.files_changed, write.additions, write.deletions),
(2, 2, 0)
);
assert!(diffs[1].1, "the write patch is a blob");
let keep = diffs[2].0.expect("the keep diff");
assert_eq!(keep.files_changed, 0);
assert!(!diffs[2].1, "an empty diff has no patch blob");
let change = diffs[3].0.expect("the change diff");
assert_eq!(
(change.files_changed, change.additions, change.deletions),
(1, 1, 1)
);
let run_diffs: Vec<_> = records
.iter()
.filter_map(|stored| match &stored.record {
PlatformRecord::RunDiff(record) => Some(record.clone()),
_ => None,
})
.collect();
assert_eq!(run_diffs.len(), 1, "{run_diffs:?}");
let run_diff = &run_diffs[0];
assert_eq!(run_diff.base_sha, branches[0].base_sha);
assert_eq!(
run_diff.head_sha.as_deref(),
Some(checkpoints.last().expect("checkpoints").1.as_str())
);
let summary = run_diff.diff_summary.expect("the run diff summary");
assert_eq!((summary.files_changed, summary.additions), (2, 2));
let patch = harness
.blobs
.read(&run_diff.patch_blob.expect("the run patch is a blob"))
.await
.expect("the blob reads")
.expect("the blob exists");
let patch = String::from_utf8_lossy(&patch);
assert!(patch.contains("+two"), "{patch}");
assert!(patch.contains("+line"), "{patch}");
}
/// The node of every checkpoint record, in record order, with its firing.
async fn checkpoint_nodes(harness: &Harness) -> Vec<(String, u64)> {
let inspection = harness.inspection().await;
let history: Vec<(u64, String)> = inspection
.executions
.iter()
.filter_map(|execution| execution.engine.as_ref())
.flat_map(|engine| engine.history.iter())
.map(|record| (record.firing, record.node.to_string()))
.collect();
harness
.records
.records(&harness.run_id)
.into_iter()
.filter_map(|stored| match stored.record {
PlatformRecord::Checkpoint(record) => {
let node = history
.iter()
.find(|(firing, _)| *firing == record.firing)
.map(|(_, node)| node.clone())?;
Some((node, record.firing))
}
_ => None,
})
.collect()
}
/// A stage that fails on its own terms is committed like a successful one,
/// and its failure route runs on the committed files.
#[tokio::test]
@ -633,6 +816,7 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() {
..RuntimeSpec::default()
},
provider: SandboxProviderKind::LOCAL,
retention: Retention::Always,
cancel: CancellationToken::new(),
controls: RunControls::new(),
interviewer,

View file

@ -16,7 +16,7 @@ use std::time::{Duration, Instant};
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{Execution, RunRequest};
use fabro_petri::engine::{Execution, Retention, RunRequest};
use fabro_petri::interview::{Approval, FabroInterviewer, QuestionNotice, QuestionSink};
use fabro_petri::runtime::RuntimeSpec;
use fabro_types::SandboxProviderKind;
@ -115,6 +115,7 @@ pub(crate) fn run_request(
store,
runtime,
provider: SandboxProviderKind::LOCAL,
retention: Retention::Always,
cancel: CancellationToken::new(),
controls: RunControls::new(),
observers: vec![interviewer.observer()],