From cd5f56a4ede3ad0f14dc88738e3e509c96f57ab7 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 1 Oct 2026 17:12:03 -0400 Subject: [PATCH 1/2] Enforce resolved network policy for Docker and Daytona runs --- Cargo.lock | 31 +-- .../src/commands/run/petri_worker.rs | 8 + .../fabro-server/src/server/petri_runs.rs | 6 +- .../fabro-server/tests/it/scenario/petri.rs | 39 ++- lib/components/fabro-petri/Cargo.toml | 1 + lib/components/fabro-petri/src/engine.rs | 20 +- lib/components/fabro-petri/tests/daytona.rs | 127 ++++++--- lib/components/fabro-petri/tests/hooks.rs | 4 +- lib/components/fabro-petri/tests/network.rs | 246 ++++++++++++++++++ .../fabro-petri/tests/support/mod.rs | 3 +- 10 files changed, 433 insertions(+), 52 deletions(-) create mode 100644 lib/components/fabro-petri/tests/network.rs diff --git a/Cargo.lock b/Cargo.lock index 9eeadd5c9..21d47f73a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2529,6 +2529,7 @@ dependencies = [ "fabro-http", "fabro-interview", "fabro-llm", + "fabro-macros", "fabro-petri", "fabro-static", "fabro-store", @@ -5160,7 +5161,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "petri-attractor-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "globset", @@ -5191,7 +5192,7 @@ dependencies = [ [[package]] name = "petri-driver" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5211,7 +5212,7 @@ dependencies = [ [[package]] name = "petri-engine" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "petri-ir", "serde", @@ -5223,7 +5224,7 @@ dependencies = [ [[package]] name = "petri-execution" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "petri-driver", @@ -5247,7 +5248,7 @@ dependencies = [ [[package]] name = "petri-executor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "libc", @@ -5262,7 +5263,7 @@ dependencies = [ [[package]] name = "petri-executor-sandbox" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "petri-executor", @@ -5284,7 +5285,7 @@ dependencies = [ [[package]] name = "petri-frontend" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "marked-yaml", "petri-ir", @@ -5298,7 +5299,7 @@ dependencies = [ [[package]] name = "petri-frontend-attractor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "minijinja", "petri-frontend", @@ -5315,7 +5316,7 @@ dependencies = [ [[package]] name = "petri-frontend-fabro" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "petri-frontend", "petri-frontend-attractor", @@ -5331,7 +5332,7 @@ dependencies = [ [[package]] name = "petri-frontend-native" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "petri-frontend", "petri-ir", @@ -5342,7 +5343,7 @@ dependencies = [ [[package]] name = "petri-ir" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "regex", "serde", @@ -5355,7 +5356,7 @@ dependencies = [ [[package]] name = "petri-runtime" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "petri-driver", @@ -5376,7 +5377,7 @@ dependencies = [ [[package]] name = "petri-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "petri-executor", @@ -5392,7 +5393,7 @@ dependencies = [ [[package]] name = "petri-store" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5407,7 +5408,7 @@ dependencies = [ [[package]] name = "petri-testkit" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#98144f107fcdd290a44699a8e76abf6e6f63fe2a" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#91d1b77b3927cb3fefcc4a412fdf89ce04857277" dependencies = [ "async-trait", "petri-driver", diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index f26eb70d9..cc491e007 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -263,6 +263,14 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { .environment .resources .clone(), + network: worker + .run_state + .spec + .settings + .run + .environment + .network + .clone(), cancel: cancel_token.clone(), controls: controls.clone(), interviewer: Arc::new(petri_interviewer), diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index a7499ab4e..eafe55510 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -140,8 +140,9 @@ fn settings_layer_toml(state: &AppState) -> Option { /// under `daytona`, and `env`. The rest is the platform's (`cwd`, /// `network`, `lifecycle`, `labels`, `image.dockerfile`, resources the host /// and Docker providers run without, an image the host runs without) and -/// stays with the server's own resolution; handing it to Petri would only -/// warn `ignored.workflow_toml.environments..` on every admit. +/// stays with the server's own resolution. The resolved network policy is +/// passed through `RunRequest` at execution; handing it to the frontend would +/// only warn `ignored.workflow_toml.environments..` on every admit. fn petri_environments(catalog: &MergeMap) -> MergeMap { MergeMap( catalog @@ -497,6 +498,7 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { runtime, provider: run_state.spec.settings.run.environment.provider.clone(), resources: run_state.spec.settings.run.environment.resources.clone(), + network: run_state.spec.settings.run.environment.network.clone(), cancel, // The in-process test path drives no pause: the server's transport // for it names the worker. A steer or an interrupt is answered in diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index 7541469ae..a9706bea3 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -1229,12 +1229,21 @@ const CATALOG_IMAGE: &str = "ghcr.io/lithoscomputer/ubuntu-22.04:slim"; /// A Docker environment in the server's catalog, with the image it runs. async fn create_docker_environment(app: &axum::Router, id: &str, image: &str) { + create_docker_environment_with_network(app, id, image, "allow_all").await; +} + +async fn create_docker_environment_with_network( + app: &axum::Router, + id: &str, + image: &str, + mode: &str, +) { let environment = serde_json::json!({ "id": id, "provider": "docker", "image": { "docker": image, "dockerfile": null }, "resources": { "cpu": null, "memory": null, "disk": null }, - "network": { "mode": "allow_all", "allow": [] }, + "network": { "mode": mode, "allow": [] }, "lifecycle": { "preserve": false, "stop_on_terminal": true, "auto_stop": null }, "labels": {}, "env": {} @@ -1331,10 +1340,21 @@ async fn admitted_root_graph(app: &axum::Router, run_id: &str) -> serde_json::Va /// a Docker daemon, the run's container runs that image. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn a_bundle_naming_a_catalog_environment_runs_on_docker_with_its_image() { + assert_catalog_environment_runs("allow_all", "bridge").await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +#[ignore = "requires Docker; creates and deletes a container"] +async fn a_catalog_network_block_reaches_the_created_docker_container() { + assert!(fabro_test::docker_available()); + assert_catalog_environment_runs("block", "none").await; +} + +async fn assert_catalog_environment_runs(mode: &str, network_mode: &str) { let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n"); let state = test_app_state_with_options(settings, 5); let app = test_app_with_scheduler(Arc::clone(&state)); - create_docker_environment(&app, "docker-small", CATALOG_IMAGE).await; + create_docker_environment_with_network(&app, "docker-small", CATALOG_IMAGE, mode).await; let version_id = register_version(&app, &[ ("workflow.fabro", COMMAND_DOT), @@ -1381,6 +1401,16 @@ async fn a_bundle_naming_a_catalog_environment_runs_on_docker_with_its_image() { instance["image"], CATALOG_IMAGE, "the container runs the catalog's image: {instance}" ); + let container_id = instance["runtime"]["id"].as_str().expect("container id"); + let inspection = Command::new("docker") + .args([ + "inspect", + "--format", + "{{.HostConfig.NetworkMode}}", + container_id, + ]) + .output() + .expect("Docker inspect runs"); let output = Command::new("docker") .args([ "ps", @@ -1401,6 +1431,11 @@ async fn a_bundle_naming_a_catalog_environment_runs_on_docker_with_its_image() { .stderr(Stdio::null()) .status(); } + assert!(inspection.status.success(), "{inspection:?}"); + assert_eq!( + String::from_utf8(inspection.stdout).unwrap().trim(), + network_mode + ); } /// A bundle's own `[environments.]` table wins over the server's, key diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 86ae46e56..f9cd9a456 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -57,6 +57,7 @@ tokio-util.workspace = true tracing.workspace = true [dev-dependencies] +fabro-macros = { path = "../../foundation/fabro-macros" } fabro-petri = { path = ".", features = ["test-support"] } fabro-tool = { path = "../fabro-tool" } httpmock = "0.8" diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 88811caa7..8820fd276 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -47,7 +47,9 @@ use std::path::PathBuf; use std::sync::Arc; -use fabro_types::settings::run::EnvironmentResourcesSettings; +use fabro_types::settings::run::{ + EnvironmentNetworkMode, EnvironmentNetworkSettings, EnvironmentResourcesSettings, +}; use fabro_types::settings::size::Size; use fabro_types::{FailureReason, RunId, SandboxProviderKind}; use petri_execution::host::{self, HostError, HostRun}; @@ -60,6 +62,7 @@ use petri_runtime::driver::lifecycle::ExecutionHooks; pub use petri_runtime::executor::Retention; use petri_runtime::executor::SecretProvider; use petri_runtime::{DaytonaResources, LostSandbox, RunOptions, SandboxBackend}; +use sandbox_driver::NetworkPolicy; use tokio::fs; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; @@ -98,6 +101,8 @@ pub struct RunRequest { pub provider: SandboxProviderKind, /// Resolved environment resources for the Daytona runner snapshot. pub resources: EnvironmentResourcesSettings, + /// Resolved run policy, applied to every sandbox at execution. + pub network: EnvironmentNetworkSettings, /// Fires to cancel the run. pub cancel: CancellationToken, /// The run's pause, unpause and steer controls, which the caller keeps @@ -176,6 +181,16 @@ pub enum Conclusion { }, } +fn network_policy(settings: &EnvironmentNetworkSettings) -> NetworkPolicy { + match settings.mode { + EnvironmentNetworkMode::AllowAll => NetworkPolicy::AllowAll, + EnvironmentNetworkMode::Block => NetworkPolicy::Block, + EnvironmentNetworkMode::CidrAllowList => NetworkPolicy::CidrAllowList { + cidrs: settings.allow.clone(), + }, + } +} + /// Execute the run to its end and report what the record says. pub async fn run(request: RunRequest) -> Result { let backend = if request.runtime.dry_run { @@ -192,6 +207,9 @@ pub async fn run(request: RunRequest) -> Result { options.run_key = Some(key.clone()); options.retention = RETENTION; options.sandbox.backend = backend; + if !request.runtime.dry_run { + options.sandbox.network = network_policy(&request.network); + } if backend == SandboxBackend::Daytona { options.sandbox.daytona_resources = daytona_resources(&request.resources)?; } diff --git a/lib/components/fabro-petri/tests/daytona.rs b/lib/components/fabro-petri/tests/daytona.rs index 8d4c72e5d..bba242305 100644 --- a/lib/components/fabro-petri/tests/daytona.rs +++ b/lib/components/fabro-petri/tests/daytona.rs @@ -10,7 +10,9 @@ use fabro_petri::engine::{self, RunStatus}; use fabro_petri::providers::{DaytonaCredentials, SandboxProviderConfig}; use fabro_petri::runtime::RuntimeSpec; use fabro_types::SandboxProviderKind; -use fabro_types::settings::run::EnvironmentResourcesSettings; +use fabro_types::settings::run::{ + EnvironmentNetworkMode, EnvironmentNetworkSettings, EnvironmentResourcesSettings, +}; use httpmock::prelude::*; use petri_store::MemoryRunStore; use serde_json::json; @@ -70,30 +72,7 @@ async fn assert_snapshot_request( memory: u64, disk: Option, ) { - let server = MockServer::start_async().await; - server - .mock_async(|when, then| { - when.method(GET).path("/api-keys/current"); - then.status(200) - .header("content-type", "application/json") - .json_body(json!({ - "name": "test-key", - "organizationId": "test-org", - "permissions": [ - "write:snapshots", "delete:snapshots", - "write:sandboxes", "delete:sandboxes" - ] - })); - }) - .await; - server - .mock_async(|when, then| { - when.method(GET).path("/sandbox"); - then.status(200) - .header("content-type", "application/json") - .json_body(json!({"items": [], "nextCursor": null})); - }) - .await; + let server = provider_server().await; server .mock_async(|when, then| { when.method(GET).path_matches(r"^/snapshots/[^/]+$"); @@ -126,6 +105,45 @@ async fn assert_snapshot_request( .json_body(json!({"message": "snapshot creation stopped by test"})); }) .await; + let outcome = run_daytona(&server, resources, EnvironmentNetworkSettings::default()).await; + + assert_eq!(create.calls_async().await, 1, "{outcome:?}"); + assert_eq!(outcome.status, RunStatus::Failed); +} + +async fn provider_server() -> MockServer { + let server = MockServer::start_async().await; + server + .mock_async(|when, then| { + when.method(GET).path("/api-keys/current"); + then.status(200) + .header("content-type", "application/json") + .json_body(json!({ + "name": "test-key", + "organizationId": "test-org", + "permissions": [ + "write:snapshots", "delete:snapshots", + "write:sandboxes", "delete:sandboxes" + ] + })); + }) + .await; + server + .mock_async(|when, then| { + when.method(GET).path("/sandbox"); + then.status(200) + .header("content-type", "application/json") + .json_body(json!({"items": [], "nextCursor": null})); + }) + .await; + server +} + +async fn run_daytona( + server: &MockServer, + resources: EnvironmentResourcesSettings, + network: EnvironmentNetworkSettings, +) -> engine::RunOutcome { let credentials = DaytonaCredentials::new("test-key".to_string()) .with_api_url(Some(server.base_url())) .with_http_client(Some(fabro_test::test_http_client())); @@ -160,11 +178,60 @@ async fn assert_snapshot_request( ); request.provider = SandboxProviderKind::DAYTONA; request.resources = resources; + request.network = network; - let outcome = engine::run(request) + engine::run(request) .await - .expect("the run records its failure"); - - assert_eq!(create.calls_async().await, 1, "{outcome:?}"); - assert_eq!(outcome.status, RunStatus::Failed); + .expect("the run records its failure") +} + +#[tokio::test] +async fn daytona_create_requests_enforce_block_and_preserve_allow_all() { + for mode in [ + EnvironmentNetworkMode::Block, + EnvironmentNetworkMode::AllowAll, + ] { + let server = provider_server().await; + server + .mock_async(|when, then| { + when.method(GET).path_matches(r"^/snapshots/[^/]+$"); + then.status(200) + .header("content-type", "application/json") + .json_body(json!({ + "id": "runner", "name": "runner", "general": false, + "state": "active", "sandboxClass": "container", + "cpu": 2, "gpu": 0, "mem": 4, "disk": 3, + "createdAt": "2026-01-01T00:00:00Z", "updatedAt": "2026-01-01T00:00:00Z", + "size": null, "entrypoint": null, "errorReason": null, + "lastUsedAt": null, "sourceSandboxId": null + })); + }) + .await; + let create = server + .mock_async(|when, then| { + when.method(POST).path("/sandbox").json_body_includes( + json!({ + "networkBlockAll": mode == EnvironmentNetworkMode::Block + }) + .to_string(), + ); + // Observe the real SDK request, then refuse creation. No cloud + // sandbox or simulated shell is needed to prove this contract. + then.status(400) + .header("content-type", "application/json") + .json_body(json!({"message": "creation stopped by test"})); + }) + .await; + let outcome = run_daytona( + &server, + EnvironmentResourcesSettings::default(), + EnvironmentNetworkSettings { + mode, + allow: Vec::new(), + }, + ) + .await; + assert_eq!(create.calls_async().await, 1, "{mode}: {outcome:?}"); + assert_eq!(outcome.status, RunStatus::Failed); + } } diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 2636732ed..0a118a3a9 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -38,7 +38,7 @@ use fabro_petri::source::{RunSource, SourceRevision}; use fabro_petri::test_support::{MemoryBlobs, MemoryPlatformRecords}; use fabro_store::{ArtifactStore, PlatformRecord, PlatformRecordKind}; use fabro_types::settings::run::{ - EnvironmentResourcesSettings, RunCheckpointSettings, RunNamespace, + EnvironmentNetworkSettings, EnvironmentResourcesSettings, RunCheckpointSettings, RunNamespace, }; use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind}; use object_store::local::LocalFileSystem; @@ -202,6 +202,7 @@ impl Harness { }, provider, resources: EnvironmentResourcesSettings::default(), + network: EnvironmentNetworkSettings::default(), cancel: CancellationToken::new(), controls: RunControls::new(), interviewer, @@ -848,6 +849,7 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() { }, provider: SandboxProviderKind::LOCAL, resources: EnvironmentResourcesSettings::default(), + network: EnvironmentNetworkSettings::default(), cancel: CancellationToken::new(), controls: RunControls::new(), interviewer, diff --git a/lib/components/fabro-petri/tests/network.rs b/lib/components/fabro-petri/tests/network.rs new file mode 100644 index 000000000..5a424b69b --- /dev/null +++ b/lib/components/fabro-petri/tests/network.rs @@ -0,0 +1,246 @@ +//! Live Docker coverage through the same engine assembly as a run worker. + +mod support; + +use std::env; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; + +use fabro_petri::check::Launch; +use fabro_petri::engine::{self, RunStatus}; +use fabro_petri::providers::SandboxProviderConfig; +use fabro_petri::prune::{self, PruneRequest}; +use fabro_petri::runtime::RuntimeSpec; +use fabro_types::settings::run::{EnvironmentNetworkMode, EnvironmentNetworkSettings}; +use fabro_types::{RunId, SandboxProviderKind}; +use petri_store::MemoryRunStore; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; +use tokio::process::Command; +use tokio::time::{self, Instant}; +use tokio_util::task::AbortOnDropHandle; + +#[fabro_macros::e2e_test()] +async fn docker_block_prevents_canary_access_while_allow_all_preserves_it() { + let listener = TcpListener::bind("0.0.0.0:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let hits = Arc::new(AtomicUsize::new(0)); + let seen = hits.clone(); + let _canary = AbortOnDropHandle::new(tokio::spawn(async move { + loop { + let (mut stream, _) = listener.accept().await.unwrap(); + let mut request = [0; 4096]; + if stream.read(&mut request).await.unwrap() > 0 { + seen.fetch_add(1, Ordering::SeqCst); + stream.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 14\r\nConnection: close\r\n\r\nnetwork-canary").await.unwrap(); + } + } + })); + for (mode, network_mode, expected_hits) in [ + (EnvironmentNetworkMode::AllowAll, "bridge", 1), + (EnvironmentNetworkMode::Block, "none", 1), + ] { + let root = tempfile::tempdir().unwrap(); + let run_id = RunId::new().to_string(); + let store = Arc::new(MemoryRunStore::new()); + let runtime = RuntimeSpec::default(); + let curl = format!( + "curl --noproxy '*' --fail --silent --max-time 2 http://host.docker.internal:{port}/canary" + ); + let script = if mode == EnvironmentNetworkMode::Block { + format!("if {curl}; then exit 19; fi") + } else { + format!("test $({curl}) = network-canary") + }; + let workflow = format!( + r#"digraph Network {{ + start [shape=Mdiamond] + probe [shape=parallelogram, script="{script}"] + exit [shape=Msquare] + start -> probe -> exit + }}"# + ); + let graphs = support::admit( + &[ + ("workflow.toml", support::SETTINGS), + ("workflow.fabro", &workflow), + ], + Launch::default(), + &runtime, + ); + let mut request = support::run_request( + &run_id, + root.path(), + graphs, + store.clone(), + runtime, + support::no_questions(Arc::new(support::Silent)), + ); + request.provider = SandboxProviderKind::DOCKER; + request.network = EnvironmentNetworkSettings { + mode, + allow: Vec::new(), + }; + let outcome = engine::run(request).await; + // Inspect Docker independently of Fabro's settings and driver status. + let containers = Command::new("docker") + .args([ + "ps", + "-aq", + "--filter", + &format!("label=petri.run={run_id}"), + ]) + .output() + .await + .unwrap(); + let ids = String::from_utf8(containers.stdout).unwrap(); + let inspected = Command::new("docker") + .args([ + "inspect", + "--format", + "{{.HostConfig.NetworkMode}}", + ids.trim(), + ]) + .output() + .await + .unwrap(); + // Cleanup before assertions, including when execution or inspection failed. + let cleanup = prune::prune(PruneRequest { + sandbox: SandboxProviderConfig::default(), + run_id, + run_dir: root.path().to_path_buf(), + store, + provider: SandboxProviderKind::DOCKER, + }) + .await; + assert!(cleanup.unwrap().is_clean()); + let outcome = outcome.unwrap(); + assert_eq!(outcome.status, RunStatus::Success, "{mode}: {outcome:?}"); + assert!(inspected.status.success(), "{inspected:?}"); + assert_eq!( + String::from_utf8(inspected.stdout).unwrap().trim(), + network_mode + ); + assert_eq!(hits.load(Ordering::SeqCst), expected_hits, "{mode}"); + } +} + +/// Runs only when explicitly requested with live credentials. It reuses the +/// runner snapshot and deletes both task-owned sandboxes through Petri. +#[fabro_macros::e2e_test(live("DAYTONA_API_KEY"), live("FABRO_TEST_DAYTONA_RUNNER_SNAPSHOT"))] +async fn daytona_block_prevents_outbound_https_while_allow_all_preserves_it() { + use fabro_petri::providers::{self, DaytonaCredentials}; + use sandbox_driver::{NetworkPolicy, SandboxFilter, SandboxState, SnapshotFilter}; + + let credentials = DaytonaCredentials::from_api_key( + fabro_test::require_env("DAYTONA_API_KEY").expect("guard checked the credential"), + provider_env, + ); + let provider = providers::connect_daytona(&credentials).await.unwrap(); + let snapshots = provider.snapshots().unwrap(); + let before = snapshots.list(&SnapshotFilter::default()).await.unwrap(); + // Live validation requires the standard runner snapshot to exist already. + // Its exact name is supplied by the operator after checking the pinned + // Petri runner, so the test does not create or delete shared snapshots. + let runner = fabro_test::require_env("FABRO_TEST_DAYTONA_RUNNER_SNAPSHOT") + .expect("set the existing pinned Petri runner snapshot name"); + assert!( + before + .iter() + .any(|s| s.name.as_deref() == Some(runner.as_str())) + ); + let config = SandboxProviderConfig::from_lookup(Some(credentials), provider_env); + for (mode, expected) in [ + (EnvironmentNetworkMode::AllowAll, NetworkPolicy::AllowAll), + (EnvironmentNetworkMode::Block, NetworkPolicy::Block), + ] { + let root = tempfile::tempdir().unwrap(); + let run_id = RunId::new().to_string(); + let store = Arc::new(MemoryRunStore::new()); + let runtime = RuntimeSpec { + sandbox: config.clone(), + ..Default::default() + }; + let curl = + "curl --noproxy '*' --fail --silent --max-time 5 https://example.com/ -o /dev/null"; + let script = if mode == EnvironmentNetworkMode::Block { + format!("if {curl}; then exit 19; fi") + } else { + curl.to_string() + }; + let workflow = format!( + r#"digraph Network {{ + start [shape=Mdiamond] + probe [shape=parallelogram, script="{script}"] + exit [shape=Msquare] + start -> probe -> exit + }}"# + ); + let graphs = support::admit( + &[ + ("workflow.toml", support::SETTINGS), + ("workflow.fabro", &workflow), + ], + Launch::default(), + &runtime, + ); + let mut request = support::run_request( + &run_id, + root.path(), + graphs, + store.clone(), + runtime, + support::no_questions(Arc::new(support::Silent)), + ); + request.provider = SandboxProviderKind::DAYTONA; + request.network = EnvironmentNetworkSettings { + mode, + allow: Vec::new(), + }; + let outcome = engine::run(request).await; + let mut filter = SandboxFilter::default(); + filter + .labels + .insert("petri.run".to_string(), run_id.clone()); + let observed = provider.list(&filter).await; + let cleanup = prune::prune(PruneRequest { + sandbox: config.clone(), + run_id, + run_dir: root.path().to_path_buf(), + store, + provider: SandboxProviderKind::DAYTONA, + }) + .await; + assert!(cleanup.unwrap().is_clean()); + let observed = observed.unwrap(); + assert_eq!(observed.len(), 1); + assert_eq!(observed[0].snapshot.as_deref(), Some(runner.as_str())); + assert_eq!(observed[0].network.as_ref(), Some(&expected)); + let outcome = outcome.unwrap(); + assert_eq!(outcome.status, RunStatus::Success, "{mode}: {outcome:?}"); + let deadline = Instant::now() + Duration::from_mins(2); + loop { + let remaining = provider.list(&filter).await.unwrap(); + if remaining + .iter() + .all(|sandbox| sandbox.state == SandboxState::Deleted) + { + break; + } + assert!( + Instant::now() < deadline, + "test sandbox deletion did not settle: {remaining:?}" + ); + time::sleep(Duration::from_millis(500)).await; + } + } +} + +#[expect( + clippy::disallowed_methods, + reason = "live tests explicitly pass the operator's provider configuration" +)] +fn provider_env(name: &str) -> Option { + env::var(name).ok() +} diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index eaefa21e1..93ed5be9a 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -19,7 +19,7 @@ use fabro_petri::engine::{Execution, RunRequest}; use fabro_petri::interview::{Approval, FabroInterviewer, QuestionNotice, QuestionSink}; use fabro_petri::runtime::RuntimeSpec; use fabro_types::SandboxProviderKind; -use fabro_types::settings::run::EnvironmentResourcesSettings; +use fabro_types::settings::run::{EnvironmentNetworkSettings, EnvironmentResourcesSettings}; use petri_execution::inspect; use petri_store::{Access, LogId, RunKey, RunStore}; use tokio::time::sleep; @@ -86,6 +86,7 @@ pub(crate) fn run_request( runtime, provider: SandboxProviderKind::LOCAL, resources: EnvironmentResourcesSettings::default(), + network: EnvironmentNetworkSettings::default(), cancel: CancellationToken::new(), controls: RunControls::new(), observers: vec![interviewer.observer()], From 500cd25814700dc9e156fb32a63348a4c99a9da9 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Fri, 2 Oct 2026 10:17:17 -0400 Subject: [PATCH 2/2] Keep network policy off host runs and tidy the network tests The host provider manages no networking and refuses any policy but its default, so sending the resolved AllowAll to local runs failed every non-dry-run local execution. Apply the run's policy only on container backends; dry runs, which always use the host backend, are covered by the same check. Also fold the Docker environment test helper into one that takes a typed network mode, share the probe setup between the live Docker and Daytona network tests, count canary hits per mode, and bind the run environment once in the worker. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 28 +--- .../fabro-server/src/server/petri_runs.rs | 4 +- .../fabro-server/tests/it/scenario/petri.rs | 38 +++-- lib/components/fabro-petri/src/engine.rs | 4 +- lib/components/fabro-petri/tests/network.rs | 150 +++++++++--------- 5 files changed, 110 insertions(+), 114 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index cc491e007..ba2a834ff 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -241,36 +241,16 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { .with_source(source) .with_publisher(publisher) .with_test_gates(test_checkpoint_gates()); + let environment = &worker.run_state.spec.settings.run.environment; let request = RunRequest { run_id: run_id.to_string(), run_dir: worker.run_dir.join("petri"), execution, store, runtime, - provider: worker - .run_state - .spec - .settings - .run - .environment - .provider - .clone(), - resources: worker - .run_state - .spec - .settings - .run - .environment - .resources - .clone(), - network: worker - .run_state - .spec - .settings - .run - .environment - .network - .clone(), + provider: environment.provider.clone(), + resources: environment.resources.clone(), + network: environment.network.clone(), cancel: cancel_token.clone(), controls: controls.clone(), interviewer: Arc::new(petri_interviewer), diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index eafe55510..99ce39575 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -140,9 +140,9 @@ fn settings_layer_toml(state: &AppState) -> Option { /// under `daytona`, and `env`. The rest is the platform's (`cwd`, /// `network`, `lifecycle`, `labels`, `image.dockerfile`, resources the host /// and Docker providers run without, an image the host runs without) and -/// stays with the server's own resolution. The resolved network policy is -/// passed through `RunRequest` at execution; handing it to the frontend would +/// stays with the server's own resolution; handing it to Petri here would /// only warn `ignored.workflow_toml.environments..` on every admit. +/// The resolved network policy reaches Petri through `RunRequest` instead. fn petri_environments(catalog: &MergeMap) -> MergeMap { MergeMap( catalog diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index a9706bea3..03b28cbeb 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -33,6 +33,7 @@ use fabro_server::test_support::{ use fabro_static::EnvVars; use fabro_store::platform_records::{PlatformRecord, PlatformRecordKind, PlatformRecordStore}; use fabro_test::{TwinScenario, TwinScenarios, twin_openai}; +use fabro_types::settings::run::EnvironmentNetworkMode; use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; use tower::ServiceExt; @@ -1018,7 +1019,13 @@ async fn the_server_attaches_to_the_container_petri_created() { .vault_entries([(EnvVars::OPENAI_API_KEY, namespace.clone())]) .build(); let app = test_app_with_scheduler(Arc::clone(&state)); - create_docker_environment(&app, "docker", CATALOG_IMAGE).await; + create_docker_environment( + &app, + "docker", + CATALOG_IMAGE, + EnvironmentNetworkMode::AllowAll, + ) + .await; let version_id = register_version(&app, &[ ("workflow.fabro", COMMAND_DOT), @@ -1228,15 +1235,11 @@ fn get(path: &str) -> Request { const CATALOG_IMAGE: &str = "ghcr.io/lithoscomputer/ubuntu-22.04:slim"; /// A Docker environment in the server's catalog, with the image it runs. -async fn create_docker_environment(app: &axum::Router, id: &str, image: &str) { - create_docker_environment_with_network(app, id, image, "allow_all").await; -} - -async fn create_docker_environment_with_network( +async fn create_docker_environment( app: &axum::Router, id: &str, image: &str, - mode: &str, + mode: EnvironmentNetworkMode, ) { let environment = serde_json::json!({ "id": id, @@ -1340,21 +1343,22 @@ async fn admitted_root_graph(app: &axum::Router, run_id: &str) -> serde_json::Va /// a Docker daemon, the run's container runs that image. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn a_bundle_naming_a_catalog_environment_runs_on_docker_with_its_image() { - assert_catalog_environment_runs("allow_all", "bridge").await; + assert_catalog_environment_runs(EnvironmentNetworkMode::AllowAll, "bridge").await; } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[ignore = "requires Docker; creates and deletes a container"] async fn a_catalog_network_block_reaches_the_created_docker_container() { + // The shared body returns early without Docker; an explicit run must not. assert!(fabro_test::docker_available()); - assert_catalog_environment_runs("block", "none").await; + assert_catalog_environment_runs(EnvironmentNetworkMode::Block, "none").await; } -async fn assert_catalog_environment_runs(mode: &str, network_mode: &str) { +async fn assert_catalog_environment_runs(mode: EnvironmentNetworkMode, network_mode: &str) { let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n"); let state = test_app_state_with_options(settings, 5); let app = test_app_with_scheduler(Arc::clone(&state)); - create_docker_environment_with_network(&app, "docker-small", CATALOG_IMAGE, mode).await; + create_docker_environment(&app, "docker-small", CATALOG_IMAGE, mode).await; let version_id = register_version(&app, &[ ("workflow.fabro", COMMAND_DOT), @@ -1433,7 +1437,9 @@ async fn assert_catalog_environment_runs(mode: &str, network_mode: &str) { } assert!(inspection.status.success(), "{inspection:?}"); assert_eq!( - String::from_utf8(inspection.stdout).unwrap().trim(), + String::from_utf8(inspection.stdout) + .expect("Docker inspect prints UTF-8") + .trim(), network_mode ); } @@ -1445,7 +1451,13 @@ async fn a_bundles_own_environment_table_overrides_the_servers() { let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n"); let state = test_app_state_with_options(settings, 5); let app = test_app_with_scheduler(Arc::clone(&state)); - create_docker_environment(&app, "docker-small", CATALOG_IMAGE).await; + create_docker_environment( + &app, + "docker-small", + CATALOG_IMAGE, + EnvironmentNetworkMode::AllowAll, + ) + .await; let version_id = register_version(&app, &[ ("workflow.fabro", COMMAND_DOT), diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 8820fd276..369149be5 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -207,7 +207,9 @@ pub async fn run(request: RunRequest) -> Result { options.run_key = Some(key.clone()); options.retention = RETENTION; options.sandbox.backend = backend; - if !request.runtime.dry_run { + // The host provider manages no networking and refuses any policy but its + // default, so only container backends receive the run's policy. + if backend != SandboxBackend::Host { options.sandbox.network = network_policy(&request.network); } if backend == SandboxBackend::Daytona { diff --git a/lib/components/fabro-petri/tests/network.rs b/lib/components/fabro-petri/tests/network.rs index 5a424b69b..8d6d3b128 100644 --- a/lib/components/fabro-petri/tests/network.rs +++ b/lib/components/fabro-petri/tests/network.rs @@ -3,18 +3,20 @@ mod support; use std::env; +use std::path::Path; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; use fabro_petri::check::Launch; -use fabro_petri::engine::{self, RunStatus}; -use fabro_petri::providers::SandboxProviderConfig; +use fabro_petri::engine::{self, RunRequest, RunStatus}; +use fabro_petri::providers::{self, DaytonaCredentials, SandboxProviderConfig}; use fabro_petri::prune::{self, PruneRequest}; use fabro_petri::runtime::RuntimeSpec; use fabro_types::settings::run::{EnvironmentNetworkMode, EnvironmentNetworkSettings}; use fabro_types::{RunId, SandboxProviderKind}; use petri_store::MemoryRunStore; +use sandbox_driver::{NetworkPolicy, SandboxFilter, SandboxState, SnapshotFilter}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; use tokio::process::Command; @@ -39,49 +41,24 @@ async fn docker_block_prevents_canary_access_while_allow_all_preserves_it() { })); for (mode, network_mode, expected_hits) in [ (EnvironmentNetworkMode::AllowAll, "bridge", 1), - (EnvironmentNetworkMode::Block, "none", 1), + (EnvironmentNetworkMode::Block, "none", 0), ] { let root = tempfile::tempdir().unwrap(); let run_id = RunId::new().to_string(); let store = Arc::new(MemoryRunStore::new()); - let runtime = RuntimeSpec::default(); let curl = format!( "curl --noproxy '*' --fail --silent --max-time 2 http://host.docker.internal:{port}/canary" ); - let script = if mode == EnvironmentNetworkMode::Block { - format!("if {curl}; then exit 19; fi") - } else { - format!("test $({curl}) = network-canary") - }; - let workflow = format!( - r#"digraph Network {{ - start [shape=Mdiamond] - probe [shape=parallelogram, script="{script}"] - exit [shape=Msquare] - start -> probe -> exit - }}"# - ); - let graphs = support::admit( - &[ - ("workflow.toml", support::SETTINGS), - ("workflow.fabro", &workflow), - ], - Launch::default(), - &runtime, - ); - let mut request = support::run_request( + let request = probe_request( &run_id, root.path(), - graphs, - store.clone(), - runtime, - support::no_questions(Arc::new(support::Silent)), - ); - request.provider = SandboxProviderKind::DOCKER; - request.network = EnvironmentNetworkSettings { + &store, + RuntimeSpec::default(), + SandboxProviderKind::DOCKER, mode, - allow: Vec::new(), - }; + &format!("test $({curl}) = network-canary"), + &curl, + ); let outcome = engine::run(request).await; // Inspect Docker independently of Fabro's settings and driver status. let containers = Command::new("docker") @@ -122,7 +99,7 @@ async fn docker_block_prevents_canary_access_while_allow_all_preserves_it() { String::from_utf8(inspected.stdout).unwrap().trim(), network_mode ); - assert_eq!(hits.load(Ordering::SeqCst), expected_hits, "{mode}"); + assert_eq!(hits.swap(0, Ordering::SeqCst), expected_hits, "{mode}"); } } @@ -130,9 +107,6 @@ async fn docker_block_prevents_canary_access_while_allow_all_preserves_it() { /// runner snapshot and deletes both task-owned sandboxes through Petri. #[fabro_macros::e2e_test(live("DAYTONA_API_KEY"), live("FABRO_TEST_DAYTONA_RUNNER_SNAPSHOT"))] async fn daytona_block_prevents_outbound_https_while_allow_all_preserves_it() { - use fabro_petri::providers::{self, DaytonaCredentials}; - use sandbox_driver::{NetworkPolicy, SandboxFilter, SandboxState, SnapshotFilter}; - let credentials = DaytonaCredentials::from_api_key( fabro_test::require_env("DAYTONA_API_KEY").expect("guard checked the credential"), provider_env, @@ -158,46 +132,21 @@ async fn daytona_block_prevents_outbound_https_while_allow_all_preserves_it() { let root = tempfile::tempdir().unwrap(); let run_id = RunId::new().to_string(); let store = Arc::new(MemoryRunStore::new()); - let runtime = RuntimeSpec { - sandbox: config.clone(), - ..Default::default() - }; let curl = "curl --noproxy '*' --fail --silent --max-time 5 https://example.com/ -o /dev/null"; - let script = if mode == EnvironmentNetworkMode::Block { - format!("if {curl}; then exit 19; fi") - } else { - curl.to_string() - }; - let workflow = format!( - r#"digraph Network {{ - start [shape=Mdiamond] - probe [shape=parallelogram, script="{script}"] - exit [shape=Msquare] - start -> probe -> exit - }}"# - ); - let graphs = support::admit( - &[ - ("workflow.toml", support::SETTINGS), - ("workflow.fabro", &workflow), - ], - Launch::default(), - &runtime, - ); - let mut request = support::run_request( + let request = probe_request( &run_id, root.path(), - graphs, - store.clone(), - runtime, - support::no_questions(Arc::new(support::Silent)), - ); - request.provider = SandboxProviderKind::DAYTONA; - request.network = EnvironmentNetworkSettings { + &store, + RuntimeSpec { + sandbox: config.clone(), + ..Default::default() + }, + SandboxProviderKind::DAYTONA, mode, - allow: Vec::new(), - }; + curl, + curl, + ); let outcome = engine::run(request).await; let mut filter = SandboxFilter::default(); filter @@ -237,6 +186,59 @@ async fn daytona_block_prevents_outbound_https_while_allow_all_preserves_it() { } } +/// A one-stage run whose probe script must succeed under `AllowAll`, and +/// whose `reach` command must fail under `Block`. +#[expect( + clippy::too_many_arguments, + reason = "each live test varies every input" +)] +fn probe_request( + run_id: &str, + root: &Path, + store: &Arc, + runtime: RuntimeSpec, + provider: SandboxProviderKind, + mode: EnvironmentNetworkMode, + allowed_probe: &str, + reach: &str, +) -> RunRequest { + let script = if mode == EnvironmentNetworkMode::Block { + format!("if {reach}; then exit 19; fi") + } else { + allowed_probe.to_string() + }; + let workflow = format!( + r#"digraph Network {{ + start [shape=Mdiamond] + probe [shape=parallelogram, script="{script}"] + exit [shape=Msquare] + start -> probe -> exit + }}"# + ); + let graphs = support::admit( + &[ + ("workflow.toml", support::SETTINGS), + ("workflow.fabro", &workflow), + ], + Launch::default(), + &runtime, + ); + let mut request = support::run_request( + run_id, + root, + graphs, + store.clone(), + runtime, + support::no_questions(Arc::new(support::Silent)), + ); + request.provider = provider; + request.network = EnvironmentNetworkSettings { + mode, + allow: Vec::new(), + }; + request +} + #[expect( clippy::disallowed_methods, reason = "live tests explicitly pass the operator's provider configuration"