From d60744b4d69f8ba59059c61c55026c70bace2a30 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 09:10:29 -0400 Subject: [PATCH] Checkpoint and recover Petri workspaces inside Docker and Daytona sandboxes A Petri run on Docker or Daytona keeps its workspace inside the scope's sandbox. Fabro's hooks now take the environment Petri hands them at `scope_acquired`, run `git` inside the scope through it (the path Petri's own sandbox-placed hooks take), and commit each stage on the run branch with the same message and trailers as the host path. The commit leaves the sandbox as a Git bundle, created against the newest ancestor the snapshot repository already holds, split into 8 MiB parts (the plugin transport reads one file up to 16 MiB), read out through the environment's file transfer, fetched into the bare snapshot repository on the host and named there under the checkpoint's ref. The repository holds every checkpoint whatever the provider, and the platform records name the same commits. The host path is unchanged; both sites share one runner and the same commands. Recovery is split: `recovery::plan` decides, over the records and the snapshot repository alone, what every live workspace must sit on and reconciles a lost record; `recover` applies it to host workspaces on the server, as before, and reports a sandbox workspace's target as deferred. The worker's hooks read the same plan at the scope's first acquisition after a resume and bring the sandbox workspace to it before any attempt runs there: verified or reset in a retained sandbox that still holds the commit, else restored from a bundle of the checkpoint written into the sandbox. Petri replaces a lease's lost sandbox on Fabro's request (`LostSandbox::Replace`), so a removed container comes back fresh and restored. The checkpoint records are written for every provider now. A Docker variant of the in-process hooks test moves a 20 MiB file through the split transfer; the same test runs on Daytona when live credentials are present. Three CLI scenarios run on a Docker environment: every stage's checkpoint published from the container, a retained container whose workspace drifted reset on restart, and a removed container replaced and restored from the snapshot. Co-Authored-By: Claude Fable 5.1 --- .../src/commands/run/petri_worker.rs | 5 +- lib/apps/fabro-cli/tests/it/scenario/mod.rs | 1 + lib/apps/fabro-cli/tests/it/scenario/petri.rs | 43 +- .../tests/it/scenario/petri_docker.rs | 391 +++++++++++ lib/components/fabro-petri/src/checkpoint.rs | 644 +++++++++++++++--- lib/components/fabro-petri/src/engine.rs | 16 +- lib/components/fabro-petri/src/hooks.rs | 291 ++++++-- lib/components/fabro-petri/src/lib.rs | 10 +- lib/components/fabro-petri/src/recovery.rs | 209 ++++-- lib/components/fabro-petri/tests/hooks.rs | 185 ++++- 10 files changed, 1569 insertions(+), 226 deletions(-) create mode 100644 lib/apps/fabro-cli/tests/it/scenario/petri_docker.rs 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 9c9bb0f70..998fe5d2c 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -28,8 +28,9 @@ //! loss as its error once the run has settled. //! //! Fabro's hooks ride the run with their platform records over the same -//! client: the checkpoint commit in the run's host workspace before every -//! durable finish, and its record after every route. +//! client: the checkpoint commit in the run's workspace, on the host or +//! inside its sandbox, before every durable finish, and its record after +//! every route. //! //! The runtime's settings layer is left empty here: the run's graphs were //! lowered and admitted at create time with the server's layer, and nothing diff --git a/lib/apps/fabro-cli/tests/it/scenario/mod.rs b/lib/apps/fabro-cli/tests/it/scenario/mod.rs index faff5ba69..0c7f732e0 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/mod.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/mod.rs @@ -9,6 +9,7 @@ mod auth; mod exec; mod lifecycle; mod petri; +mod petri_docker; mod petri_tools; mod server_lifecycle; mod smoke; diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index fb0a8e899..ecb841fd2 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -48,9 +48,9 @@ use crate::cmd::support::created_run_id; use crate::support::{TEST_DEV_TOKEN, TEST_SESSION_SECRET, seed_dev_token_auth}; const HOST_PLUGIN: &str = "sandbox-driver-host"; -const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS"; -const RUN_TIMEOUT: Duration = Duration::from_mins(1); -const POLL: Duration = Duration::from_millis(50); +pub(super) const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS"; +pub(super) const RUN_TIMEOUT: Duration = Duration::from_mins(1); +pub(super) const POLL: Duration = Duration::from_millis(50); /// The host plugin as Petri's lookup finds it: the override variable, else /// the executable on `PATH`. `None`, after saying so, when the test should @@ -146,7 +146,7 @@ impl RunningServer { /// Start the server process over this storage; the same call brings /// it back after a kill. - async fn launch(&mut self) { + pub(super) async fn launch(&mut self) { assert!(self.child.is_none(), "the server is already running"); let mut cmd = Command::new(env!("CARGO_BIN_EXE_fabro")); apply_test_isolation(&mut cmd, self.home_root.path()); @@ -198,7 +198,7 @@ impl RunningServer { /// Kill the server outright, as a crash would; its workers live on in /// their own process groups. - fn kill(&mut self) { + pub(super) fn kill(&mut self) { let mut child = self.child.take().expect("the server is running"); child.kill().expect("the server dies"); let _ = child.wait(); @@ -249,7 +249,7 @@ impl RunningServer { /// Where the run's worker ran Petri: the run's scratch under the /// server's storage. - fn petri_run_dir(&self, run_id: &str) -> PathBuf { + pub(super) fn petri_run_dir(&self, run_id: &str) -> PathBuf { let run_id: RunId = run_id.parse().expect("the run id parses"); Storage::new(&self.storage_dir) .run_scratch(&run_id) @@ -258,7 +258,7 @@ impl RunningServer { } /// The worker's own log for the run. - fn worker_log(&self, run_id: &str) -> PathBuf { + pub(super) fn worker_log(&self, run_id: &str) -> PathBuf { let run_id: RunId = run_id.parse().expect("the run id parses"); Storage::new(&self.storage_dir) .run_scratch(&run_id) @@ -269,18 +269,18 @@ impl RunningServer { /// Hold the worker's checkpoint at `point` (`commit` or `record`) for /// `node` until [`release`](Self::release). - fn hold(&self, point: &str, node: &str) { + pub(super) fn hold(&self, point: &str, node: &str) { std::fs::write(self.gates_dir.join(format!("{point}.{node}.hold")), "") .expect("the hold file writes"); } - fn release(&self, point: &str, node: &str) { + pub(super) fn release(&self, point: &str, node: &str) { std::fs::write(self.gates_dir.join(format!("{point}.{node}.release")), "") .expect("the release file writes"); } /// Wait until the worker's log says its checkpoint is held at a gate. - fn wait_until_held(&self, run_id: &str, point: &str, node: &str) { + pub(super) fn wait_until_held(&self, run_id: &str, point: &str, node: &str) { let log = self.worker_log(run_id); let needle = format!("checkpoint held at a test gate point=\"{point}\" node=\"{node}\""); let deadline = Instant::now() + RUN_TIMEOUT; @@ -332,7 +332,7 @@ impl RunningServer { } /// The run's checkpoint records, in seq order, as `(node position, sha)`. - async fn checkpoints(&self, run_id: &str) -> Vec<(CheckpointKey, String)> { + pub(super) async fn checkpoints(&self, run_id: &str) -> Vec<(CheckpointKey, String)> { let run_id: RunId = run_id.parse().expect("the run id parses"); self.platform_records() .await @@ -403,7 +403,7 @@ fn write_petri_workspace(context: &fabro_test::TestContext, script: &str) -> Pat /// A workspace holding the given workflow with a `workflow.toml` that names /// Petri. -fn write_petri_workflow(context: &fabro_test::TestContext, dot: &str) -> PathBuf { +pub(super) fn write_petri_workflow(context: &fabro_test::TestContext, dot: &str) -> PathBuf { let workspace = context.temp_dir.join("petri-workspace"); std::fs::create_dir_all(&workspace).expect("the workspace creates"); std::fs::write(workspace.join("workflow.fabro"), dot).expect("the workflow writes"); @@ -433,6 +433,17 @@ pub(super) fn run_detached_with( server: &RunningServer, workspace: &Path, extra: &[&str], +) -> String { + run_detached_in(context, server, workspace, "local", extra) +} + +/// `fabro run --detach` on the server's environment `environment`. +pub(super) fn run_detached_in( + context: &fabro_test::TestContext, + server: &RunningServer, + workspace: &Path, + environment: &str, + extra: &[&str], ) -> String { let target = server.target(); seed_dev_token_auth( @@ -445,7 +456,7 @@ pub(super) fn run_detached_with( .current_dir(workspace) .args(["--server", &target, "--detach"]) .args(extra) - .args(["--environment", "local", "workflow.toml"]) + .args(["--environment", environment, "workflow.toml"]) .output() .expect("the detached run executes"); assert!( @@ -599,7 +610,7 @@ fn worker_pid(run_id: &str) -> Option { .find_map(|line| line.trim().parse().ok()) } -fn wait_for_worker(run_id: &str) -> u32 { +pub(super) fn wait_for_worker(run_id: &str) -> u32 { let deadline = Instant::now() + RUN_TIMEOUT; loop { if let Some(pid) = worker_pid(run_id) { @@ -1352,7 +1363,7 @@ fn three_stage_bundle(context: &fabro_test::TestContext, gate: &Path) -> PathBuf /// under the host plugin, and a machine crash takes it with everything /// else, where a killed worker alone would leave it writing into the /// workspace. -fn crash(server: &mut RunningServer, worker: u32, gate: Option<&Path>) { +pub(super) fn crash(server: &mut RunningServer, worker: u32, gate: Option<&Path>) { server.kill(); fabro_proc::sigkill_process_group(worker); let deadline = Instant::now() + Duration::from_secs(10); @@ -1382,7 +1393,7 @@ fn crash(server: &mut RunningServer, worker: u32, gate: Option<&Path>) { /// Wait for the run to succeed after a restart, with the server's stderr /// on failure. -async fn wait_for_success(server: &RunningServer, run_id: &str) { +pub(super) async fn wait_for_success(server: &RunningServer, run_id: &str) { let status = wait_for_status(server, run_id, &["succeeded", "failed"]).await; assert_eq!( status, diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri_docker.rs b/lib/apps/fabro-cli/tests/it/scenario/petri_docker.rs new file mode 100644 index 000000000..9261d9ca5 --- /dev/null +++ b/lib/apps/fabro-cli/tests/it/scenario/petri_docker.rs @@ -0,0 +1,391 @@ +//! Petri runs on a Docker environment through a real server and its +//! worker: the workspace lives inside the run's container, every stage's +//! checkpoint is committed there and published to the snapshot repository +//! on the host, and a restart brings the container's workspace back to the +//! snapshot its durable state names, in the retained container or in a +//! fresh one when the old one is gone. +//! +//! The runs take their scope through the sandbox-driver Docker plugin on +//! this machine's daemon, so the tests skip, and say why, when the +//! executable is not found or no daemon answers, unless +//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set and the plugin is missing. The +//! server, the detached run and the crash come from `petri.rs`. + +#![expect( + clippy::disallowed_methods, + reason = "these scenarios locate the plugin through the process environment and drive the Docker daemon with its CLI" +)] +#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")] + +use std::env; +use std::path::{Path, PathBuf}; +use std::process::{Command, Stdio}; + +use fabro_petri::checkpoint::CheckpointKey; +use fabro_static::EnvVars; +use fabro_test::{expect_reqwest_json, test_context}; +use serde_json::json; + +use super::petri::{ + REQUIRE_ENV, RunningServer, crash, run_detached_in, wait_for_status, wait_for_success, + wait_for_worker, write_petri_workflow, +}; +use crate::support::TEST_DEV_TOKEN; + +const DOCKER_PLUGIN: &str = "sandbox-driver-docker"; +/// The server-side environment the runs select. +const ENVIRONMENT: &str = "docker"; + +/// The Docker plugin as Petri's lookup finds it, with a daemon that +/// answers. `None`, after saying so, when the test should skip; a panic +/// when the environment forbids a skip and the plugin is missing. +fn docker_plugin() -> Option { + let found = env::var_os(EnvVars::PETRI_SANDBOX_DOCKER_PLUGIN) + .map(PathBuf::from) + .or_else(|| { + env::split_paths(&env::var_os(EnvVars::PATH)?) + .map(|dir| dir.join(DOCKER_PLUGIN)) + .find(|candidate| candidate.is_file()) + }); + let Some(found) = found else { + assert!( + env::var_os(REQUIRE_ENV).is_none(), + "{REQUIRE_ENV} is set, but {DOCKER_PLUGIN} is not on PATH and {} is unset", + EnvVars::PETRI_SANDBOX_DOCKER_PLUGIN + ); + eprintln!( + "skipping: {DOCKER_PLUGIN} is not on PATH and {} is unset", + EnvVars::PETRI_SANDBOX_DOCKER_PLUGIN + ); + return None; + }; + let daemon = Command::new("docker") + .args(["version", "--format", "{{.Server.Version}}"]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .is_ok_and(|status| status.success()); + if !daemon { + eprintln!("skipping: no Docker daemon answers"); + return None; + } + Some(found) +} + +/// A server with a Docker environment beside the default local one. +async fn docker_server() -> RunningServer { + let server = RunningServer::start().await; + let body = json!({ + "id": ENVIRONMENT, + "provider": "docker", + "image": { "docker": null, "dockerfile": null }, + "resources": { "cpu": null, "memory": null, "disk": null }, + "network": { "mode": "allow_all", "allow": [] }, + "lifecycle": { "preserve": false, "stop_on_terminal": true, "auto_stop": null }, + "labels": {}, + "env": {} + }); + let response = fabro_test::test_http_client() + .post(format!("{}/api/v1/environments", server.api_base_url)) + .bearer_auth(TEST_DEV_TOKEN) + .json(&body) + .send() + .await + .expect("the environment create sends"); + expect_reqwest_json( + response, + fabro_http::StatusCode::CREATED, + "POST /api/v1/environments", + ) + .await; + server +} + +/// The run's container on the daemon, by Petri's run label: the one the +/// run's scope lives in. +fn container_of(run_id: &str) -> Option { + let output = Command::new("docker") + .args([ + "ps", + "-aq", + "--filter", + &format!("label=petri.run={run_id}"), + ]) + .output() + .expect("docker ps runs"); + let ids: Vec = String::from_utf8_lossy(&output.stdout) + .lines() + .map(str::trim) + .filter(|line| !line.is_empty()) + .map(str::to_owned) + .collect(); + assert!(ids.len() <= 1, "one container per run: {ids:?}"); + ids.into_iter().next() +} + +/// `sh -c script` inside the container's workspace. +fn docker_exec(container: &str, script: &str) -> String { + let output = Command::new("docker") + .args(["exec", "-w", "/workspace", container, "sh", "-c", script]) + .output() + .expect("docker exec runs"); + assert!( + output.status.success(), + "docker exec failed: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_owned() +} + +fn docker_rm(container: &str) { + let status = Command::new("docker") + .args(["rm", "-f", container]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .expect("docker rm runs"); + assert!(status.success(), "the container is removed"); +} + +/// Remove whatever the run left on the daemon, so a failed assertion does +/// not leak a container. +fn cleanup(run_id: &str) { + if let Some(container) = container_of(run_id) { + docker_rm(&container); + } +} + +/// The snapshot repository of the run's one workspace, on the host. +fn snapshot_repository(server: &RunningServer, run_id: &str) -> PathBuf { + let snapshots = server.petri_run_dir(run_id).join("snapshots"); + let mut repositories: Vec = std::fs::read_dir(&snapshots) + .expect("the snapshots directory lists") + .map(|entry| entry.expect("an entry reads").path()) + .filter(|path| path.extension().is_some_and(|extension| extension == "git")) + .collect(); + assert_eq!(repositories.len(), 1, "one workspace: {repositories:?}"); + repositories.remove(0) +} + +/// `git` in a repository on the host, its stdout. +fn git(repository: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .args(args) + .current_dir(repository) + .output() + .expect("git runs"); + assert!( + output.status.success(), + "git {args:?} failed: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_owned() +} + +/// The commits the snapshot repository holds, oldest first, as +/// `(sha, subject, key)`. +fn snapshot_commits(repository: &Path) -> Vec<(String, String, Option)> { + let log = git(repository, &[ + "log", + "--topo-order", + "--reverse", + "--all", + "--format=%H%x00%s%x00%B%x1e", + ]); + log.split('\u{1e}') + .filter(|entry| !entry.trim().is_empty()) + .map(|entry| { + let mut parts = entry.trim_start().splitn(3, '\0'); + let sha = parts.next().unwrap_or_default().to_string(); + let subject = parts.next().unwrap_or_default().to_string(); + let body = parts.next().unwrap_or_default(); + (sha, subject, CheckpointKey::from_message(body)) + }) + .collect() +} + +fn subjects(commits: &[(String, String, Option)]) -> Vec<&str> { + commits + .iter() + .map(|(_, subject, _)| subject.as_str()) + .collect() +} + +/// Two command stages: `one` writes a file; `two` checks it is the one +/// `one` wrote, that nothing else is in the workspace beside the +/// repository's own files, and writes another. +fn two_stage_bundle(context: &fabro_test::TestContext) -> PathBuf { + write_petri_workflow( + context, + "digraph Stages {\n graph [goal=\"Two stages\", default_max_retries=0]\n start \ + [shape=Mdiamond]\n exit [shape=Msquare]\n one [shape=parallelogram, script=\"echo one > \ + one.txt\"]\n two [shape=parallelogram, script=\"test \\\"$(cat one.txt)\\\" = one && \ + test ! -e stray.txt && echo two > two.txt\"]\n start -> one -> two -> exit\n}\n", + ) +} + +/// The commit subjects one run of the two-stage bundle produces. +fn two_stage_subjects(run_id: &str) -> Vec { + ["start", "one", "two", "exit"] + .iter() + .map(|node| format!("fabro({run_id}): {node} (success)")) + .collect() +} + +/// The checkpoint records and the published snapshots name the same +/// commits, and the two stages' trees hold their files. +fn assert_snapshots_complete(server: &RunningServer, run_id: &str, repository: &Path) { + let commits = snapshot_commits(repository); + assert_eq!(subjects(&commits), two_stage_subjects(run_id)); + let (one, _, _) = &commits[1]; + let (two, _, _) = &commits[2]; + assert_eq!(git(repository, &["show", &format!("{one}:one.txt")]), "one"); + assert_eq!(git(repository, &["show", &format!("{two}:two.txt")]), "two"); + let refs = git(repository, &[ + "for-each-ref", + "--format=%(objectname)", + "refs/checkpoints/", + ]); + let mut published: Vec<&str> = refs.lines().collect(); + published.sort_unstable(); + let checkpoints = futures_lite_block_on(server.checkpoints(run_id)); + let mut recorded: Vec<&str> = checkpoints.iter().map(|(_, sha)| sha.as_str()).collect(); + recorded.sort_unstable(); + assert_eq!( + published, recorded, + "every record names a published snapshot" + ); + assert_eq!(checkpoints.len(), 4, "{checkpoints:?}"); +} + +/// Wait on a future from a synchronous helper inside a multi-thread test. +fn futures_lite_block_on(future: impl std::future::Future) -> T { + tokio::task::block_in_place(|| tokio::runtime::Handle::current().block_on(future)) +} + +/// What the worker's log says it did to the sandbox workspace at resume. +fn restore_actions(server: &RunningServer, run_id: &str) -> Vec { + let log = std::fs::read_to_string(server.worker_log(run_id)).unwrap_or_default(); + log.lines() + .filter(|line| line.contains("sandbox workspace brought to its durable snapshot")) + .filter_map(|line| { + line.split_whitespace() + .find_map(|word| word.strip_prefix("action=").map(str::to_owned)) + }) + .collect() +} + +/// A run on Docker: every stage is committed inside the container, each +/// checkpoint is published to the snapshot repository on the host, and +/// nothing of the workspace is on the host. +#[tokio::test(flavor = "multi_thread")] +async fn a_docker_run_publishes_every_stages_checkpoint_from_the_container() { + if docker_plugin().is_none() { + return; + } + let context = test_context!(); + let server = docker_server().await; + let workspace = two_stage_bundle(&context); + let run_id = run_detached_in(&context, &server, &workspace, ENVIRONMENT, &[ + "--auto-approve", + ]); + wait_for_success(&server, &run_id).await; + + let repository = snapshot_repository(&server, &run_id); + assert_snapshots_complete(&server, &run_id, &repository); + assert!( + !server.petri_run_dir(&run_id).join("scopes").exists(), + "no workspace is on the host" + ); + assert!( + container_of(&run_id).is_some(), + "the container is retained after the run" + ); + cleanup(&run_id); + server.shutdown(); +} + +/// A worker killed after the first stage's durable finish, with the +/// container's workspace changed behind Petri's back: the restart resumes +/// on the retained container, the workspace is reset to the snapshot, and +/// the second stage sees the first stage's files and nothing else. +#[tokio::test(flavor = "multi_thread")] +async fn a_retained_container_whose_workspace_drifted_is_reset_on_restart() { + if docker_plugin().is_none() { + return; + } + let context = test_context!(); + let mut server = docker_server().await; + let workspace = two_stage_bundle(&context); + server.hold("record", "one"); + let run_id = run_detached_in(&context, &server, &workspace, ENVIRONMENT, &[ + "--auto-approve", + ]); + + wait_for_status(&server, &run_id, &["running"]).await; + let worker = wait_for_worker(&run_id); + server.wait_until_held(&run_id, "record", "one"); + let container = container_of(&run_id).expect("the run's container exists"); + docker_exec( + &container, + "echo junk > one.txt && echo stray > stray.txt && git status --porcelain", + ); + crash(&mut server, worker, None); + + server.release("record", "one"); + server.launch().await; + let resumed = wait_for_worker(&run_id); + assert_ne!(resumed, worker); + wait_for_success(&server, &run_id).await; + + assert_eq!( + container_of(&run_id).as_deref(), + Some(container.as_str()), + "the run continued in its retained container" + ); + assert_eq!(restore_actions(&server, &run_id), vec!["Reset".to_string()]); + let repository = snapshot_repository(&server, &run_id); + assert_snapshots_complete(&server, &run_id, &repository); + cleanup(&run_id); + server.shutdown(); +} + +/// A worker killed after the first stage's durable finish, with the +/// container removed while the run is down: the restart gets a fresh +/// container, the workspace is restored into it from the snapshot +/// repository, and the second stage sees the first stage's files. +#[tokio::test(flavor = "multi_thread")] +async fn a_lost_container_is_replaced_and_its_workspace_restored_from_the_snapshot() { + if docker_plugin().is_none() { + return; + } + let context = test_context!(); + let mut server = docker_server().await; + let workspace = two_stage_bundle(&context); + server.hold("record", "one"); + let run_id = run_detached_in(&context, &server, &workspace, ENVIRONMENT, &[ + "--auto-approve", + ]); + + wait_for_status(&server, &run_id, &["running"]).await; + let worker = wait_for_worker(&run_id); + server.wait_until_held(&run_id, "record", "one"); + let container = container_of(&run_id).expect("the run's container exists"); + crash(&mut server, worker, None); + docker_rm(&container); + assert_eq!(container_of(&run_id), None, "the container is gone"); + + server.release("record", "one"); + server.launch().await; + wait_for_success(&server, &run_id).await; + + let fresh = container_of(&run_id).expect("a fresh container was created"); + assert_ne!(fresh, container); + assert_eq!(restore_actions(&server, &run_id), vec![ + "Restored".to_string() + ]); + let repository = snapshot_repository(&server, &run_id); + assert_snapshots_complete(&server, &run_id, &repository); + cleanup(&run_id); + server.shutdown(); +} diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index b3cc30dcf..9e10922bb 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -14,26 +14,38 @@ //! at `scopes//work`, the layout `HostExecutor::workspace_for` //! names. This module reaches it there and runs `git` on the host, which //! is where the worker, and the server at recovery, run. A Docker or -//! Daytona workspace lives inside its sandbox, out of reach of this module: -//! the hooks record that no snapshot was taken and recovery resumes such a -//! run on the retained sandbox as it was left. +//! Daytona workspace lives inside its sandbox: there `git` runs inside the +//! scope through the environment Petri hands the hooks at +//! `scope_acquired`, the same capability a step spawns its process with, +//! and the same commands run on both sites through one runner +//! ([`Site`]). Only the transfer differs: a sandbox commit leaves its +//! sandbox as a Git bundle and a restore enters one the same way. //! //! # The snapshot repository //! -//! Every checkpoint commit is also pushed to a bare repository beside the -//! run's workspaces, `snapshots/.git`, under an immutable ref -//! per checkpoint (`refs/checkpoints///`). A -//! workspace that is gone at recovery is restored from it, and the refs -//! are what recovery reconciles a missing record from. +//! Every checkpoint commit is also published to a bare repository beside +//! the run's workspaces, `snapshots/.git`, under an immutable +//! ref per checkpoint (`refs/checkpoints///`). +//! A host workspace pushes to it; a sandbox workspace bundles the commit +//! (`git bundle create`, against the newest ancestor the repository already +//! holds), the bundle is read out of the sandbox through the environment's +//! file transfer in parts the transport accepts, and the repository fetches +//! it. A workspace that is gone at recovery is restored from the +//! repository: a host directory fetches from it, a sandbox receives a +//! bundle of the checkpoint and fetches from that. The refs are what +//! recovery reconciles a missing record from, whatever the provider. use std::path::{Path, PathBuf}; use std::process::Stdio; +use std::sync::Arc; 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::settings::run::RunCheckpointSettings; +use petri_runtime::executor::{EnvError, ExecEnv, OutputMode, ProcessSpec, Sig}; +use petri_runtime::ir::LogStream; use tokio::process::Command; use tokio::{fs, time}; @@ -52,6 +64,13 @@ 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/"; +/// 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"; +/// The largest piece of a bundle read out of a sandbox at once: half the +/// plugin transport's 16 MiB cap on one file read. +const TRANSFER_PART_BYTES: u64 = 8 * 1024 * 1024; + /// Directories never committed, the legacy executor's list: build output /// and dependency caches a stage regenerates. pub const EXCLUDE_DIRS: &[&str] = &[ @@ -120,6 +139,11 @@ impl CheckpointKey { } } + /// The key as a file name fragment. + fn transfer_name(self) -> String { + format!("{}-{}-{}", self.execution, self.firing, self.attempt) + } + fn from_ref(name: &str) -> Option { let mut parts = name.strip_prefix(REFS_PREFIX)?.split('/'); let execution = parts.next()?.parse().ok()?; @@ -173,6 +197,40 @@ pub enum CheckpointError { }, #[error("the restored workspace is at {actual}, not the snapshot {expected}")] RestoreMismatch { expected: String, actual: String }, + #[error("the snapshot bundle could not be {action} the sandbox")] + Transfer { + /// `read out of` or `written into`. + action: &'static str, + #[source] + source: EnvError, + }, +} + +/// Where a workspace's `git` runs: in a directory on this host, or inside +/// a scope's sandbox through the environment Petri handed the hooks. +#[derive(Clone)] +pub enum Site { + Host(PathBuf), + Sandbox(Arc), +} + +impl std::fmt::Debug for Site { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Host(path) => f.debug_tuple("Host").field(path).finish(), + Self::Sandbox(env) => f + .debug_tuple("Sandbox") + .field(&env.workspace_path()) + .finish(), + } + } +} + +/// What one `git` run produced, on either site. +struct GitOutput { + success: bool, + stdout: Vec, + stderr: Vec, } /// A checkpoint commit: the commit, and whether an earlier attempt of the @@ -246,6 +304,11 @@ impl RunWorkspaces { .unwrap_or(false) } + /// The host site of a workspace. + fn host(&self, workspace: &str) -> Site { + Site::Host(self.workspace_path(workspace)) + } + /// Commit the workspace's files on the run branch as the snapshot of /// `key`, and publish it. An earlier commit of the same key that the /// workspace still sits on, unchanged, is reused. @@ -256,17 +319,49 @@ impl RunWorkspaces { node: &str, status: &str, ) -> Result { - let path = self.workspace_path(workspace); if !self.workspace_exists(workspace).await { return Err(CheckpointError::WorkspaceMissing { workspace: workspace.to_string(), - path, + path: self.workspace_path(workspace), }); } - self.ensure_repository(&path).await?; + self.commit_at(&self.host(workspace), workspace, key, node, status) + .await + } + + /// [`commit`](Self::commit) for a workspace inside a sandbox: `git` + /// runs in the scope through `env`, and the commit reaches the + /// snapshot repository as a bundle. + pub async fn commit_in( + &self, + env: &Arc, + workspace: &str, + key: CheckpointKey, + node: &str, + status: &str, + ) -> Result { + self.commit_at( + &Site::Sandbox(Arc::clone(env)), + workspace, + key, + node, + status, + ) + .await + } + + async fn commit_at( + &self, + site: &Site, + workspace: &str, + key: CheckpointKey, + node: &str, + status: &str, + ) -> Result { + self.ensure_repository(site).await?; if let Some(existing) = self.published_sha(workspace, key).await? { - if self.head(&path).await?.as_deref() == Some(existing.as_str()) - && self.is_clean(&path).await? + if self.head(site).await?.as_deref() == Some(existing.as_str()) + && self.is_clean(site).await? { return Ok(Snapshot { sha: existing, @@ -290,11 +385,11 @@ impl RunWorkspaces { .iter() .map(|glob| format!(":(glob,exclude){glob}")), ); - self.git(&path, "add", &add).await?; + self.git(site, "add", &add).await?; let message = self.message(key, node, status); let user_name = format!("user.name={}", self.author.name); let user_email = format!("user.email={}", self.author.email); - self.git(&path, "commit", &[ + self.git(site, "commit", &[ "-c", &user_name, "-c", @@ -306,13 +401,16 @@ impl RunWorkspaces { &message, ]) .await?; - let sha = self.git(&path, "rev-parse", &["rev-parse", "HEAD"]).await?; - self.publish(workspace, &path, key, &sha).await?; + let sha = self.git(site, "rev-parse", &["rev-parse", "HEAD"]).await?; + match site { + 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 }) } /// The commit of `key`, from the snapshot repository first, else from - /// the workspace's own history by the trailers. + /// the host workspace's own history by the trailers. pub async fn find( &self, workspace: &str, @@ -321,12 +419,12 @@ impl RunWorkspaces { if let Some(sha) = self.published_sha(workspace, key).await? { return Ok(Some(sha)); } - let path = self.workspace_path(workspace); - if !self.workspace_exists(workspace).await || self.head(&path).await?.is_none() { + let site = self.host(workspace); + if !self.workspace_exists(workspace).await || self.head(&site).await?.is_none() { return Ok(None); } let listed = self - .git(&path, "log", &[ + .git(&site, "log", &[ "log", "--format=%H", "--extended-regexp", @@ -350,7 +448,7 @@ impl RunWorkspaces { return Ok(Vec::new()); } let listed = self - .git(&repository, "for-each-ref", &[ + .git(&Site::Host(repository), "for-each-ref", &[ "for-each-ref", "--format=%(refname) %(objectname)", REFS_PREFIX, @@ -376,7 +474,7 @@ impl RunWorkspaces { ancestor: &str, descendant: &str, ) -> Result { - let repository = self.snapshot_repository(workspace); + let repository = Site::Host(self.snapshot_repository(workspace)); Ok(self .git_status(&repository, "merge-base", &[ "merge-base", @@ -390,21 +488,74 @@ impl RunWorkspaces { /// The workspace's `HEAD`, or `None` when it has no commit. pub async fn workspace_head(&self, workspace: &str) -> Result, CheckpointError> { - let path = self.workspace_path(workspace); - self.head(&path).await + self.head(&self.host(workspace)).await + } + + /// [`workspace_head`](Self::workspace_head) for a workspace inside a + /// sandbox. + pub async fn workspace_head_in( + &self, + env: &Arc, + ) -> Result, CheckpointError> { + self.head(&Site::Sandbox(Arc::clone(env))).await } /// Whether the workspace sits on `sha` with nothing changed since. pub async fn matches(&self, workspace: &str, sha: &str) -> Result { - let path = self.workspace_path(workspace); - Ok(self.head(&path).await?.as_deref() == Some(sha) && self.is_clean(&path).await?) + self.matches_at(&self.host(workspace), sha).await + } + + /// [`matches`](Self::matches) for a workspace inside a sandbox. + pub async fn matches_in( + &self, + env: &Arc, + sha: &str, + ) -> Result { + self.matches_at(&Site::Sandbox(Arc::clone(env)), sha).await + } + + async fn matches_at(&self, site: &Site, sha: &str) -> Result { + Ok(self.head(site).await?.as_deref() == Some(sha) && self.is_clean(site).await?) + } + + /// Whether a sandbox workspace's repository holds the commit `sha`, so + /// a reset can reach it without a transfer. + pub async fn has_commit_in( + &self, + env: &Arc, + sha: &str, + ) -> Result { + let site = Site::Sandbox(Arc::clone(env)); + if self + .git_status(&site, "rev-parse", &["rev-parse", "--git-dir"]) + .await? + .is_none() + { + return Ok(false); + } + Ok(self + .git_status(&site, "cat-file", &[ + "cat-file", + "-e", + &format!("{sha}^{{commit}}"), + ]) + .await? + .is_some()) } /// Bring the workspace back to `sha`: tracked files reset, untracked /// files removed, the excluded caches left alone. pub async fn reset(&self, workspace: &str, sha: &str) -> Result<(), CheckpointError> { - let path = self.workspace_path(workspace); - self.git(&path, "reset", &["reset", "-q", "--hard", sha]) + self.reset_at(&self.host(workspace), sha).await + } + + /// [`reset`](Self::reset) for a workspace inside a sandbox. + pub async fn reset_in(&self, env: &Arc, sha: &str) -> Result<(), CheckpointError> { + self.reset_at(&Site::Sandbox(Arc::clone(env)), sha).await + } + + async fn reset_at(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> { + self.git(site, "reset", &["reset", "-q", "--hard", sha]) .await?; let mut clean = vec!["clean".to_string(), "-fdq".to_string()]; for dir in EXCLUDE_DIRS { @@ -415,7 +566,7 @@ impl RunWorkspaces { clean.push("-e".to_string()); clean.push(glob.clone()); } - self.git(&path, "clean", &clean).await?; + self.git(site, "clean", &clean).await?; Ok(()) } @@ -434,10 +585,11 @@ impl RunWorkspaces { path: path.clone(), source, })?; - self.git(&path, "init", &["init", "-q"]).await?; + let site = Site::Host(path); + self.git(&site, "init", &["init", "-q"]).await?; let repository = self.snapshot_repository(workspace); let repository = repository.to_string_lossy().into_owned(); - self.git(&path, "fetch", &[ + self.git(&site, "fetch", &[ "fetch", "-q", &repository, @@ -445,7 +597,7 @@ impl RunWorkspaces { ]) .await?; let branch = self.run_branch(); - self.git(&path, "checkout", &[ + self.git(&site, "checkout", &[ "checkout", "-q", "-B", @@ -453,7 +605,75 @@ impl RunWorkspaces { "FETCH_HEAD", ]) .await?; - let actual = self.git(&path, "rev-parse", &["rev-parse", "HEAD"]).await?; + self.verify_restored(&site, sha).await + } + + /// [`restore`](Self::restore) into a sandbox: the snapshot enters the + /// scope as a bundle of the checkpoint's ref, and the workspace, fresh + /// or stale, is fetched from it and forced onto the run branch at `sha`. + pub async fn restore_in( + &self, + env: &Arc, + workspace: &str, + key: CheckpointKey, + sha: &str, + ) -> Result<(), CheckpointError> { + let site = Site::Sandbox(Arc::clone(env)); + let repository = self.snapshot_repository(workspace); + let bundle = self.transfer_path(&format!("restore-{}.bundle", key.transfer_name())); + let staged = self.run_dir.join("snapshots").join(format!( + "{workspace}.restore-{}.bundle", + key.transfer_name() + )); + self.git(&Site::Host(repository), "bundle create", &[ + "bundle", + "create", + &staged.to_string_lossy(), + &key.snapshot_ref(), + ]) + .await?; + let bytes = fs::read(&staged) + .await + .map_err(|source| CheckpointError::Io { + path: staged.clone(), + source, + })?; + let _ = fs::remove_file(&staged).await; + env.write_file(Path::new(&bundle), &bytes) + .await + .map_err(|source| CheckpointError::Transfer { + action: "written into", + source, + })?; + let restored = async { + self.git(&site, "init", &["init", "-q"]).await?; + self.git(&site, "fetch", &[ + "fetch", + "-q", + &bundle, + &key.snapshot_ref(), + ]) + .await?; + let branch = self.run_branch(); + self.git(&site, "checkout", &[ + "checkout", + "-q", + "-f", + "-B", + &branch, + "FETCH_HEAD", + ]) + .await?; + self.reset_at(&site, "HEAD").await?; + self.verify_restored(&site, sha).await + } + .await; + self.remove_transfer(&site, &bundle).await; + restored + } + + async fn verify_restored(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> { + let actual = self.git(site, "rev-parse", &["rev-parse", "HEAD"]).await?; if actual != sha { return Err(CheckpointError::RestoreMismatch { expected: sha.to_string(), @@ -502,17 +722,17 @@ impl RunWorkspaces { /// A repository on the run branch, initialised when the workspace has /// none. - async fn ensure_repository(&self, path: &Path) -> Result<(), CheckpointError> { + async fn ensure_repository(&self, site: &Site) -> Result<(), CheckpointError> { if self - .git_status(path, "rev-parse", &["rev-parse", "--git-dir"]) + .git_status(site, "rev-parse", &["rev-parse", "--git-dir"]) .await? .is_none() { - self.git(path, "init", &["init", "-q"]).await?; + self.git(site, "init", &["init", "-q"]).await?; } let branch = self.run_branch(); let current = self - .git_status(path, "symbolic-ref", &[ + .git_status(site, "symbolic-ref", &[ "symbolic-ref", "-q", "--short", @@ -520,19 +740,17 @@ impl RunWorkspaces { ]) .await?; if current.as_deref() != Some(branch.as_str()) { - self.git(path, "checkout", &["checkout", "-q", "-B", &branch]) + self.git(site, "checkout", &["checkout", "-q", "-B", &branch]) .await?; } Ok(()) } - async fn publish( + /// The bare snapshot repository of the workspace, created on first use. + async fn ensure_snapshot_repository( &self, workspace: &str, - path: &Path, - key: CheckpointKey, - sha: &str, - ) -> Result<(), CheckpointError> { + ) -> Result { let repository = self.snapshot_repository(workspace); if !fs::try_exists(&repository).await.unwrap_or(false) { fs::create_dir_all(&repository) @@ -541,12 +759,27 @@ impl RunWorkspaces { path: repository.clone(), source, })?; - self.git(&repository, "init --bare", &["init", "-q", "--bare"]) - .await?; + self.git(&Site::Host(repository.clone()), "init --bare", &[ + "init", "-q", "--bare", + ]) + .await?; } + Ok(repository) + } + + /// Publish a host workspace's commit: a push into the snapshot + /// repository. + async fn publish( + &self, + workspace: &str, + path: &Path, + key: CheckpointKey, + sha: &str, + ) -> Result<(), CheckpointError> { + let repository = self.ensure_snapshot_repository(workspace).await?; let refspec = format!("{sha}:{}", key.snapshot_ref()); let repository = repository.to_string_lossy().into_owned(); - self.git(path, "push", &[ + self.git(&Site::Host(path.to_path_buf()), "push", &[ "push", "-q", "--force", @@ -557,6 +790,176 @@ impl RunWorkspaces { Ok(()) } + /// Publish a sandbox workspace's commit: a bundle of the run branch + /// since the newest ancestor the snapshot repository already holds, + /// read out of the sandbox in parts, fetched into the repository, and + /// named there under the checkpoint's ref. + async fn publish_from_sandbox( + &self, + env: &Arc, + workspace: &str, + key: CheckpointKey, + sha: &str, + ) -> Result<(), CheckpointError> { + let site = Site::Sandbox(Arc::clone(env)); + let repository = self.ensure_snapshot_repository(workspace).await?; + let branch = format!("refs/heads/{}", self.run_branch()); + // The bundle carries only what the repository lacks when the + // commit's parent is already there; the whole history otherwise. + let parent = self + .git_status(&site, "rev-parse", &[ + "rev-parse", + "-q", + "--verify", + "HEAD~1", + ]) + .await?; + let basis = match parent { + Some(parent) + if self + .git_status(&Site::Host(repository.clone()), "cat-file", &[ + "cat-file", + "-e", + &format!("{parent}^{{commit}}"), + ]) + .await? + .is_some() => + { + Some(parent) + } + _ => None, + }; + let revision = match &basis { + Some(parent) => format!("{parent}..{branch}"), + None => branch.clone(), + }; + let bundle = self.transfer_path(&format!("publish-{}.bundle", key.transfer_name())); + let published = async { + self.sh(&site, "prepare the transfer directory", &[ + "mkdir -p -- \"$(dirname -- \"$1\")\"", + "sh", + &bundle, + ]) + .await?; + self.git(&site, "bundle create", &[ + "bundle", "create", &bundle, &revision, + ]) + .await?; + let bytes = self.read_out(env, &site, &bundle).await?; + let staged = self.run_dir.join("snapshots").join(format!( + "{workspace}.publish-{}.bundle", + key.transfer_name() + )); + fs::write(&staged, &bytes) + .await + .map_err(|source| CheckpointError::Io { + path: staged.clone(), + source, + })?; + let fetched = self + .git(&Site::Host(repository.clone()), "fetch", &[ + "fetch", + "-q", + &staged.to_string_lossy(), + &branch, + ]) + .await; + let _ = fs::remove_file(&staged).await; + fetched?; + self.git(&Site::Host(repository.clone()), "update-ref", &[ + "update-ref", + &key.snapshot_ref(), + sha, + ]) + .await?; + Ok(()) + } + .await; + self.remove_transfer(&site, &bundle).await; + published + } + + /// Where a transfer file of this run waits inside a sandbox. + fn transfer_path(&self, name: &str) -> String { + format!("{TRANSFER_DIR}/{}/{name}", self.run_id) + } + + /// Read a file out of the sandbox in parts the transport accepts: the + /// file is split beside itself, each part comes through the + /// environment's file read, and the parts are removed as they go. + async fn read_out( + &self, + env: &Arc, + site: &Site, + path: &str, + ) -> Result, CheckpointError> { + let script = format!( + "split -b {TRANSFER_PART_BYTES} -a 4 -- \"$1\" \"$1.part.\" && rm -f -- \"$1\" && ls \ + -1 -- \"$1\".part.*" + ); + let listed = self + .sh(site, "split the bundle", &[&script, "sh", path]) + .await?; + let mut bytes = Vec::new(); + for part in listed + .lines() + .map(str::trim) + .filter(|line| !line.is_empty()) + { + let read = env.read_file(Path::new(part)).await.map_err(|source| { + CheckpointError::Transfer { + action: "read out of", + source, + } + })?; + let Some(read) = read else { + return Err(CheckpointError::Command { + action: "read the bundle".to_string(), + status: "missing".to_string(), + detail: format!("`{part}` is not in the sandbox"), + }); + }; + bytes.extend(read); + self.remove_transfer(site, part).await; + } + Ok(bytes) + } + + /// Remove a transfer file from the sandbox, best effort. + async fn remove_transfer(&self, site: &Site, path: &str) { + let _ = self + .sh(site, "remove the bundle", &["rm -f -- \"$1\"", "sh", path]) + .await; + } + + /// Run a shell command in the sandbox; a non-zero exit is the error. + async fn sh( + &self, + site: &Site, + action: &str, + args: &[&str], + ) -> Result { + let Site::Sandbox(env) = site else { + return Err(CheckpointError::Command { + action: action.to_string(), + status: "no sandbox".to_string(), + detail: "a shell transfer runs in a sandbox only".to_string(), + }); + }; + let mut all = vec!["-c"]; + all.extend(args); + let output = self.run_sandbox(env, "sh", &all, action).await?; + if output.success { + Ok(String::from_utf8_lossy(&output.stdout).trim().to_owned()) + } else { + Err(CheckpointError::Command { + action: action.to_string(), + status: "failed".to_string(), + detail: detail(&output.stderr), + }) + } + } + async fn published_sha( &self, workspace: &str, @@ -566,7 +969,7 @@ impl RunWorkspaces { if !fs::try_exists(&repository).await.unwrap_or(false) { return Ok(None); } - self.git_status(&repository, "rev-parse", &[ + self.git_status(&Site::Host(repository), "rev-parse", &[ "rev-parse", "-q", "--verify", @@ -575,71 +978,85 @@ impl RunWorkspaces { .await } - async fn head(&self, path: &Path) -> Result, CheckpointError> { - self.git_status(path, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"]) + async fn head(&self, site: &Site) -> Result, CheckpointError> { + self.git_status(site, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"]) .await } - async fn is_clean(&self, path: &Path) -> Result { - let status = self.git(path, "status", &["status", "--porcelain"]).await?; + async fn is_clean(&self, site: &Site) -> Result { + let status = self.git(site, "status", &["status", "--porcelain"]).await?; Ok(status.trim().is_empty()) } - /// Run `git` in `cwd`; a non-zero exit is the error. + /// Run `git` at `site`; a non-zero exit is the error. async fn git>( &self, - cwd: &Path, + site: &Site, action: &str, args: &[S], ) -> Result { - let output = self.run(cwd, action, args).await?; - if output.status.success() { + let output = self.run(site, action, args).await?; + if output.success { Ok(String::from_utf8_lossy(&output.stdout).trim().to_owned()) } else { Err(CheckpointError::Command { action: action.to_string(), - status: output.status.to_string(), + status: "non-zero exit".to_string(), detail: detail(&output.stderr), }) } } - /// Run `git` in `cwd`; a non-zero exit is `None`, for the queries whose + /// Run `git` at `site`; a non-zero exit is `None`, for the queries whose /// answer it is (an unborn `HEAD`, a missing ref, no repository). async fn git_status>( &self, - cwd: &Path, + site: &Site, action: &str, args: &[S], ) -> Result, CheckpointError> { - let output = self.run(cwd, action, args).await?; + let output = self.run(site, action, args).await?; Ok(output - .status - .success() + .success .then(|| String::from_utf8_lossy(&output.stdout).trim().to_owned())) } + /// Run `git` with the arguments and the configuration every checkpoint + /// command carries, at either site. async fn run>( &self, - cwd: &Path, + site: &Site, action: &str, args: &[S], - ) -> Result { + ) -> Result { + let mut all: Vec<&str> = vec![ + "-c", + "core.hooksPath=/dev/null", + "-c", + "commit.gpgsign=false", + "-c", + "gc.auto=0", + "-c", + "advice.detachedHead=false", + "-c", + "init.defaultBranch=main", + ]; + all.extend(args.iter().map(AsRef::as_ref)); + match site { + Site::Host(cwd) => self.run_host(cwd, &all, action).await, + Site::Sandbox(env) => self.run_sandbox(env, "git", &all, action).await, + } + } + + async fn run_host( + &self, + cwd: &Path, + args: &[&str], + action: &str, + ) -> Result { let mut command = Command::new("git"); command - .args([ - "-c", - "core.hooksPath=/dev/null", - "-c", - "commit.gpgsign=false", - "-c", - "gc.auto=0", - "-c", - "advice.detachedHead=false", - "-c", - "init.defaultBranch=main", - ]) - .args(args.iter().map(AsRef::as_ref)) + .args(args) .current_dir(cwd) .env("GIT_TERMINAL_PROMPT", "0") .env_remove("GIT_DIR") @@ -650,7 +1067,11 @@ impl RunWorkspaces { .stderr(Stdio::piped()) .kill_on_drop(true); match time::timeout(self.timeout, command.output()).await { - Ok(Ok(output)) => Ok(output), + Ok(Ok(output)) => Ok(GitOutput { + success: output.status.success(), + stdout: output.stdout, + stderr: output.stderr, + }), Ok(Err(source)) => Err(CheckpointError::Spawn { action: action.to_string(), source, @@ -661,6 +1082,73 @@ impl RunWorkspaces { }), } } + + /// Run `program` inside the sandbox, in its workspace, with the + /// checkpoint's deadline on the process and both streams captured. + async fn run_sandbox( + &self, + env: &Arc, + program: &str, + args: &[&str], + action: &str, + ) -> Result { + let spec = ProcessSpec::new(program, args) + .with_output(OutputMode::Bytes) + .with_timeout(Some(self.timeout)) + .with_env( + [("GIT_TERMINAL_PROMPT".into(), "0".into())] + .into_iter() + .collect(), + ); + let mut handle = env + .spawn(spec) + .await + .map_err(|error| CheckpointError::Command { + action: action.to_string(), + status: "spawn failed".to_string(), + detail: error.to_string(), + })?; + let Some(mut chunks) = handle.bytes() else { + let _ = handle.signal(Sig::Kill).await; + let _ = handle.wait().await; + return Err(CheckpointError::Command { + action: action.to_string(), + status: "no output stream".to_string(), + detail: "the sandbox offered no byte stream for the command".to_string(), + }); + }; + let drain = tokio::spawn(async move { + let mut stdout = Vec::new(); + let mut stderr = Vec::new(); + while let Some(chunk) = chunks.recv().await { + match chunk.stream { + LogStream::Stdout => stdout.extend(chunk.bytes), + LogStream::Stderr => stderr.extend(chunk.bytes), + } + } + (stdout, stderr) + }); + let status = handle + .wait() + .await + .map_err(|error| CheckpointError::Command { + action: action.to_string(), + status: "wait failed".to_string(), + detail: error.to_string(), + })?; + let (stdout, stderr) = drain.await.unwrap_or_default(); + if status.timed_out { + return Err(CheckpointError::TimedOut { + action: action.to_string(), + timeout: self.timeout, + }); + } + Ok(GitOutput { + success: status.is_success(), + stdout, + stderr, + }) + } } /// The tail of git's stderr for an error message: what the run's record diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index cfb6e06f8..badfab013 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -29,9 +29,11 @@ //! invocation is cancelled politely and Petri records why. //! //! A resume here is Petri's own: the run continues from its records, and -//! sandbox leases are reconciled by label. What the workspaces look like +//! sandbox leases are reconciled by label. What a host workspace looks like //! when it does is the server's business before it relaunches the worker -//! ([`recovery`](crate::recovery)). +//! ([`recovery`](crate::recovery)); a Docker or Daytona workspace is brought +//! to its snapshot by Fabro's hooks when its scope is acquired, and a lease +//! whose sandbox is gone gets a fresh one to restore into. //! //! No stage or agent event is projected into Fabro's tables here; the //! caller appends only the run lifecycle events Fabro's read side needs to @@ -50,7 +52,7 @@ use petri_execution::{ }; use petri_runtime::driver::lifecycle::ExecutionHooks; use petri_runtime::executor::{Retention, SecretProvider}; -use petri_runtime::{RunOptions, SandboxBackend}; +use petri_runtime::{LostSandbox, RunOptions, SandboxBackend}; use tokio::fs; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; @@ -165,6 +167,13 @@ pub async fn run(request: RunRequest) -> Result { options.run_key = Some(key.clone()); options.retention = Retention::Always; 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 + // one instead of failing the run. + if request.hooks.is_some() && backend != SandboxBackend::Host { + options.sandbox.lost_sandbox = LostSandbox::Replace; + } + let resumed = matches!(request.execution, Execution::Resume); let mut runtime = request .runtime .runtime(true) @@ -188,6 +197,7 @@ pub async fn run(request: RunRequest) -> Result { key.clone(), request.run_dir.clone(), Arc::clone(&request.store), + resumed, )) }); if let Some(hooks) = &fabro_hooks { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 34c16603e..a28e2fb5f 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -38,14 +38,21 @@ //! //! # Where the workspace is //! -//! The commit runs on the host, in the workspace Petri's host backend keeps -//! under the run directory (`crate::checkpoint`). A run on Docker or -//! Daytona has no workspace this process can reach; its hooks record that -//! no snapshot was taken and leave the run to continue as before. +//! On the local provider the commit runs on the host, in the workspace +//! Petri's host backend keeps under the run directory (`crate::checkpoint`). +//! 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 +//! 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. -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeMap, HashMap, HashSet}; use std::path::PathBuf; -use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError}; use std::time::Duration; @@ -58,10 +65,11 @@ use fabro_util::error::collect_chain; use petri_execution::{CancelReason, CoordinatorHandle, InvocationId, RunKey, RunStore}; use petri_runtime::driver::lifecycle::{ AdmitAttempt, AttemptDecision, ExecutionHooks, HookContext, Note, PrepareError, PrepareResult, - Prepared, Recorded, ResultOrigin, RunFinished, ScopeReleased, Transition, TransitionError, - TransitionReport, + Prepared, Recorded, ResultOrigin, RunFinished, ScopeAcquired, ScopeAcquiredError, + ScopeReleased, Transition, TransitionError, TransitionReport, }; -use petri_runtime::ir::{FailureInfo, ScopeId, Status}; +use petri_runtime::executor::ExecEnv; +use petri_runtime::ir::{ExecutionId, FailureInfo, ScopeId, Status}; use serde_json::json; use tokio::sync::{Mutex as AsyncMutex, OnceCell}; use tokio::{fs, time}; @@ -69,6 +77,7 @@ use tracing::{debug, info, warn}; use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces}; use crate::platform_records::PlatformRecords; +use crate::recovery::{self, Plan, RestoreTarget}; use crate::workspace::{self, WorkspaceLookup}; /// The note kind the hooks record on a firing about its checkpoint. @@ -84,7 +93,7 @@ pub struct HooksSpec { pub author: GitAuthor, pub checkpoint: RunCheckpointSettings, /// Whether the run's workspaces are on this host (the local sandbox - /// provider). A run elsewhere takes no snapshot. + /// provider). A run elsewhere snapshots inside its sandboxes. pub host_workspaces: bool, /// A test's gate directory: a checkpoint point named by a `.hold` file /// there waits for its `.release` file. `None` outside tests. @@ -118,37 +127,54 @@ impl HooksSpec { } } +/// A scope's sandbox environment as the hooks keep it: the workspace id +/// the executor named, and the environment `git` runs in. +type AcquiredEnv = (String, Arc); + /// 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, + workspaces: RunWorkspaces, + lookup: WorkspaceLookup, + 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<()>, /// 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>, - unreachable_noted: AtomicBool, + 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>, + /// Whether the run continues from its records: a sandbox workspace is + /// then brought to its snapshot when its scope is first acquired. + 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, } impl FabroHooks { /// Wrap `inner` (the hooks `Runtime::installed_hooks` returned) for the /// run whose records are in `store` under `run_key`, with its - /// workspaces under `run_dir`. + /// 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. #[must_use] pub fn new( spec: HooksSpec, @@ -157,6 +183,7 @@ impl FabroHooks { run_key: RunKey, run_dir: PathBuf, store: Arc, + resumed: bool, ) -> Self { let workspaces = RunWorkspaces::new(run_dir, run_id.to_string(), spec.author, &spec.checkpoint); @@ -165,7 +192,7 @@ impl FabroHooks { run_id, records: spec.records, workspaces, - lookup: WorkspaceLookup::new(store, run_key), + lookup: WorkspaceLookup::new(Arc::clone(&store), run_key), host_workspaces: spec.host_workspaces, test_gates: spec.test_gates, handle: OnceLock::new(), @@ -175,7 +202,10 @@ impl FabroHooks { inherited: Mutex::default(), workspace_locks: Mutex::default(), failure: Mutex::default(), - unreachable_noted: AtomicBool::new(false), + envs: Mutex::default(), + resumed, + restore: OnceCell::new(), + store, } } @@ -268,21 +298,9 @@ impl FabroHooks { origin: ResultOrigin, ) -> Result, String> { if !self.host_workspaces { - if !self.unreachable_noted.swap(true, Ordering::SeqCst) { - warn!( - run_id = %self.run_id, - "the run's workspaces are not on this host; no checkpoint snapshot is taken" - ); - } - return Ok(Some(Note::new( - CHECKPOINT_NOTE, - json!({ - "execution": key.execution, - "firing": key.firing, - "attempt": key.attempt, - "skipped": "the workspace is not on this host", - }), - ))); + return self + .snapshot_in_sandbox(context, scope, key, node, status, origin) + .await; } let workspace = self.workspace_of(context, scope).await?; if !self.workspaces.workspace_exists(&workspace).await { @@ -343,6 +361,136 @@ impl FabroHooks { } } + /// [`snapshot`](Self::snapshot) for a workspace inside the scope's + /// sandbox, through the environment kept at `scope_acquired`. + async fn snapshot_in_sandbox( + &self, + context: &HookContext, + scope: ScopeId, + key: CheckpointKey, + node: &str, + status: &Status, + origin: ResultOrigin, + ) -> Result, String> { + let held = lock(&self.envs).get(&(context.execution, scope)).cloned(); + let Some((workspace, env)) = held else { + // A skipped node or a driver-made outcome may precede the scope's + // environment; nothing of the stage's exists to snapshot. + if origin == ResultOrigin::Driver || matches!(status, Status::Skipped) { + return Ok(Some(Note::new( + CHECKPOINT_NOTE, + json!({ + "execution": key.execution, + "firing": key.firing, + "attempt": key.attempt, + "skipped": "the scope has no environment yet", + }), + ))); + } + return Err(format!( + "scope {scope} of execution {} has no sandbox environment to snapshot in", + context.execution + )); + }; + self.gate("commit", node).await; + let serialized = self.workspace_lock(&workspace); + let _held = serialized.lock().await; + match self + .workspaces + .commit_in(&env, &workspace, key, node, status.tag()) + .await + { + Ok(snapshot) => { + debug!( + run_id = %self.run_id, + node, + execution = key.execution, + firing = key.firing, + attempt = key.attempt, + reused = snapshot.reused, + "checkpoint committed in the sandbox" + ); + lock(&self.committed).insert(key, (workspace.clone(), snapshot.sha.clone())); + Ok(Some(Note::new( + CHECKPOINT_NOTE, + json!({ + "execution": key.execution, + "firing": key.firing, + "attempt": key.attempt, + "workspace": workspace, + "git_commit_sha": snapshot.sha, + "reused": snapshot.reused, + }), + ))) + } + Err(error) => Err(format!( + "the checkpoint commit of `{node}` in the sandbox failed: {}", + collect_chain(&error).join(": ") + )), + } + } + + /// 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( + &self, + ) -> Result<&Mutex>, ScopeAcquiredError> { + self.restore + .get_or_try_init(|| async { + let plan = recovery::plan( + Arc::clone(&self.store), + self.records.as_ref(), + &self.run_id, + &self.workspaces, + ) + .await + .map_err(|error| { + ScopeAcquiredError::new(format!( + "the run's restore plan could not be read: {}", + collect_chain(&error).join(": ") + )) + })?; + match plan { + Plan::Resume { targets } => Ok(Mutex::new(targets)), + Plan::Start => Ok(Mutex::default()), + Plan::Failed { reason } => Err(ScopeAcquiredError::new(reason)), + } + }) + .await + } + + /// Bring a sandbox workspace to the snapshot the resumed run's durable + /// state names, once, at its first acquisition. + async fn restore_sandbox( + &self, + workspace: &str, + env: &Arc, + ) -> Result<(), ScopeAcquiredError> { + let targets = self.restore_targets().await?; + let target = lock(targets).remove(workspace); + let Some(target) = target else { + return Ok(()); + }; + let serialized = self.workspace_lock(workspace); + let _held = serialized.lock().await; + let action = recovery::bring_sandbox_to(&self.workspaces, env, workspace, &target) + .await + .map_err(|error| { + ScopeAcquiredError::new(format!( + "the sandbox workspace `{workspace}` could not be brought to its snapshot: {}", + collect_chain(&error).join(": ") + )) + })?; + info!( + run_id = %self.run_id, + workspace, + sha = target.sha, + action = ?action, + "sandbox workspace brought to its durable snapshot" + ); + Ok(()) + } + /// The checkpoint's platform record, once per operation identity. async fn record( &self, @@ -360,7 +508,13 @@ impl FabroHooks { let (workspace, sha) = if let Some(committed) = committed { committed } else { - let workspace = self.workspace_of(context, scope).await?; + let acquired = lock(&self.envs) + .get(&(context.execution, scope)) + .map(|(workspace, _)| workspace.clone()); + let workspace = match acquired { + Some(workspace) => workspace, + None => self.workspace_of(context, scope).await?, + }; let serialized = self.workspace_lock(&workspace); let held = serialized.lock().await; let found = self.workspaces.find(&workspace, key).await; @@ -542,19 +696,17 @@ impl ExecutionHooks for FabroHooks { attempt: transition.view.attempt.raw(), }; let mut problems = Vec::new(); - if self.host_workspaces { - self.gate("record", &node).await; - if let Err(problem) = self.record(context, scope, key).await { - warn!( - run_id = %self.run_id, - node, - execution = key.execution, - firing = key.firing, - error = %problem, - "the checkpoint record was not written" - ); - problems.push(problem); - } + self.gate("record", &node).await; + if let Err(problem) = self.record(context, scope, key).await { + warn!( + run_id = %self.run_id, + node, + execution = key.execution, + firing = key.firing, + error = %problem, + "the checkpoint record was not written" + ); + problems.push(problem); } let mut report = self.inner.transition(context, transition).await?; report.problems.extend(problems); @@ -578,6 +730,29 @@ impl ExecutionHooks for FabroHooks { outcome = ?released.outcome, "scope released; running the sandbox cleanup hooks" ); - self.inner.scope_released(context, released).await + let scope = released.scope; + let notes = self.inner.scope_released(context, released).await; + lock(&self.envs).remove(&(context.execution, scope)); + notes + } + + async fn scope_acquired( + &self, + context: &HookContext, + 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 { + return Ok(()); + } + self.restore_sandbox(&workspace, &acquired.env).await } } diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index fe3a64df2..ae740aea5 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -36,11 +36,13 @@ //! - [`hooks`]: Fabro's `ExecutionHooks`, the checkpoint commit in //! `prepare_result` and its platform record in `transition`, around Petri's //! own hook service for `[[run.hooks]]`; -//! - [`checkpoint`]: the Git snapshots of a run's host workspaces and the -//! snapshot repository they are published to; +//! - [`checkpoint`]: the Git snapshots of a run's workspaces, on the host or +//! inside a Docker or Daytona sandbox, and the snapshot repository they are +//! published to; //! - [`recovery`]: the resume-on-restart protocol, which brings every live -//! workspace to the snapshot its durable state names before the run goes back -//! to a worker; +//! workspace to the snapshot its durable state names: a host workspace before +//! the run goes back to a worker, a sandbox workspace in the worker when its +//! scope is acquired; //! - [`platform_records`]: Fabro's platform records as the adapters reach them, //! in the server's database or over its API from a worker; //! - [`host_tools`]: Fabro's run tools on every native agent session of a run, diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index c5da156df..eed0f9ac9 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -24,9 +24,14 @@ //! the caller's workspace, and the workspace is brought to the newest of //! the live executions' snapshots on it. //! -//! A run whose workspaces are not on this host (Docker, Daytona) is resumed -//! on its retained sandbox as it was left: the snapshot side of the -//! protocol reaches only host workspaces. +//! The decision is [`plan`], over the records and the snapshot repository +//! alone, both on this host whatever the provider. Applying it differs: a +//! host workspace is brought to its snapshot here, before the worker is +//! relaunched; a Docker or Daytona workspace lives inside a sandbox only +//! the worker's run reaches, so its target is deferred, and the worker's +//! hooks read the same plan and apply it through the scope's environment +//! at `scope_acquired`, before the first attempt runs there +//! ([`bring_sandbox_to`]). use std::collections::BTreeMap; use std::path::PathBuf; @@ -40,8 +45,9 @@ use fabro_types::{RunId, SandboxProviderKind}; use petri_execution::host::{self, HostError}; use petri_execution::inspect::{self, ExecutionInspection, InspectError}; use petri_execution::{Access, InvocationId, RunKey, RunStore}; +use petri_runtime::executor::ExecEnv; use petri_store::StoreError; -use tracing::{info, warn}; +use tracing::info; use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces}; use crate::platform_records::{PlatformRecordError, PlatformRecords}; @@ -98,6 +104,29 @@ pub enum WorkspaceAction { Reset, /// It was gone and was recreated from the snapshot repository. Restored, + /// It lives in a sandbox this process does not reach: the worker's + /// hooks bring it to the snapshot when its scope is acquired. + Deferred, +} + +/// The snapshot a workspace must sit on before work resumes in it. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RestoreTarget { + pub key: CheckpointKey, + pub sha: String, +} + +/// What recovery decided, before any workspace was touched. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Plan { + /// The store never held the run: it starts from its admitted graphs. + Start, + /// The run continues; each live workspace, by id, and its snapshot. + Resume { + targets: BTreeMap, + }, + /// The run cannot continue and is reported failed. + Failed { reason: String }, } #[derive(Clone, Debug, PartialEq, Eq)] @@ -149,32 +178,38 @@ struct Target { key: CheckpointKey, } -/// Decide how the run continues, and bring its workspaces to their -/// snapshots. -pub async fn recover(request: RecoveryRequest) -> Result { - let key = RunKey::new(request.run_id.to_string()); - let logs = match request.store.open(&key, Access::Read).await { +/// Decide how the run continues: the snapshot every live workspace must sit +/// on, from the records and the snapshot repository, with a lost record +/// reconciled from the repository. Nothing is touched. +pub async fn plan( + store: Arc, + records: &dyn PlatformRecords, + run_id: &RunId, + workspaces: &RunWorkspaces, +) -> Result { + let key = RunKey::new(run_id.to_string()); + let logs = match store.open(&key, Access::Read).await { Ok(logs) => logs, - Err(StoreError::NotFound { .. }) => return Ok(Recovery::Start), + Err(StoreError::NotFound { .. }) => return Ok(Plan::Start), Err(error) => return Err(RecoveryError::Open(error)), }; // A record with no root invocation (the worker died between creating // the run and declaring it) has nothing to reconcile; the worker's // resume reports it as such. - let records = petri_execution::read_coordinator_log(&*logs) + let coordinator = petri_execution::read_coordinator_log(&*logs) .await .map_err(RecoveryError::Log)?; - if records.is_empty() { - return Ok(Recovery::Resume { - workspaces: Vec::new(), + if coordinator.is_empty() { + return Ok(Plan::Resume { + targets: BTreeMap::new(), }); } let state = host::stored_state(&*logs) .await .map_err(RecoveryError::State)?; if !state.invocations.contains_key(&InvocationId::ROOT) { - return Ok(Recovery::Resume { - workspaces: Vec::new(), + return Ok(Plan::Resume { + targets: BTreeMap::new(), }); } let inspection = inspect::inspect_run(&*logs) @@ -183,26 +218,11 @@ pub async fn recover(request: RecoveryRequest) -> Result Result Result Result Result { + let workspaces = RunWorkspaces::new( + request.run_dir.clone(), + request.run_id.to_string(), + request.author.clone(), + &request.checkpoint, + ); + let targets = match plan( + Arc::clone(&request.store), + &*request.records, + &request.run_id, + &workspaces, + ) + .await? + { + Plan::Start => return Ok(Recovery::Start), + Plan::Failed { reason } => return Ok(Recovery::Failed { reason }), + Plan::Resume { targets } => targets, + }; + let mut recovered = Vec::new(); - for (workspace, targets) in candidates { - let sha = newest(&workspaces, &workspace, &targets).await?; - let action = bring_to(&workspaces, &workspace, &sha, &targets).await?; + for (workspace, target) in targets { + let action = if request.host_workspaces { + bring_to(&workspaces, &workspace, &target).await? + } else { + WorkspaceAction::Deferred + }; info!( run_id = %request.run_id, workspace, - sha, + sha = target.sha, action = ?action, - "workspace brought to its durable snapshot" + "workspace's durable snapshot decided" ); recovered.push(RecoveredWorkspace { workspace, - sha, + sha: target.sha, action, }); } @@ -420,30 +470,71 @@ async fn newest( Ok(chosen.clone()) } -/// Verify, reset or restore the workspace onto `sha`. +/// Verify, reset or restore the host workspace onto its target. async fn bring_to( workspaces: &RunWorkspaces, workspace: &str, - sha: &str, - targets: &[(Target, String)], + target: &RestoreTarget, ) -> Result { let failed = |source| RecoveryError::Workspace { workspace: workspace.to_string(), source, }; if workspaces.workspace_exists(workspace).await { - if workspaces.matches(workspace, sha).await.map_err(failed)? { + if workspaces + .matches(workspace, &target.sha) + .await + .map_err(failed)? + { return Ok(WorkspaceAction::Verified); } - workspaces.reset(workspace, sha).await.map_err(failed)?; + workspaces + .reset(workspace, &target.sha) + .await + .map_err(failed)?; return Ok(WorkspaceAction::Reset); } - let key = targets - .iter() - .find(|(_, candidate)| candidate == sha) - .map_or(targets[0].0.key, |(target, _)| target.key); workspaces - .restore(workspace, key, sha) + .restore(workspace, target.key, &target.sha) + .await + .map_err(failed)?; + Ok(WorkspaceAction::Restored) +} + +/// Verify, reset or restore a sandbox workspace onto its target, through +/// the scope's environment: a retained sandbox that still holds the commit +/// is verified or reset in place; a fresh one, or one whose repository +/// lost the commit, is restored from a bundle of the snapshot. +pub async fn bring_sandbox_to( + workspaces: &RunWorkspaces, + env: &Arc, + workspace: &str, + target: &RestoreTarget, +) -> Result { + let failed = |source| RecoveryError::Workspace { + workspace: workspace.to_string(), + source, + }; + if workspaces + .has_commit_in(env, &target.sha) + .await + .map_err(failed)? + { + if workspaces + .matches_in(env, &target.sha) + .await + .map_err(failed)? + { + return Ok(WorkspaceAction::Verified); + } + workspaces + .reset_in(env, &target.sha) + .await + .map_err(failed)?; + return Ok(WorkspaceAction::Reset); + } + workspaces + .restore_in(env, workspace, target.key, &target.sha) .await .map_err(failed)?; Ok(WorkspaceAction::Restored) diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 3f227d033..6546f6932 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -43,6 +43,8 @@ mod support; const HOST_PLUGIN: &str = "sandbox-driver-host"; const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN"; +const DOCKER_PLUGIN: &str = "sandbox-driver-docker"; +const DOCKER_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_DOCKER_PLUGIN"; const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS"; /// The host plugin as Petri's lookup finds it: the override variable, else @@ -99,6 +101,39 @@ fn admit(workflow: &str, settings: &str) -> AdmittedGraphs { } } +/// The Docker plugin as Petri's lookup finds it, with a daemon that +/// answers. `None`, after saying so, when the test should skip; a panic +/// when the environment forbids a skip and the plugin is missing. +fn docker_plugin() -> Option { + let found = env::var_os(DOCKER_PLUGIN_OVERRIDE) + .map(PathBuf::from) + .or_else(|| { + env::split_paths(&env::var_os("PATH")?) + .map(|dir| dir.join(DOCKER_PLUGIN)) + .find(|candidate| candidate.is_file()) + }); + let Some(found) = found else { + assert!( + env::var_os(REQUIRE_ENV).is_none(), + "{REQUIRE_ENV} is set, but {DOCKER_PLUGIN} is not on PATH and {DOCKER_PLUGIN_OVERRIDE} \ + is unset" + ); + eprintln!("skipping: {DOCKER_PLUGIN} is not on PATH and {DOCKER_PLUGIN_OVERRIDE} is unset"); + return None; + }; + let daemon = std::process::Command::new("docker") + .args(["version", "--format", "{{.Server.Version}}"]) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .is_ok_and(|status| status.success()); + if !daemon { + eprintln!("skipping: no Docker daemon answers"); + return None; + } + Some(found) +} + /// One run's pieces: the store, its platform records, where it ran. struct Harness { run_id: RunId, @@ -120,37 +155,70 @@ impl Harness { } } - fn hooks(&self) -> HooksSpec { + fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec { HooksSpec { records: Arc::clone(&self.records) as Arc, author: GitAuthor::default(), checkpoint: RunCheckpointSettings::default(), - host_workspaces: true, + host_workspaces: *provider == SandboxProviderKind::LOCAL, test_gates: None, } } /// Run the bundle to its end through the engine module, as the worker - /// does, and report what the record says. + /// does, on the local provider, and report what the record says. async fn run(&self, workflow: &str, settings: &str) -> engine::RunOutcome { + self.run_on(SandboxProviderKind::LOCAL, workflow, settings) + .await + } + + /// [`run`](Self::run) on `provider`. + async fn run_on( + &self, + provider: SandboxProviderKind, + workflow: &str, + settings: &str, + ) -> engine::RunOutcome { let (interviewer, observers) = no_questions(); + let hooks = self.hooks(&provider); let request = RunRequest { run_id: self.run_id.to_string(), run_dir: self.run_dir.clone(), execution: Execution::Start(admit(workflow, settings)), store: Arc::clone(&self.store) as Arc, runtime: RuntimeSpec::default(), - provider: SandboxProviderKind::LOCAL, + provider, cancel: CancellationToken::new(), interviewer, observers, secrets: None, blobs: None, - hooks: Some(self.hooks()), + hooks: Some(hooks), }; engine::run(request).await.expect("the run executes") } + /// The commits the snapshot repository of `workspace` holds, oldest + /// first, as `(sha, subject, key)`: every checkpoint's history, whatever + /// site committed it. + async fn snapshot_commits( + &self, + workspace: &str, + ) -> Vec<(String, String, Option)> { + let repository = self.workspaces().snapshot_repository(workspace); + // Topological, so the linear run history reads parents first even + // when commits share a timestamp. + let log = git(&repository, &[ + "log", + "--topo-order", + "--reverse", + "--all", + "--format=%H%x00%s%x00%B%x1e", + ]) + .await; + parse_log(&log) + } + async fn inspection(&self) -> RunInspection { let logs = self .store @@ -247,6 +315,10 @@ async fn git(path: &Path, args: &[&str]) -> String { /// The commits on the run branch, oldest first, as `(sha, subject, key)`. async fn commits(path: &Path) -> Vec<(String, String, Option)> { let log = git(path, &["log", "--reverse", "--format=%H%x00%s%x00%B%x1e"]).await; + parse_log(&log) +} + +fn parse_log(log: &str) -> Vec<(String, String, Option)> { log.split('\u{1e}') .filter(|entry| !entry.trim().is_empty()) .map(|entry| { @@ -562,7 +634,7 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() { observers, secrets: None, blobs: None, - hooks: Some(harness.hooks()), + hooks: Some(harness.hooks(&SandboxProviderKind::LOCAL)), }; let outcome = engine::run(request).await.expect("the run executes"); assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); @@ -689,3 +761,104 @@ async fn parallel_branches_checkpoint_the_shared_workspace_in_turn() { assert_eq!(inspection.executions.len(), 3, "the root and two branches"); assert_eq!(harness.workspace().await, "invocation-0-scope-0"); } + +/// On Docker the workspace lives inside the container: every finished +/// stage is committed there, the commit leaves the container as a bundle, +/// and the snapshot repository on the host holds each checkpoint under +/// its ref, with the platform records naming the same commits. +#[tokio::test] +async fn a_docker_run_commits_inside_the_container_and_publishes_every_checkpoint() { + if docker_plugin().is_none() { + return; + } + assert_sandbox_run_publishes_every_checkpoint(SandboxProviderKind::DOCKER).await; +} + +/// The same protocol on Daytona: the sandbox-driver facets are provider +/// neutral, so the commit, the bundle and the restore take one path. Live: +/// it needs `DAYTONA_API_KEY` and the Daytona plugin, and provisions a +/// sandbox. +#[tokio::test] +#[ignore = "requires live Daytona credentials and provisions a sandbox"] +async fn a_daytona_run_commits_inside_the_sandbox_and_publishes_every_checkpoint() { + assert!( + env::var_os("DAYTONA_API_KEY").is_some(), + "DAYTONA_API_KEY must be set to run this live test" + ); + assert_sandbox_run_publishes_every_checkpoint(SandboxProviderKind::DAYTONA).await; +} + +/// Bytes of incompressible data the first stage writes: past the plugin +/// transport's 16 MiB cap on one file read, so its bundle leaves the +/// sandbox in more than one part. +const LARGE_FILE_BYTES: usize = 20 * 1024 * 1024; + +/// A two-stage run on `provider`, whose workspace lives inside a sandbox: +/// nothing of it is on the host, every checkpoint is published, and the +/// bundles carried the stages' files, a large one in parts. +async fn assert_sandbox_run_publishes_every_checkpoint(provider: SandboxProviderKind) { + let harness = Harness::new(); + let workflow = workflow( + &format!( + " write [shape=parallelogram, script=\"echo one > out.txt && head -c \ + {LARGE_FILE_BYTES} /dev/urandom > large.bin\"]\n check [shape=parallelogram, \ + script=\"test \\\"$(cat out.txt)\\\" = one && git log --format=%s | head -1 | grep -q \ + write\"]" + ), + " start -> write -> check -> exit", + ); + let outcome = harness.run_on(provider, &workflow, SETTINGS).await; + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); + + let checkpoints = harness.checkpoints(); + assert_eq!(checkpoints.len(), 4, "{checkpoints:?}"); + let workspace = "invocation-0-scope-0"; + assert!( + !harness.workspaces().workspace_exists(workspace).await, + "nothing of the workspace is on the host" + ); + let published = harness + .workspaces() + .published(workspace) + .await + .expect("the snapshot repository lists"); + let mut by_key: Vec<(CheckpointKey, String)> = published + .iter() + .map(|snapshot| (snapshot.key, snapshot.sha.clone())) + .collect(); + let mut recorded = checkpoints.clone(); + recorded.sort(); + by_key.sort(); + assert_eq!(by_key, recorded, "every record names a published snapshot"); + + let commits = harness.snapshot_commits(workspace).await; + let subjects: Vec<&str> = commits + .iter() + .map(|(_, subject, _)| subject.as_str()) + .collect(); + let run_id = harness.run_id.to_string(); + assert_eq!(subjects, vec![ + format!("fabro({run_id}): start (success)"), + format!("fabro({run_id}): write (success)"), + format!("fabro({run_id}): check (success)"), + format!("fabro({run_id}): exit (success)"), + ]); + let (write_sha, _, _) = &commits[1]; + let repository = harness.workspaces().snapshot_repository(workspace); + assert_eq!( + git(&repository, &["show", &format!("{write_sha}:out.txt")]).await, + "one", + "the bundle carried the stage's files" + ); + assert_eq!( + git(&repository, &[ + "cat-file", + "-s", + &format!("{write_sha}:large.bin") + ]) + .await, + LARGE_FILE_BYTES.to_string(), + "the large file came through the split transfer whole" + ); +}