mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Enforce resolved network policy for Docker and Daytona runs
This commit is contained in:
parent
a1ee6e5f47
commit
cd5f56a4ed
10 changed files with 433 additions and 52 deletions
31
Cargo.lock
generated
31
Cargo.lock
generated
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -140,8 +140,9 @@ fn settings_layer_toml(state: &AppState) -> Option<String> {
|
|||
/// 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.<id>.<key>` 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.<id>.<key>` on every admit.
|
||||
fn petri_environments(catalog: &MergeMap<EnvironmentLayer>) -> MergeMap<EnvironmentLayer> {
|
||||
MergeMap(
|
||||
catalog
|
||||
|
|
@ -497,6 +498,7 @@ pub(crate) async fn execute(state: Arc<AppState>, 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
|
||||
|
|
|
|||
|
|
@ -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.<id>]` table wins over the server's, key
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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<RunOutcome, RunError> {
|
||||
let backend = if request.runtime.dry_run {
|
||||
|
|
@ -192,6 +207,9 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
|||
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)?;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<u64>,
|
||||
) {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
246
lib/components/fabro-petri/tests/network.rs
Normal file
246
lib/components/fabro-petri/tests/network.rs
Normal file
|
|
@ -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<String> {
|
||||
env::var(name).ok()
|
||||
}
|
||||
|
|
@ -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()],
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue