diff --git a/Cargo.lock b/Cargo.lock index b7e525b78..4b953af02 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6044,7 +6044,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "petri-attractor-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "globset", @@ -6075,7 +6075,7 @@ dependencies = [ [[package]] name = "petri-driver" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -6095,7 +6095,7 @@ dependencies = [ [[package]] name = "petri-engine" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "petri-ir", "serde", @@ -6107,7 +6107,7 @@ dependencies = [ [[package]] name = "petri-execution" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "petri-driver", @@ -6131,7 +6131,7 @@ dependencies = [ [[package]] name = "petri-executor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "libc", @@ -6146,7 +6146,7 @@ dependencies = [ [[package]] name = "petri-executor-sandbox" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "petri-executor", @@ -6168,7 +6168,7 @@ dependencies = [ [[package]] name = "petri-frontend" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "marked-yaml", "petri-ir", @@ -6182,7 +6182,7 @@ dependencies = [ [[package]] name = "petri-frontend-attractor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "minijinja", "petri-frontend", @@ -6199,7 +6199,7 @@ dependencies = [ [[package]] name = "petri-frontend-fabro" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "petri-frontend", "petri-frontend-attractor", @@ -6215,7 +6215,7 @@ dependencies = [ [[package]] name = "petri-frontend-native" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "petri-frontend", "petri-ir", @@ -6226,7 +6226,7 @@ dependencies = [ [[package]] name = "petri-ir" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "regex", "serde", @@ -6239,7 +6239,7 @@ dependencies = [ [[package]] name = "petri-runtime" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "petri-driver", @@ -6260,7 +6260,7 @@ dependencies = [ [[package]] name = "petri-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "petri-executor", @@ -6276,7 +6276,7 @@ dependencies = [ [[package]] name = "petri-store" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -6291,7 +6291,7 @@ dependencies = [ [[package]] name = "petri-testkit" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=83345a8cba9b1352666f615c555860ba03a26499#83345a8cba9b1352666f615c555860ba03a26499" +source = "git+https://github.com/lithoscomputer/petri.git?rev=4d4bdd694aa573514968f8e8321123146d2d2458#4d4bdd694aa573514968f8e8321123146d2d2458" dependencies = [ "async-trait", "petri-driver", diff --git a/Cargo.toml b/Cargo.toml index aa14ca5e1..1f785630d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -132,13 +132,13 @@ pebble-cli-core = { git = "https://github.com/lithoscomputer/pebble", rev = "a39 # lithos-llm and sandbox-driver revisions as this file, so the workspace links # one copy of each. Only `fabro-petri` may depend on these packages; the keys # carry the `petri_` prefix so the crate names say where they come from. -petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-runtime" } -petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-execution" } -petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-store" } -petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-attractor-steps" } -petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-frontend-attractor" } -petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-frontend-fabro" } -petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "83345a8cba9b1352666f615c555860ba03a26499", package = "petri-testkit" } +petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-runtime" } +petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-execution" } +petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-store" } +petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-attractor-steps" } +petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-frontend-attractor" } +petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-frontend-fabro" } +petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "4d4bdd694aa573514968f8e8321123146d2d2458", package = "petri-testkit" } sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] } fork = "0.2" exec = "0.3" 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 d72260531..a67ad14cf 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -37,8 +37,9 @@ //! the worker exits with that 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 ecfc89fc9..27198d1fa 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/mod.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/mod.rs @@ -10,6 +10,7 @@ mod exec; mod lifecycle; mod petri; mod petri_controls; +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 01c2c6d3a..f0ad0ec44 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 @@ -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 @@ -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!( @@ -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/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index a4e02eae2..39f6dec53 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -3775,6 +3775,7 @@ fn worker_launch_spec( run_dir: &std::path::Path, agent_fabro_tools_enabled: bool, github_app_private_key: Option, + daytona_api_key: Option, ) -> anyhow::Result { let current_exe = std::env::current_exe().context("reading current executable path")?; let executable = @@ -3813,6 +3814,7 @@ fn worker_launch_spec( fabro_log, active_config_path: state.active_config_path().to_path_buf(), github_app_private_key, + daytona_api_key, fabro_home: fabro_config::Home::from_env().root().to_path_buf(), }) } @@ -4458,7 +4460,21 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { return; } - let github_app_private_key = match state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await { + // A Daytona run's worker hands the vault's key to Petri's Daytona + // plugin through its own environment; any other run's worker never + // sees it. + let wants_daytona = + run_state.spec.settings.run.environment.provider == SandboxProviderKind::DAYTONA; + let secrets = async { + let github_app_private_key = state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await?; + let daytona_api_key = if wants_daytona { + state.vault_secret(EnvVars::DAYTONA_API_KEY).await? + } else { + None + }; + Ok::<_, SecretStoreError>((github_app_private_key, daytona_api_key)) + }; + let (github_app_private_key, daytona_api_key) = match secrets.await { Ok(value) => value, Err(err) => { fail_run_before_execution( @@ -4482,6 +4498,7 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { &run_dir_for_build, agent_fabro_tools_enabled, github_app_private_key, + daytona_api_key, ) }) .await diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 75333bfe1..724422b1b 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -2402,6 +2402,7 @@ fn worker_command_forwards_github_app_private_key_from_vault() { storage_dir.path(), false, Some("test-private-key".to_string()), + None, ) .unwrap(); let cmd = LocalWorkerRuntime::command_for_spec(&spec); @@ -2410,6 +2411,35 @@ fn worker_command_forwards_github_app_private_key_from_vault() { command_env_value(&cmd, EnvVars::GITHUB_APP_PRIVATE_KEY), EnvOverride::Set("test-private-key".to_string()) ); + assert_eq!( + command_env_value(&cmd, EnvVars::DAYTONA_API_KEY), + EnvOverride::Unchanged + ); +} + +/// A Daytona run's worker carries the vault's key for Petri's Daytona +/// plugin. +#[cfg(unix)] +#[test] +fn worker_command_forwards_daytona_api_key_from_vault() { + let storage_dir = tempfile::tempdir().unwrap(); + let state = worker_command_test_state(storage_dir.path(), &["dev-token"], Some(TEST_DEV_TOKEN)); + let spec = worker_launch_spec( + state.as_ref(), + RunId::new(), + RunExecutionMode::Start, + storage_dir.path(), + false, + None, + Some("dtn_test-key".to_string()), + ) + .unwrap(); + let cmd = LocalWorkerRuntime::command_for_spec(&spec); + + assert_eq!( + command_env_value(&cmd, EnvVars::DAYTONA_API_KEY), + EnvOverride::Set("dtn_test-key".to_string()) + ); } #[cfg(unix)] @@ -2661,6 +2691,7 @@ fn worker_command( run_dir, agent_fabro_tools_enabled, None, + None, )?; Ok(LocalWorkerRuntime::command_for_spec(&spec)) } diff --git a/lib/apps/fabro-server/src/spawn_env.rs b/lib/apps/fabro-server/src/spawn_env.rs index 351939d7c..8510c5a94 100644 --- a/lib/apps/fabro-server/src/spawn_env.rs +++ b/lib/apps/fabro-server/src/spawn_env.rs @@ -62,6 +62,21 @@ const WORKER_ENV_ALLOWLIST: &[&str] = &[ EnvVars::PETRI_SANDBOX_PLUGIN_DEV, EnvVars::PETRI_SANDBOX_DOCKER_HOST_ADDRESS, EnvVars::PETRI_SANDBOX_ACTION_HOST_IMAGE, + // The Docker daemon selection: the worker's Docker plugin reads these + // from its own process, so the worker's sandboxes go to the daemon the + // server uses (a remote or TLS daemon, a named context), not the + // default socket. + EnvVars::DOCKER_HOST, + EnvVars::DOCKER_TLS_VERIFY, + EnvVars::DOCKER_CERT_PATH, + EnvVars::DOCKER_API_VERSION, + EnvVars::DOCKER_CONFIG, + EnvVars::DOCKER_CONTEXT, + // Daytona's control-plane selection, the non-secret half: the plugin + // reads them from the worker. The API key comes from the vault, set on + // the command by the launch (`WorkerLaunchSpec::daytona_api_key`). + EnvVars::DAYTONA_API_URL, + EnvVars::DAYTONA_ORGANIZATION_ID, // A test's checkpoint gates: the worker's hooks hold at a named point // until the test releases them, so a crash can be placed there. EnvVars::FABRO_TEST_CHECKPOINT_GATES, @@ -164,6 +179,27 @@ mod tests { "/opt/petri/sandbox-driver-host".to_string(), ), ("PETRI_SANDBOX_PLUGIN_DEV".to_string(), "1".to_string()), + ( + "DOCKER_HOST".to_string(), + "tcp://build-daemon.internal:2376".to_string(), + ), + ("DOCKER_TLS_VERIFY".to_string(), "1".to_string()), + ( + "DOCKER_CERT_PATH".to_string(), + "/etc/docker/certs".to_string(), + ), + ("DOCKER_API_VERSION".to_string(), "1.47".to_string()), + ( + "DOCKER_CONFIG".to_string(), + "/etc/docker/client".to_string(), + ), + ("DOCKER_CONTEXT".to_string(), "build".to_string()), + ( + "DAYTONA_API_URL".to_string(), + "https://daytona.internal/api".to_string(), + ), + ("DAYTONA_ORGANIZATION_ID".to_string(), "org-1".to_string()), + ("DAYTONA_API_KEY".to_string(), "leak".to_string()), ]); let mut cmd = env_command(); apply_allowlist(&mut cmd, WORKER_ENV_ALLOWLIST, &|name| { @@ -208,6 +244,43 @@ mod tests { actual.get("PETRI_SANDBOX_PLUGIN_DEV").map(String::as_str), Some("1") ); + // The Docker daemon selection crosses whole, so the worker's Docker + // plugin drives the daemon the server uses. + assert_eq!( + actual.get("DOCKER_HOST").map(String::as_str), + Some("tcp://build-daemon.internal:2376") + ); + assert_eq!( + actual.get("DOCKER_TLS_VERIFY").map(String::as_str), + Some("1") + ); + assert_eq!( + actual.get("DOCKER_CERT_PATH").map(String::as_str), + Some("/etc/docker/certs") + ); + assert_eq!( + actual.get("DOCKER_API_VERSION").map(String::as_str), + Some("1.47") + ); + assert_eq!( + actual.get("DOCKER_CONFIG").map(String::as_str), + Some("/etc/docker/client") + ); + assert_eq!( + actual.get("DOCKER_CONTEXT").map(String::as_str), + Some("build") + ); + // Daytona's non-secret selectors cross; its key is the vault's, + // never the server's environment. + assert_eq!( + actual.get("DAYTONA_API_URL").map(String::as_str), + Some("https://daytona.internal/api") + ); + assert_eq!( + actual.get("DAYTONA_ORGANIZATION_ID").map(String::as_str), + Some("org-1") + ); + assert!(!actual.contains_key("DAYTONA_API_KEY")); assert_eq!(actual.get("CLICOLOR").map(String::as_str), Some("0")); assert_eq!(actual.get("CLICOLOR_FORCE").map(String::as_str), Some("1")); // Bedrock SigV4 chain inputs cross into the worker so it can re-resolve diff --git a/lib/apps/fabro-server/src/worker_runtime.rs b/lib/apps/fabro-server/src/worker_runtime.rs index 7b4fa39b6..21271494b 100644 --- a/lib/apps/fabro-server/src/worker_runtime.rs +++ b/lib/apps/fabro-server/src/worker_runtime.rs @@ -48,6 +48,9 @@ pub(crate) struct WorkerLaunchSpec { pub(crate) fabro_log: Option, pub(crate) active_config_path: PathBuf, pub(crate) github_app_private_key: Option, + /// The vault's Daytona API key, for a run on a Daytona environment: + /// Petri's Daytona plugin reads it from the worker's process. + pub(crate) daytona_api_key: Option, /// The Fabro home the server resolved, so a Petri run's skills step /// reads the same home whatever the worker's environment says. pub(crate) fabro_home: PathBuf, @@ -109,6 +112,9 @@ impl LocalWorkerRuntime { if let Some(pem) = spec.github_app_private_key.as_deref() { cmd.env(EnvVars::GITHUB_APP_PRIVATE_KEY, pem); } + if let Some(key) = spec.daytona_api_key.as_deref() { + cmd.env(EnvVars::DAYTONA_API_KEY, key); + } #[cfg(unix)] fabro_proc::pre_exec_setpgid(cmd.as_std_mut()); 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 a1ed6daff..3f4d6065b 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -33,9 +33,11 @@ //! a resume alike, so a run that was paused resumes paused. //! //! 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 @@ -54,7 +56,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}; @@ -173,6 +175,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) @@ -196,6 +205,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 1b9308bda..0a4d392af 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 cd7b2bd6a..2c17eed9e 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -44,6 +44,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 @@ -100,6 +102,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, @@ -121,38 +156,71 @@ 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(), controls: RunControls::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 @@ -249,6 +317,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| { @@ -565,7 +637,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:?}"); @@ -692,3 +764,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" + ); +} diff --git a/lib/foundation/fabro-static/src/env_vars.rs b/lib/foundation/fabro-static/src/env_vars.rs index a14c8ce7e..ff3f29fbb 100644 --- a/lib/foundation/fabro-static/src/env_vars.rs +++ b/lib/foundation/fabro-static/src/env_vars.rs @@ -74,6 +74,29 @@ impl EnvVars { Self::PETRI_SANDBOX_ACTION_HOST_IMAGE, ]; + // The Docker daemon selection the Docker CLI and its client libraries + // read: which daemon, over which transport, with which TLS material, + // client configuration and context. Petri's Docker plugin forwards them + // from the process that launches it, so a run's worker must carry the + // server's. + pub const DOCKER_HOST: &'static str = "DOCKER_HOST"; + pub const DOCKER_TLS_VERIFY: &'static str = "DOCKER_TLS_VERIFY"; + pub const DOCKER_CERT_PATH: &'static str = "DOCKER_CERT_PATH"; + pub const DOCKER_API_VERSION: &'static str = "DOCKER_API_VERSION"; + pub const DOCKER_CONFIG: &'static str = "DOCKER_CONFIG"; + pub const DOCKER_CONTEXT: &'static str = "DOCKER_CONTEXT"; + + /// Every Docker daemon selection variable, in one list for the process + /// boundaries that forward them. + pub const DOCKER_VARS: &'static [&'static str] = &[ + Self::DOCKER_HOST, + Self::DOCKER_TLS_VERIFY, + Self::DOCKER_CERT_PATH, + Self::DOCKER_API_VERSION, + Self::DOCKER_CONFIG, + Self::DOCKER_CONTEXT, + ]; + // LLM providers and tool integrations pub const ANTHROPIC_API_KEY: &'static str = "ANTHROPIC_API_KEY"; pub const AWS_BEARER_TOKEN_BEDROCK: &'static str = "AWS_BEARER_TOKEN_BEDROCK"; @@ -239,6 +262,12 @@ mod tests { EnvVars::PETRI_SANDBOX_PLUGIN_DEV, EnvVars::PETRI_SANDBOX_DOCKER_HOST_ADDRESS, EnvVars::PETRI_SANDBOX_ACTION_HOST_IMAGE, + EnvVars::DOCKER_HOST, + EnvVars::DOCKER_TLS_VERIFY, + EnvVars::DOCKER_CERT_PATH, + EnvVars::DOCKER_API_VERSION, + EnvVars::DOCKER_CONFIG, + EnvVars::DOCKER_CONTEXT, EnvVars::ANTHROPIC_API_KEY, EnvVars::ANTHROPIC_BASE_URL, EnvVars::AWS_BEARER_TOKEN_BEDROCK, diff --git a/lib/foundation/fabro-test/src/lib.rs b/lib/foundation/fabro-test/src/lib.rs index 44a665d13..0ddafda47 100644 --- a/lib/foundation/fabro-test/src/lib.rs +++ b/lib/foundation/fabro-test/src/lib.rs @@ -177,7 +177,10 @@ pub fn isolated_env(home_dir: &Path) -> HashMap { if let Some(path) = std::env::var_os(EnvVars::PATH).and_then(|value| value.into_string().ok()) { env.insert(EnvVars::PATH.to_string(), path); } - for name in EnvVars::PETRI_SANDBOX_PLUGIN_VARS { + for name in EnvVars::PETRI_SANDBOX_PLUGIN_VARS + .iter() + .chain(EnvVars::DOCKER_VARS) + { if let Some(value) = std::env::var_os(name).and_then(|value| value.into_string().ok()) { env.insert((*name).to_string(), value); } @@ -223,8 +226,12 @@ fn apply_test_isolation_with_lookup( } // Petri resolves its sandbox-driver plugins from these, in the server a // test starts and in the workers that server launches; a developer's - // plugin override reaches them like `PATH` does. - for name in EnvVars::PETRI_SANDBOX_PLUGIN_VARS { + // plugin override reaches them like `PATH` does, and so does the Docker + // daemon selection the Docker plugin needs. + for name in EnvVars::PETRI_SANDBOX_PLUGIN_VARS + .iter() + .chain(EnvVars::DOCKER_VARS) + { if let Some(value) = lookup(name) { cmd.env(name, value); }