Merge branch 'petri-integration-sandbox' into petri-integration

# Conflicts:
#	lib/apps/fabro-cli/tests/it/scenario/mod.rs
This commit is contained in:
Bryan Helmkamp 2026-09-18 09:11:19 -04:00
commit b5fdc015f0
No known key found for this signature in database
18 changed files with 1754 additions and 248 deletions

30
Cargo.lock generated
View file

@ -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",

View file

@ -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"

View file

@ -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

View file

@ -10,6 +10,7 @@ mod exec;
mod lifecycle;
mod petri;
mod petri_controls;
mod petri_docker;
mod petri_tools;
mod server_lifecycle;
mod smoke;

View file

@ -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,

View file

@ -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<PathBuf> {
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<String> {
let output = Command::new("docker")
.args([
"ps",
"-aq",
"--filter",
&format!("label=petri.run={run_id}"),
])
.output()
.expect("docker ps runs");
let ids: Vec<String> = 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<PathBuf> = 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<CheckpointKey>)> {
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<CheckpointKey>)]) -> 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<String> {
["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<T>(future: impl std::future::Future<Output = T>) -> 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<String> {
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();
}

View file

@ -3775,6 +3775,7 @@ fn worker_launch_spec(
run_dir: &std::path::Path,
agent_fabro_tools_enabled: bool,
github_app_private_key: Option<String>,
daytona_api_key: Option<String>,
) -> anyhow::Result<WorkerLaunchSpec> {
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<AppState>, 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<AppState>, run_id: RunId) {
&run_dir_for_build,
agent_fabro_tools_enabled,
github_app_private_key,
daytona_api_key,
)
})
.await

View file

@ -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))
}

View file

@ -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

View file

