diff --git a/Cargo.lock b/Cargo.lock index d4b82ee9c..d15f1dee6 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..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,28 +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(), + 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 a7499ab4e..99ce39575 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; 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 @@ -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..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,13 +1235,18 @@ 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) { +async fn create_docker_environment( + app: &axum::Router, + id: &str, + image: &str, + mode: EnvironmentNetworkMode, +) { 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 +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(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(EnvironmentNetworkMode::Block, "none").await; +} + +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(&app, "docker-small", CATALOG_IMAGE).await; + create_docker_environment(&app, "docker-small", CATALOG_IMAGE, mode).await; let version_id = register_version(&app, &[ ("workflow.fabro", COMMAND_DOT), @@ -1381,6 +1405,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 +1435,13 @@ 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) + .expect("Docker inspect prints UTF-8") + .trim(), + network_mode + ); } /// A bundle's own `[environments.]` table wins over the server's, key @@ -1410,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/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..369149be5 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,11 @@ pub async fn run(request: RunRequest) -> Result { options.run_key = Some(key.clone()); options.retention = RETENTION; options.sandbox.backend = backend; + // 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 { 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..8d6d3b128 --- /dev/null +++ b/lib/components/fabro-petri/tests/network.rs @@ -0,0 +1,248 @@ +//! Live Docker coverage through the same engine assembly as a run worker. + +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, 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; +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", 0), + ] { + let root = tempfile::tempdir().unwrap(); + let run_id = RunId::new().to_string(); + let store = Arc::new(MemoryRunStore::new()); + let curl = format!( + "curl --noproxy '*' --fail --silent --max-time 2 http://host.docker.internal:{port}/canary" + ); + let request = probe_request( + &run_id, + root.path(), + &store, + RuntimeSpec::default(), + SandboxProviderKind::DOCKER, + mode, + &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") + .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.swap(0, 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() { + 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 curl = + "curl --noproxy '*' --fail --silent --max-time 5 https://example.com/ -o /dev/null"; + let request = probe_request( + &run_id, + root.path(), + &store, + RuntimeSpec { + sandbox: config.clone(), + ..Default::default() + }, + SandboxProviderKind::DAYTONA, + mode, + curl, + curl, + ); + 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; + } + } +} + +/// 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" +)] +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()],