From 67ba595b01bff63794144f365487a9a361417684 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 16:00:11 -0400 Subject: [PATCH] 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 --- .../src/commands/run/petri_worker.rs | 1 + .../fabro-server/src/server/petri_runs.rs | 1 + lib/components/fabro-petri/README.md | 53 +- lib/components/fabro-petri/VIEWS.md | 5 +- lib/components/fabro-petri/src/checkpoint.rs | 278 ++++++- lib/components/fabro-petri/src/engine.rs | 78 +- lib/components/fabro-petri/src/hooks.rs | 687 ++++++++++++++++-- .../fabro-petri/src/test_support.rs | 50 +- lib/components/fabro-petri/tests/hooks.rs | 212 +++++- .../fabro-petri/tests/support/mod.rs | 3 +- 10 files changed, 1239 insertions(+), 129 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 23576c3ce..55a333437 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -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), diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 5a14d2c13..6b762914e 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -318,6 +318,7 @@ pub(crate) async fn execute(state: Arc, 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. diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index ccfd36793..99c1565d7 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -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/` 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 diff --git a/lib/components/fabro-petri/VIEWS.md b/lib/components/fabro-petri/VIEWS.md index 824269921..d660b91a9 100644 --- a/lib/components/fabro-petri/VIEWS.md +++ b/lib/components/fabro-petri/VIEWS.md @@ -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 | diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index 9e10922bb..9d8ad5cb4 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -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, } -/// 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, +} + +/// 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, +} + +/// 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 { - 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, 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 { + 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, 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, +/// `\t\t`, 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::().unwrap_or(0); + summary.deletions += deletions.parse::().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); diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 3f4d6065b..034ba838c 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -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 { 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 { 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 { 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 { Ok(outcome) } +/// When Petri keeps a run's workspaces after their scope is released, from +/// the run's environment settings: +/// +/// - `[environments..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, diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index a28e2fb5f..6b2244d72 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -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, 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, /// 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, 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); +/// 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, - run_id: RunId, - records: Arc, - workspaces: RunWorkspaces, - lookup: WorkspaceLookup, - host_workspaces: bool, - test_gates: Option, - handle: OnceLock, + inner: Arc, + run_id: RunId, + records: Arc, + /// Where an artifact's bytes and a diff's patch go; `None` records + /// summaries alone. + blobs: Option>, + workspaces: RunWorkspaces, + lookup: WorkspaceLookup, + identity: GitIdentity, + artifact_globs: Result, + host_workspaces: bool, + test_gates: Option, + handle: OnceLock, /// The workspace and commit of every checkpoint this process made. - committed: Mutex>, + committed: Mutex>, /// Which checkpoints have their platform record, loaded from the store /// once and kept up to date with every append. - recorded: Mutex>, - recorded_loaded: OnceCell<()>, + recorded: Mutex>, + recorded_loaded: OnceCell<()>, + /// The workspace and commit of the checkpoint recorded last: the head + /// the run's diff is measured to. + last_checkpoint: Mutex>, + /// The run branch as recorded, once: read from the store, or written + /// by the commit that created the branch. + branch: OnceCell, + /// Every artifact collected so far, by path and digest, loaded from the + /// store once and kept up to date with every append. + collected: Mutex>, + collected_loaded: OnceCell<()>, /// Inherited workspaces resolved through the run's records. - inherited: Mutex>>, + inherited: Mutex>>, /// 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>>>, + workspace_locks: Mutex>>>, /// The checkpoint failure that ended the run, when one did. - failure: Mutex>, - /// 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>, + failure: Mutex>, + /// 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>, /// 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>>, - store: Arc, + restore: OnceCell>>, + store: Arc, } 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, resumed: bool, + blobs: Option>, ) -> 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 { + 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, 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, 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, Option), 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, 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 { + 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, 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(mutex: &Mutex) -> 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")); + } +} diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs index 873fbd035..9779fd2ea 100644 --- a/lib/components/fabro-petri/src/test_support.rs +++ b/lib/components/fabro-petri/src/test_support.rs @@ -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>>, +} + +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 { + let hash = BlobHash::new(bytes); + lock(&self.rows).insert(hash, bytes.to_vec()); + Ok(hash) + } + + async fn read(&self, hash: &BlobHash) -> anyhow::Result> { + 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 { diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index e22e4f19f..9b07b9861 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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 { /// One run's pieces: the store, its platform records, where it ran. struct Harness { - run_id: RunId, - run_dir: PathBuf, - store: Arc, - records: Arc, - _root: tempfile::TempDir, + run_id: RunId, + run_dir: PathBuf, + store: Arc, + records: Arc, + blobs: Arc, + /// The `[run.artifacts] include` patterns the hooks collect under. + artifacts: Vec, + _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, 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, 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), 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, diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index 7cc8fd0c3..4db4604a0 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -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()],