@ -48,6 +48,9 @@ pub(crate) struct WorkerLaunchSpec {
pub(crate) fabro_log: Option<String>,
pub(crate) active_config_path: PathBuf,
pub(crate) github_app_private_key: Option<String>,
/// 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<String>,
/// 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());

View file

@ -14,26 +14,38 @@
//! at `scopes/<workspace id>/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/<workspace id>.git`, under an immutable ref
//! per checkpoint (`refs/checkpoints/<execution>/<firing>/<attempt>`). 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/<workspace id>.git`, under an immutable
//! ref per checkpoint (`refs/checkpoints/<execution>/<firing>/<attempt>`).
//! 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<Self> {
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<dyn ExecEnv>),
}
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<u8>,
stderr: Vec<u8>,
}
/// 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<Snapshot, CheckpointError> {
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<dyn ExecEnv>,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
self.commit_at(
&Site::Sandbox(Arc::clone(env)),
workspace,
key,
node,
status,
)
.await
}
async fn commit_at(
&self,
site: &Site,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
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<bool, CheckpointError> {
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<Option<String>, 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<dyn ExecEnv>,
) -> Result<Option<String>, CheckpointError> {
self.head(&Site::Sandbox(Arc::clone(env))).await
}
/// Whether the workspace sits on `sha` with nothing changed since.
pub async fn matches(&self, workspace: &str, sha: &str) -> Result<bool, CheckpointError> {
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<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
self.matches_at(&Site::Sandbox(Arc::clone(env)), sha).await
}
async fn matches_at(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
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<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
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<dyn ExecEnv>, sha: &str) -> Result<(), CheckpointError> {
self.reset_at(&Site::Sandbox(Arc::clone(env)), sha).await
}
async fn reset_at(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> {
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<dyn ExecEnv>,
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<PathBuf, CheckpointError> {
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<dyn ExecEnv>,
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<dyn ExecEnv>,
site: &Site,
path: &str,
) -> Result<Vec<u8>, 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<String, CheckpointError> {
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<Option<String>, CheckpointError> {
self.git_status(path, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"])
async fn head(&self, site: &Site) -> Result<Option<String>, CheckpointError> {
self.git_status(site, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"])
.await
}
async fn is_clean(&self, path: &Path) -> Result<bool, CheckpointError> {
let status = self.git(path, "status", &["status", "--porcelain"]).await?;
async fn is_clean(&self, site: &Site) -> Result<bool, CheckpointError> {
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<S: AsRef<str>>(
&self,
cwd: &Path,
site: &Site,
action: &str,
args: &[S],
) -> Result<String, CheckpointError> {
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<S: AsRef<str>>(
&self,
cwd: &Path,
site: &Site,
action: &str,
args: &[S],
) -> Result<Option<String>, 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<S: AsRef<str>>(
&self,
cwd: &Path,
site: &Site,
action: &str,
args: &[S],
) -> Result<std::process::Output, CheckpointError> {
) -> Result<GitOutput, CheckpointError> {
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<GitOutput, CheckpointError> {
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<dyn ExecEnv>,
program: &str,
args: &[&str],
action: &str,
) -> Result<GitOutput, CheckpointError> {
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

View file

@ -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<RunOutcome, RunError> {
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<RunOutcome, RunError> {
key.clone(),
request.run_dir.clone(),
Arc::clone(&request.store),
resumed,
))
});
if let Some(hooks) = &fabro_hooks {

View file

@ -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<dyn ExecEnv>);
/// Fabro's `ExecutionHooks`, around the hooks the runtime installed.
pub struct FabroHooks {
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
host_workspaces: bool,
test_gates: Option<PathBuf>,
handle: OnceLock<CoordinatorHandle>,
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
host_workspaces: bool,
test_gates: Option<PathBuf>,
handle: OnceLock<CoordinatorHandle>,
/// The workspace and commit of every checkpoint this process made.
committed: Mutex<HashMap<CheckpointKey, (String, String)>>,
committed: Mutex<HashMap<CheckpointKey, (String, String)>>,
/// Which checkpoints have their platform record, loaded from the store
/// once and kept up to date with every append.
recorded: Mutex<HashSet<CheckpointKey>>,
recorded_loaded: OnceCell<()>,
recorded: Mutex<HashSet<CheckpointKey>>,
recorded_loaded: OnceCell<()>,
/// Inherited workspaces resolved through the run's records.
inherited: Mutex<HashMap<InvocationId, Option<String>>>,
inherited: Mutex<HashMap<InvocationId, Option<String>>>,
/// One lock per workspace: the branches of a parallel node and a nested
/// invocation share their caller's workspace, and Git allows one index
/// operation at a time in it.
workspace_locks: Mutex<HashMap<String, Arc<AsyncMutex<()>>>>,
workspace_locks: Mutex<HashMap<String, Arc<AsyncMutex<()>>>>,
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<Option<String>>,
unreachable_noted: AtomicBool,
failure: Mutex<Option<String>>,
/// The sandbox environment of every acquired scope, by execution and
/// scope, with the workspace id the executor named: where `git` runs
/// when the workspaces are not on this host. Dropped at release.
envs: Mutex<HashMap<(ExecutionId, ScopeId), AcquiredEnv>>,
/// 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<Mutex<BTreeMap<String, RestoreTarget>>>,
store: Arc<dyn RunStore>,
}
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<dyn RunStore>,
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<Option<Note>, 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<Option<Note>, 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<BTreeMap<String, RestoreTarget>>, 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<dyn ExecEnv>,
) -> 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
}
}

View file

@ -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,

View file

@ -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<String, RestoreTarget>,
},
/// 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<Recovery, RecoveryError> {
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<dyn RunStore>,
records: &dyn PlatformRecords,
run_id: &RunId,
workspaces: &RunWorkspaces,
) -> Result<Plan, RecoveryError> {
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<Recovery, RecoveryError
drop(logs);
if let Some(failed) = checkpoint_failure(&inspection.executions) {
return Ok(Recovery::Failed { reason: failed });
}
if !request.host_workspaces {
warn!(
run_id = %request.run_id,
"the run's workspaces are not on this host; resuming on the retained sandbox as it was left"
);
return Ok(Recovery::Resume {
workspaces: Vec::new(),
});
return Ok(Plan::Failed { reason: failed });
}
let workspaces = RunWorkspaces::new(
request.run_dir.clone(),
request.run_id.to_string(),
request.author.clone(),
&request.checkpoint,
);
let lookup = WorkspaceLookup::new(Arc::clone(&request.store), key);
let recorded = recorded_checkpoints(&*request.records, &request.run_id).await?;
let lookup = WorkspaceLookup::new(Arc::clone(&store), key);
let recorded = recorded_checkpoints(records, run_id).await?;
// The snapshot each live execution's workspace must sit on. A live
// execution is one whose log records no exit: `inspect_run` reports it
@ -240,14 +260,7 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
source,
})?;
if let Some(sha) = &sha {
reconcile_record(
&*request.records,
&request.run_id,
target.key,
&workspace,
sha,
)
.await?;
reconcile_record(records, run_id, target.key, &workspace, sha).await?;
}
sha
}
@ -261,7 +274,7 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
}
}
if !found {
return Ok(Recovery::Failed {
return Ok(Plan::Failed {
reason: format!(
"no checkpoint snapshot exists for the last durable finish of execution {} \
(firing {} attempt {}); the run cannot resume on stale files",
@ -271,20 +284,57 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
}
}
let mut targets = BTreeMap::new();
for (workspace, candidates) in candidates {
let sha = newest(workspaces, &workspace, &candidates).await?;
let key = candidates
.iter()
.find(|(_, candidate)| *candidate == sha)
.map_or(candidates[0].0.key, |(target, _)| target.key);
targets.insert(workspace, RestoreTarget { key, sha });
}
Ok(Plan::Resume { targets })
}
/// Decide how the run continues, and bring its host workspaces to their
/// snapshots; a sandbox workspace's target is deferred to the worker.
pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError> {
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<WorkspaceAction, RecoveryError> {
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<dyn ExecEnv>,
workspace: &str,
target: &RestoreTarget,
) -> Result<WorkspaceAction, RecoveryError> {
let failed = |source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
};
if workspaces
.has_commit_in(env, &target.sha)
.await
.map_err(failed)?
{
if workspaces
.matches_in(env, &target.sha)
.await
.map_err(failed)?
{
return Ok(WorkspaceAction::Verified);
}
workspaces
.reset_in(env, &target.sha)
.await
.map_err(failed)?;
return Ok(WorkspaceAction::Reset);
}
workspaces
.restore_in(env, workspace, target.key, &target.sha)
.await
.map_err(failed)?;
Ok(WorkspaceAction::Restored)

View file

@ -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<PathBuf> {
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<dyn PlatformRecords>,
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<dyn petri_store::RunStore>,
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<CheckpointKey>)> {
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<CheckpointKey>)> {
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<CheckpointKey>)> {
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"
);
}

View file

@ -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,

View file

@ -177,7 +177,10 @@ pub fn isolated_env(home_dir: &Path) -> HashMap<String, String> {
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);
}