From 0738cae424b6bf042d4f02822dfba96868c58257 Mon Sep 17 00:00:00 2001 From: jif-oai Date: Tue, 23 Jun 2026 20:31:28 +0100 Subject: [PATCH 1/7] Attribute Linux network requests to their exec --- codex-rs/Cargo.lock | 1 + .../core/src/network_policy_decision_tests.rs | 2 + codex-rs/core/src/tools/network_approval.rs | 117 ++++++++++---- .../core/src/tools/network_approval_tests.rs | 31 ++++ codex-rs/core/src/tools/orchestrator.rs | 13 +- .../src/tools/runtimes/apply_patch_tests.rs | 2 + codex-rs/core/src/tools/runtimes/mod_tests.rs | 1 + .../core/src/tools/runtimes/unified_exec.rs | 4 +- codex-rs/core/src/tools/sandboxing.rs | 10 ++ codex-rs/core/src/tools/sandboxing_tests.rs | 1 + codex-rs/core/tests/suite/network_approval.rs | 134 +++++++++++++++- codex-rs/linux-sandbox/Cargo.toml | 1 + codex-rs/linux-sandbox/src/proxy_routing.rs | 40 ++++- codex-rs/network-proxy/src/attribution.rs | 143 ++++++++++++++++++ .../network-proxy/src/attribution_tests.rs | 63 ++++++++ codex-rs/network-proxy/src/http_proxy.rs | 8 +- codex-rs/network-proxy/src/lib.rs | 3 + codex-rs/network-proxy/src/network_policy.rs | 11 +- codex-rs/network-proxy/src/proxy.rs | 26 ++++ .../src/proxy/execution_scope.rs | 41 +++++ codex-rs/network-proxy/src/runtime.rs | 54 ++++++- codex-rs/network-proxy/src/socks5.rs | 11 +- 22 files changed, 664 insertions(+), 53 deletions(-) create mode 100644 codex-rs/network-proxy/src/attribution.rs create mode 100644 codex-rs/network-proxy/src/attribution_tests.rs create mode 100644 codex-rs/network-proxy/src/proxy/execution_scope.rs diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index bc5953803ed7..49ebe98a88d3 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -3244,6 +3244,7 @@ dependencies = [ "clap", "codex-core", "codex-install-context", + "codex-network-proxy", "codex-process-hardening", "codex-protocol", "codex-sandboxing", diff --git a/codex-rs/core/src/network_policy_decision_tests.rs b/codex-rs/core/src/network_policy_decision_tests.rs index 064004184701..8888bbe9c2d8 100644 --- a/codex-rs/core/src/network_policy_decision_tests.rs +++ b/codex-rs/core/src/network_policy_decision_tests.rs @@ -161,6 +161,7 @@ fn denied_network_policy_message_requires_deny_decision() { method: Some("GET".to_string()), mode: None, protocol: "http".to_string(), + execution_id: None, decision: Some("ask".to_string()), source: Some("decider".to_string()), port: Some(80), @@ -178,6 +179,7 @@ fn denied_network_policy_message_for_denylist_block_is_explicit() { method: Some("GET".to_string()), mode: None, protocol: "http".to_string(), + execution_id: None, decision: Some("deny".to_string()), source: Some("baseline_policy".to_string()), port: Some(80), diff --git a/codex-rs/core/src/tools/network_approval.rs b/codex-rs/core/src/tools/network_approval.rs index 2b63164b4df4..3150d36968c9 100644 --- a/codex-rs/core/src/tools/network_approval.rs +++ b/codex-rs/core/src/tools/network_approval.rs @@ -30,6 +30,7 @@ use codex_protocol::protocol::WarningEvent; use indexmap::IndexMap; use std::collections::HashMap; use std::collections::HashSet; +use std::io; use std::sync::Arc; use tokio::sync::Mutex; use tokio::sync::Notify; @@ -59,6 +60,7 @@ pub(crate) struct DeferredNetworkApproval { registration_id: String, cancellation_token: CancellationToken, finish_outcome: Arc>>, + _execution_proxy: Option, } impl DeferredNetworkApproval { @@ -89,6 +91,7 @@ pub(crate) struct ActiveNetworkApproval { registration_id: Option, mode: NetworkApprovalMode, cancellation_token: CancellationToken, + execution_proxy: NetworkProxy, } impl ActiveNetworkApproval { @@ -100,11 +103,16 @@ impl ActiveNetworkApproval { self.cancellation_token.clone() } + pub(crate) fn execution_proxy(&self) -> &NetworkProxy { + &self.execution_proxy + } + pub(crate) fn into_deferred(self) -> Option { let ActiveNetworkApproval { registration_id, mode, cancellation_token, + execution_proxy, } = self; match (mode, registration_id) { (NetworkApprovalMode::Deferred, Some(registration_id)) => { @@ -112,6 +120,7 @@ impl ActiveNetworkApproval { registration_id, cancellation_token, finish_outcome: Arc::new(OnceCell::new()), + _execution_proxy: Some(execution_proxy), }) } _ => None, @@ -305,9 +314,8 @@ impl NetworkApprovalService { async fn resolve_single_active_call(&self) -> Option> { let calls = self.calls.lock().await; - // Blocked proxy requests are not attributed to a specific tool call. Only pick an owner - // when there is exactly one candidate; with concurrent calls, canceling one would be a guess. - // TODO: Carry blocked-request attribution so concurrent active calls can be handled safely. + // Shared proxy requests can still arrive without an execution ID. Only pick an owner when + // there is exactly one candidate; with concurrent calls, canceling one would be a guess. if calls.active_calls.len() == 1 { return calls.active_calls.values().next().cloned(); } @@ -315,6 +323,18 @@ impl NetworkApprovalService { None } + async fn resolve_active_call_by_execution_id( + &self, + execution_id: &str, + ) -> Option> { + self.calls + .lock() + .await + .active_calls + .get(execution_id) + .cloned() + } + async fn resolve_active_call_attribution(&self) -> ActiveNetworkApprovalAttribution { let calls = self.calls.lock().await; match calls.active_calls.len() { @@ -393,8 +413,12 @@ impl NetworkApprovalService { return; }; - self.record_outcome_for_single_active_call(NetworkApprovalOutcome::DeniedByPolicy(message)) - .await; + let outcome = NetworkApprovalOutcome::DeniedByPolicy(message); + if let Some(execution_id) = blocked.execution_id.as_deref() { + self.record_call_outcome(execution_id, outcome).await; + } else { + self.record_outcome_for_single_active_call(outcome).await; + } } async fn active_turn_context( @@ -431,28 +455,45 @@ impl NetworkApprovalService { NetworkProtocol::Socks5Tcp => NetworkApprovalProtocol::Socks5Tcp, NetworkProtocol::Socks5Udp => NetworkApprovalProtocol::Socks5Udp, }; - let (owner_call, active_environment_id) = - if let Some(environment_id) = request.environment_id.clone() { - let owner_call = match self.resolve_active_call_attribution().await { - ActiveNetworkApprovalAttribution::Single(call) => { - (call.environment_id == environment_id).then_some(call) - } - ActiveNetworkApprovalAttribution::None - | ActiveNetworkApprovalAttribution::Ambiguous => None, - }; - (owner_call, Some(environment_id)) - } else { - match self.resolve_active_call_attribution().await { - ActiveNetworkApprovalAttribution::None => (None, None), - ActiveNetworkApprovalAttribution::Single(call) => { - let environment_id = call.environment_id.clone(); - (Some(call), Some(environment_id)) - } - ActiveNetworkApprovalAttribution::Ambiguous => { - return NetworkDecision::deny(REASON_NOT_ALLOWED); - } + let attributed_call = if let Some(execution_id) = request.execution_id.as_deref() { + let Some(call) = self.resolve_active_call_by_execution_id(execution_id).await else { + return NetworkDecision::deny(REASON_NOT_ALLOWED); + }; + Some(call) + } else { + None + }; + let (owner_call, active_environment_id) = if let Some(call) = attributed_call { + let environment_id = request + .environment_id + .clone() + .unwrap_or_else(|| call.environment_id.clone()); + let owner_call = (call.environment_id == environment_id).then_some(call); + (owner_call, Some(environment_id)) + } else if let Some(environment_id) = request.environment_id.clone() { + let owner_call = match self.resolve_active_call_attribution().await { + ActiveNetworkApprovalAttribution::Single(call) => { + (call.environment_id == environment_id).then_some(call) } + ActiveNetworkApprovalAttribution::None + | ActiveNetworkApprovalAttribution::Ambiguous => None, }; + (owner_call, Some(environment_id)) + } else { + match self.resolve_active_call_attribution().await { + ActiveNetworkApprovalAttribution::None => (None, None), + ActiveNetworkApprovalAttribution::Single(call) => { + let environment_id = call.environment_id.clone(); + (Some(call), Some(environment_id)) + } + ActiveNetworkApprovalAttribution::Ambiguous => { + return NetworkDecision::deny(REASON_NOT_ALLOWED); + } + } + }; + if request.execution_id.is_some() && owner_call.is_none() { + return NetworkDecision::deny(REASON_NOT_ALLOWED); + } let turn_context = Self::active_turn_context(session.as_ref()).await; let Some(environment_id) = active_environment_id.or_else(|| { turn_context @@ -789,19 +830,32 @@ pub(crate) async fn begin_network_approval( turn_id: &str, managed_network_active: bool, spec: Option, -) -> Option { +) -> Result, ToolError> { let NetworkApprovalSpec { network, mode, trigger, command, environment_id, - } = spec?; - if !managed_network_active || network.is_none() { - return None; + } = match spec { + Some(spec) => spec, + None => return Ok(None), + }; + let Some(network) = network else { + return Ok(None); + }; + if !managed_network_active { + return Ok(None); } let registration_id = Uuid::new_v4().to_string(); + let execution_proxy = network + .for_execution(&environment_id, registration_id.clone()) + .map_err(|err| { + ToolError::Codex(codex_protocol::error::CodexErr::Io(io::Error::other( + format!("failed to create execution-scoped network proxy: {err}"), + ))) + })?; let cancellation_token = CancellationToken::new(); session .services @@ -816,11 +870,12 @@ pub(crate) async fn begin_network_approval( ) .await; - Some(ActiveNetworkApproval { + Ok(Some(ActiveNetworkApproval { registration_id: Some(registration_id), mode, cancellation_token, - }) + execution_proxy, + })) } pub(crate) async fn finish_immediate_network_approval( diff --git a/codex-rs/core/src/tools/network_approval_tests.rs b/codex-rs/core/src/tools/network_approval_tests.rs index 7a9ba86cef87..4fdeb50ed7b5 100644 --- a/codex-rs/core/src/tools/network_approval_tests.rs +++ b/codex-rs/core/src/tools/network_approval_tests.rs @@ -285,6 +285,12 @@ fn denied_blocked_request(host: &str) -> BlockedRequest { }) } +fn denied_blocked_request_for_execution(host: &str, execution_id: &str) -> BlockedRequest { + let mut blocked = denied_blocked_request(host); + blocked.execution_id = Some(execution_id.to_string()); + blocked +} + async fn register_call_with_default_shell_trigger( service: &NetworkApprovalService, registration_id: &str, @@ -429,6 +435,7 @@ async fn deferred_finish_reuses_denial_result_after_first_consumer() { registration_id: "registration-1".to_string(), cancellation_token, finish_outcome: Arc::new(OnceCell::new()), + _execution_proxy: None, }; service .record_call_outcome( @@ -481,3 +488,27 @@ async fn record_blocked_request_ignores_ambiguous_unattributed_blocked_requests( assert_eq!(service.take_call_outcome("registration-1").await, None); assert_eq!(service.take_call_outcome("registration-2").await, None); } + +#[tokio::test] +async fn attributed_blocked_request_targets_one_of_multiple_active_calls() { + let service = NetworkApprovalService::default(); + let first = register_call_with_default_shell_trigger(&service, "registration-1").await; + let second = register_call_with_default_shell_trigger(&service, "registration-2").await; + + service + .record_blocked_request(denied_blocked_request_for_execution( + "example.com", + "registration-2", + )) + .await; + + assert!(!first.is_cancelled()); + assert!(second.is_cancelled()); + assert_eq!(service.take_call_outcome("registration-1").await, None); + assert_eq!( + service.take_call_outcome("registration-2").await, + Some(NetworkApprovalOutcome::DeniedByPolicy( + "Network access to \"example.com\" was blocked: domain is not on the allowlist for the current sandbox mode.".to_string() + )) + ); +} diff --git a/codex-rs/core/src/tools/orchestrator.rs b/codex-rs/core/src/tools/orchestrator.rs index aaa0dfe9d14a..8eaa44ce67b0 100644 --- a/codex-rs/core/src/tools/orchestrator.rs +++ b/codex-rs/core/src/tools/orchestrator.rs @@ -68,13 +68,17 @@ impl ToolOrchestrator { where T: ToolRuntime, { - let network_approval = begin_network_approval( + let network_approval = match begin_network_approval( &tool_ctx.session, &tool_ctx.turn.sub_id, managed_network_active, tool.network_approval_spec(req, tool_ctx), ) - .await; + .await + { + Ok(network_approval) => network_approval, + Err(err) => return (Err(err), None), + }; let attempt_tool_ctx = ToolCtx { session: tool_ctx.session.clone(), @@ -98,6 +102,9 @@ impl ToolOrchestrator { network_denial_cancellation_token: network_approval .as_ref() .map(ActiveNetworkApproval::cancellation_token), + network_proxy: network_approval + .as_ref() + .map(ActiveNetworkApproval::execution_proxy), }; let run_result = tool .run(req, &attempt_with_network_approval, &attempt_tool_ctx) @@ -274,6 +281,7 @@ impl ToolOrchestrator { .permissions .windows_sandbox_private_desktop, network_denial_cancellation_token: None, + network_proxy: None, }; let initial_attempt_start = Instant::now(); @@ -456,6 +464,7 @@ impl ToolOrchestrator { .permissions .windows_sandbox_private_desktop, network_denial_cancellation_token: None, + network_proxy: None, }; // Second attempt. diff --git a/codex-rs/core/src/tools/runtimes/apply_patch_tests.rs b/codex-rs/core/src/tools/runtimes/apply_patch_tests.rs index f1c6f43aa342..170068cde09d 100644 --- a/codex-rs/core/src/tools/runtimes/apply_patch_tests.rs +++ b/codex-rs/core/src/tools/runtimes/apply_patch_tests.rs @@ -232,6 +232,7 @@ async fn file_system_sandbox_context_uses_active_attempt() { windows_sandbox_level: WindowsSandboxLevel::RestrictedToken, windows_sandbox_private_desktop: true, network_denial_cancellation_token: None, + network_proxy: None, }; let sandbox = ApplyPatchRuntime::file_system_sandbox_context_for_attempt(&req, &attempt) @@ -300,6 +301,7 @@ async fn no_sandbox_attempt_has_no_file_system_context() { windows_sandbox_level: WindowsSandboxLevel::Disabled, windows_sandbox_private_desktop: false, network_denial_cancellation_token: None, + network_proxy: None, }; assert_eq!( diff --git a/codex-rs/core/src/tools/runtimes/mod_tests.rs b/codex-rs/core/src/tools/runtimes/mod_tests.rs index 9b485dbe58d3..7862a1e3c207 100644 --- a/codex-rs/core/src/tools/runtimes/mod_tests.rs +++ b/codex-rs/core/src/tools/runtimes/mod_tests.rs @@ -118,6 +118,7 @@ async fn explicit_escalation_prepares_exec_without_managed_network() -> anyhow:: windows_sandbox_level: WindowsSandboxLevel::Disabled, windows_sandbox_private_desktop: false, network_denial_cancellation_token: None, + network_proxy: None, }; let exec_request = attempt diff --git a/codex-rs/core/src/tools/runtimes/unified_exec.rs b/codex-rs/core/src/tools/runtimes/unified_exec.rs index 5df42cc979a5..0f99a117a03f 100644 --- a/codex-rs/core/src/tools/runtimes/unified_exec.rs +++ b/codex-rs/core/src/tools/runtimes/unified_exec.rs @@ -320,10 +320,10 @@ impl<'a> ToolRuntime for UnifiedExecRunt req.sandbox_permissions, &file_system_sandbox_policy, ); - let managed_network = managed_network_for_sandbox_permissions( + let managed_network = attempt.network_proxy(managed_network_for_sandbox_permissions( req.network.as_ref(), launch_sandbox_permissions, - ); + )); let env = exec_env_for_sandbox_permissions(&req.env, launch_sandbox_permissions); let (env, managed_network_context) = match managed_network { Some(network) => { diff --git a/codex-rs/core/src/tools/sandboxing.rs b/codex-rs/core/src/tools/sandboxing.rs index 8e9e274f24af..cb61f640731d 100644 --- a/codex-rs/core/src/tools/sandboxing.rs +++ b/codex-rs/core/src/tools/sandboxing.rs @@ -423,9 +423,17 @@ pub(crate) struct SandboxAttempt<'a> { pub windows_sandbox_level: codex_protocol::config_types::WindowsSandboxLevel, pub windows_sandbox_private_desktop: bool, pub network_denial_cancellation_token: Option, + pub(crate) network_proxy: Option<&'a NetworkProxy>, } impl<'a> SandboxAttempt<'a> { + pub(crate) fn network_proxy<'b>( + &'b self, + fallback: Option<&'b NetworkProxy>, + ) -> Option<&'b NetworkProxy> { + fallback.map(|fallback| self.network_proxy.unwrap_or(fallback)) + } + pub fn env_for( &self, command: SandboxCommand, @@ -433,6 +441,7 @@ impl<'a> SandboxAttempt<'a> { network: Option<&NetworkProxy>, environment_id: Option<&str>, ) -> Result { + let network = self.network_proxy(network); let request = self .manager .transform(SandboxTransformRequest { @@ -465,6 +474,7 @@ impl<'a> SandboxAttempt<'a> { network: Option<&NetworkProxy>, environment_id: Option<&str>, ) -> Result { + let network = self.network_proxy(network); let managed_network = command.managed_network.clone(); let exec_server_permissions = effective_permission_profile( self.exec_server_permissions, diff --git a/codex-rs/core/src/tools/sandboxing_tests.rs b/codex-rs/core/src/tools/sandboxing_tests.rs index 28b9de9ce156..e52186f7721a 100644 --- a/codex-rs/core/src/tools/sandboxing_tests.rs +++ b/codex-rs/core/src/tools/sandboxing_tests.rs @@ -227,6 +227,7 @@ fn exec_server_env_keeps_command_native_and_carries_sandbox_context() { windows_sandbox_level: codex_protocol::config_types::WindowsSandboxLevel::Disabled, windows_sandbox_private_desktop: false, network_denial_cancellation_token: None, + network_proxy: None, }; let managed_network = ManagedNetworkSandboxContext { loopback_ports: vec![43123], diff --git a/codex-rs/core/tests/suite/network_approval.rs b/codex-rs/core/tests/suite/network_approval.rs index 6b96c0d6355c..5c958c3f6901 100644 --- a/codex-rs/core/tests/suite/network_approval.rs +++ b/codex-rs/core/tests/suite/network_approval.rs @@ -56,6 +56,118 @@ use tempfile::TempDir; const NETWORK_TEST_HOST: &str = "codex-network-test.invalid"; const NETWORK_TEST_TARGET: &str = "http://codex-network-test.invalid:80"; +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +#[cfg_attr( + not(target_os = "linux"), + ignore = "requires the trusted Linux proxy bridge" +)] +async fn guardian_receives_exact_triggers_for_concurrent_network_requests() -> Result<()> { + skip_if_target_windows!(Ok(()), "uses the POSIX/Python network fixture"); + skip_if_host_windows!(Ok(())); + skip_if_no_network!(Ok(())); + skip_if_sandbox!(Ok(())); + + let server = start_mock_server().await; + let test = managed_network_unified_exec_test(&server).await?; + let barrier_dir = TempDir::new_in(test.cwd.path())?; + let first_marker = barrier_dir.path().join("first"); + let second_marker = barrier_dir.path().join("second"); + let network_command = |marker: &PathBuf, peer_marker: &PathBuf, host: &str| { + format!( + "touch '{}' && while [ ! -e '{}' ]; do sleep 0.01; done && python3 -c \"import urllib.request; urllib.request.build_opener(urllib.request.ProxyHandler()).open('http://{host}', timeout=10).read()\"", + marker.display(), + peer_marker.display(), + ) + }; + let first_command = network_command(&first_marker, &second_marker, "1.1.1.1"); + let second_command = network_command(&second_marker, &first_marker, "8.8.8.8"); + let responses = mount_sse_sequence( + &server, + vec![ + sse(vec![ + ev_response_created("resp-network-concurrent"), + ev_function_call( + "exec-network-first", + "exec_command", + &serde_json::to_string(&network_exec_args(&first_command))?, + ), + ev_function_call( + "exec-network-second", + "exec_command", + &serde_json::to_string(&network_exec_args(&second_command))?, + ), + ev_completed("resp-network-concurrent"), + ]), + sse(vec![ + ev_response_created("resp-network-guardian-1"), + ev_assistant_message("msg-network-guardian-1", r#"{"outcome":"deny"}"#), + ev_completed("resp-network-guardian-1"), + ]), + sse(vec![ + ev_response_created("resp-network-guardian-2"), + ev_assistant_message("msg-network-guardian-2", r#"{"outcome":"deny"}"#), + ev_completed("resp-network-guardian-2"), + ]), + sse(vec![ + ev_response_created("resp-network-done"), + ev_assistant_message("msg-network-done", "done"), + ev_completed("resp-network-done"), + ]), + ], + ) + .await; + + submit_managed_network_turn( + &test, + "run both network requests", + vec![local(test.config.cwd.clone())], + ApprovalsReviewer::AutoReview, + AskForApproval::OnRequest, + ) + .await?; + wait_for_turn_complete(&test).await; + + let mut actual_triggers = responses + .requests() + .into_iter() + .filter(|request| { + request.body_json()["client_metadata"]["x-openai-subagent"].as_str() == Some("guardian") + }) + .map(|request| { + let user_texts = request.message_input_texts("user"); + let action: Value = serde_json::from_str( + user_texts + .iter() + .find(|text| text.contains("\"tool\": \"network_access\"")) + .context("expected network access JSON in Guardian request")? + .trim(), + )?; + Ok(( + action + .pointer("/trigger/callId") + .and_then(Value::as_str) + .context("expected exact trigger call id")? + .to_string(), + action + .pointer("/trigger/command/2") + .and_then(Value::as_str) + .context("expected exact trigger command")? + .to_string(), + )) + }) + .collect::>>()?; + actual_triggers.sort_unstable(); + assert_eq!( + actual_triggers, + vec![ + ("exec-network-first".to_string(), first_command), + ("exec-network-second".to_string(), second_command), + ] + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn approved_network_host_for_one_environment_still_prompts_in_another() -> Result<()> { skip_if_target_windows!(Ok(()), "uses the POSIX/Python network fixture"); @@ -99,6 +211,8 @@ async fn approved_network_host_for_one_environment_still_prompts_in_another() -> &test, "fetch from the local environment", environments.clone(), + ApprovalsReviewer::User, + AskForApproval::OnFailure, ) .await?; let approval = expect_network_approval(&test, LOCAL_ENVIRONMENT_ID).await?; @@ -122,6 +236,8 @@ async fn approved_network_host_for_one_environment_still_prompts_in_another() -> &test, "fetch from the remote environment", environments.clone(), + ApprovalsReviewer::User, + AskForApproval::OnFailure, ) .await?; let approval = expect_network_approval(&test, REMOTE_ENVIRONMENT_ID).await?; @@ -225,12 +341,20 @@ async fn mount_exec_network_turn( } fn network_fetch_args(environment_id: &str) -> Value { + let command = format!( + "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{NETWORK_TEST_HOST}', timeout=2).read().decode(errors='replace'))\"" + ); + let mut args = network_exec_args(&command); + args["environment_id"] = json!(environment_id); + args +} + +fn network_exec_args(command: &str) -> Value { json!({ "shell": "/bin/sh", - "cmd": format!("python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{NETWORK_TEST_HOST}', timeout=2).read().decode(errors='replace'))\""), + "cmd": command, "login": false, "yield_time_ms": 1_000, - "environment_id": environment_id, }) } @@ -238,6 +362,8 @@ async fn submit_managed_network_turn( test: &TestCodex, prompt: &str, environments: Vec, + approvals_reviewer: ApprovalsReviewer, + approval_policy: AskForApproval, ) -> Result<()> { let permission_profile = PermissionProfile::workspace_write_with( &[], @@ -261,8 +387,8 @@ async fn submit_managed_network_turn( additional_context: Default::default(), thread_settings: codex_protocol::protocol::ThreadSettingsOverrides { environments: Some(turn_environment_selections), - approval_policy: Some(AskForApproval::OnRequest), - approvals_reviewer: Some(ApprovalsReviewer::User), + approval_policy: Some(approval_policy), + approvals_reviewer: Some(approvals_reviewer), sandbox_policy: Some(sandbox_policy), permission_profile, collaboration_mode: Some(codex_protocol::config_types::CollaborationMode { diff --git a/codex-rs/linux-sandbox/Cargo.toml b/codex-rs/linux-sandbox/Cargo.toml index fc7937536c05..0f7c68a152a7 100644 --- a/codex-rs/linux-sandbox/Cargo.toml +++ b/codex-rs/linux-sandbox/Cargo.toml @@ -19,6 +19,7 @@ workspace = true [target.'cfg(target_os = "linux")'.dependencies] clap = { workspace = true, features = ["derive"] } codex-install-context = { workspace = true } +codex-network-proxy = { workspace = true } codex-process-hardening = { workspace = true } codex-protocol = { workspace = true } codex-sandboxing = { workspace = true } diff --git a/codex-rs/linux-sandbox/src/proxy_routing.rs b/codex-rs/linux-sandbox/src/proxy_routing.rs index f171d91745d8..32f949cb0fd5 100644 --- a/codex-rs/linux-sandbox/src/proxy_routing.rs +++ b/codex-rs/linux-sandbox/src/proxy_routing.rs @@ -1,3 +1,5 @@ +use codex_network_proxy::PROXY_ATTRIBUTION_TOKEN_ENV_KEY; +use codex_network_proxy::write_attribution_frame; use serde::Deserialize; use serde::Serialize; use std::collections::BTreeMap; @@ -71,7 +73,13 @@ struct ProxyRoutePlan { } pub(crate) fn prepare_host_proxy_route_spec() -> io::Result { - let env: HashMap = std::env::vars().collect(); + let mut env: HashMap = std::env::vars().collect(); + let attribution_token = env.remove(PROXY_ATTRIBUTION_TOKEN_ENV_KEY); + // SAFETY: the sandbox helper is single-threaded here, before it forks bridge workers or + // executes the user command. + unsafe { + std::env::remove_var(PROXY_ATTRIBUTION_TOKEN_ENV_KEY); + } let plan = plan_proxy_routes(&env); if plan.routes.is_empty() { @@ -100,7 +108,11 @@ pub(crate) fn prepare_host_proxy_route_spec() -> io::Result { let mut host_bridge_pids = Vec::with_capacity(socket_by_endpoint.len()); for (endpoint, socket_path) in &socket_by_endpoint { - host_bridge_pids.push(spawn_host_bridge(*endpoint, socket_path)?); + host_bridge_pids.push(spawn_host_bridge( + *endpoint, + socket_path, + attribution_token.as_deref(), + )?); } spawn_proxy_socket_dir_cleanup_worker(socket_dir, host_bridge_pids)?; @@ -438,7 +450,11 @@ fn cleanup_proxy_socket_dir(socket_dir: &Path) -> io::Result<()> { } } -fn spawn_host_bridge(endpoint: SocketAddr, uds_path: &Path) -> io::Result { +fn spawn_host_bridge( + endpoint: SocketAddr, + uds_path: &Path, + attribution_token: Option<&str>, +) -> io::Result { let (read_fd, write_fd) = create_ready_pipe()?; let pid = unsafe { libc::fork() }; if pid < 0 { @@ -452,7 +468,7 @@ fn spawn_host_bridge(endpoint: SocketAddr, uds_path: &Path) -> io::Result io::Result io::Result<()> { +fn run_host_bridge( + endpoint: SocketAddr, + uds_path: &Path, + ready_fd: libc::c_int, + attribution_token: Option<&str>, +) -> io::Result<()> { harden_bridge_process()?; if uds_path.exists() { std::fs::remove_file(uds_path)?; @@ -482,13 +503,20 @@ fn run_host_bridge(endpoint: SocketAddr, uds_path: &Path, ready_fd: libc::c_int) ready_file.write_all(&[HOST_BRIDGE_READY])?; drop(ready_file); + let attribution_token = attribution_token.map(str::to_owned); loop { let (unix_stream, _) = listener.accept()?; + let attribution_token = attribution_token.clone(); std::thread::spawn(move || { - let tcp_stream = match TcpStream::connect(endpoint) { + let mut tcp_stream = match TcpStream::connect(endpoint) { Ok(stream) => stream, Err(_) => return, }; + if let Some(attribution_token) = attribution_token + && write_attribution_frame(&mut tcp_stream, &attribution_token).is_err() + { + return; + } let _ = proxy_bidirectional(tcp_stream, unix_stream); }); } diff --git a/codex-rs/network-proxy/src/attribution.rs b/codex-rs/network-proxy/src/attribution.rs new file mode 100644 index 000000000000..f526460fd622 --- /dev/null +++ b/codex-rs/network-proxy/src/attribution.rs @@ -0,0 +1,143 @@ +use crate::state::NetworkProxyState; +use rama_core::Service; +use rama_core::error::BoxError; +use rama_core::extensions::ExtensionsMut; +use rama_tcp::TcpStream; +use std::io; +use std::io::Write; +use std::sync::Arc; +use std::time::Duration; +use tokio::io::AsyncReadExt; + +/// Internal handoff from the trusted Linux proxy bridge. +#[doc(hidden)] +pub const PROXY_ATTRIBUTION_TOKEN_ENV_KEY: &str = "CODEX_NETWORK_PROXY_ATTRIBUTION"; + +const ATTRIBUTION_FRAME_MAGIC: &[u8; 8] = b"\0CDXPXY1"; +const MAX_ATTRIBUTION_TOKEN_LEN: usize = 128; +const ATTRIBUTION_FRAME_TIMEOUT: Duration = Duration::from_secs(3); + +pub(crate) struct BindConnectionAttribution { + inner: S, + state: Arc, + environment_id: Option, +} + +impl BindConnectionAttribution { + pub(crate) fn new( + inner: S, + state: Arc, + environment_id: Option, + ) -> Self { + Self { + inner, + state, + environment_id, + } + } +} + +impl Service for BindConnectionAttribution +where + S: Service, + S::Error: Into, +{ + type Output = S::Output; + type Error = BoxError; + + async fn serve(&self, mut stream: TcpStream) -> Result { + let state = match read_attribution_token(&mut stream).await? { + Some(token) => self.state.for_execution_token(&token).ok_or_else(|| { + io::Error::new( + io::ErrorKind::PermissionDenied, + "unknown network proxy attribution token", + ) + })?, + None => self.state.as_ref().clone(), + }; + if let Some(expected_environment_id) = self.environment_id.as_deref() + && state + .environment_id() + .is_some_and(|actual| actual != expected_environment_id) + { + return Err(io::Error::new( + io::ErrorKind::PermissionDenied, + "network proxy attribution environment mismatch", + ) + .into()); + } + stream.extensions_mut().insert(Arc::new(state)); + self.inner.serve(stream).await.map_err(Into::into) + } +} + +async fn read_attribution_token(stream: &mut TcpStream) -> Result, BoxError> { + let mut marker = [0_u8; 1]; + let read = stream.stream.peek(&mut marker).await?; + if read == 0 { + return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "empty proxy connection").into()); + } + if marker[0] != ATTRIBUTION_FRAME_MAGIC[0] { + return Ok(None); + } + + let token = tokio::time::timeout(ATTRIBUTION_FRAME_TIMEOUT, async { + let mut magic = [0_u8; ATTRIBUTION_FRAME_MAGIC.len()]; + stream.read_exact(&mut magic).await?; + if &magic != ATTRIBUTION_FRAME_MAGIC { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "invalid network proxy attribution frame", + )); + } + + let token_len = stream.read_u16().await? as usize; + if token_len == 0 || token_len > MAX_ATTRIBUTION_TOKEN_LEN { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "invalid network proxy attribution token length", + )); + } + let mut token = vec![0_u8; token_len]; + stream.read_exact(&mut token).await?; + String::from_utf8(token).map_err(|_| { + io::Error::new( + io::ErrorKind::InvalidData, + "network proxy attribution token is not UTF-8", + ) + }) + }) + .await + .map_err(|_| { + io::Error::new( + io::ErrorKind::TimedOut, + "network proxy attribution frame timed out", + ) + })??; + + Ok(Some(token)) +} + +/// Writes the trusted bridge preface consumed by the shared proxy ingress. +#[doc(hidden)] +pub fn write_attribution_frame(writer: &mut impl Write, token: &str) -> io::Result<()> { + if token.is_empty() || token.len() > MAX_ATTRIBUTION_TOKEN_LEN { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "invalid network proxy attribution token length", + )); + } + let token_len = u16::try_from(token.len()).map_err(|_| { + io::Error::new( + io::ErrorKind::InvalidInput, + "network proxy attribution token is too long", + ) + })?; + writer.write_all(ATTRIBUTION_FRAME_MAGIC)?; + writer.write_all(&token_len.to_be_bytes())?; + writer.write_all(token.as_bytes()) +} + +#[cfg(test)] +#[path = "attribution_tests.rs"] +mod tests; diff --git a/codex-rs/network-proxy/src/attribution_tests.rs b/codex-rs/network-proxy/src/attribution_tests.rs new file mode 100644 index 000000000000..2021f8275963 --- /dev/null +++ b/codex-rs/network-proxy/src/attribution_tests.rs @@ -0,0 +1,63 @@ +use super::BindConnectionAttribution; +use super::write_attribution_frame; +use crate::config::NetworkProxySettings; +use crate::runtime::network_proxy_state_for_policy; +use crate::state::NetworkProxyState; +use pretty_assertions::assert_eq; +use rama_core::Service; +use rama_core::error::BoxError; +use rama_core::extensions::ExtensionsRef; +use rama_core::service::service_fn; +use rama_tcp::TcpStream as RamaTcpStream; +use std::io; +use std::sync::Arc; +use tokio::io::AsyncWriteExt; +use tokio::net::TcpListener; +use tokio::net::TcpStream; + +#[test] +fn attribution_frame_has_bounded_binary_prefix() -> io::Result<()> { + let mut frame = Vec::new(); + write_attribution_frame(&mut frame, "token-1")?; + + assert_eq!(&frame[..8], b"\0CDXPXY1"); + assert_eq!(u16::from_be_bytes([frame[8], frame[9]]), 7); + assert_eq!(&frame[10..], b"token-1"); + Ok(()) +} + +#[tokio::test] +async fn framed_connection_receives_registered_execution_state() -> Result<(), BoxError> { + let state = Arc::new(network_proxy_state_for_policy( + NetworkProxySettings::default(), + )); + state.register_execution("local", "token-1"); + + let listener = TcpListener::bind("127.0.0.1:0").await?; + let addr = listener.local_addr()?; + let client = tokio::spawn(async move { + let mut stream = TcpStream::connect(addr).await?; + let mut frame = Vec::new(); + write_attribution_frame(&mut frame, "token-1")?; + stream.write_all(&frame).await + }); + + let (stream, _) = listener.accept().await?; + let service = BindConnectionAttribution::new( + service_fn(|stream: RamaTcpStream| async move { + let state = stream.extensions().get::>().cloned(); + Ok::<_, io::Error>(state) + }), + state, + Some("local".to_string()), + ); + let actual = service + .serve(RamaTcpStream::new(stream)) + .await? + .expect("connection state"); + client.await??; + + assert_eq!(actual.environment_id(), Some("local")); + assert_eq!(actual.execution_id().as_deref(), Some("token-1")); + Ok(()) +} diff --git a/codex-rs/network-proxy/src/http_proxy.rs b/codex-rs/network-proxy/src/http_proxy.rs index 180ef41c51c5..59cad9c20a5e 100644 --- a/codex-rs/network-proxy/src/http_proxy.rs +++ b/codex-rs/network-proxy/src/http_proxy.rs @@ -1,3 +1,4 @@ +use crate::attribution::BindConnectionAttribution; use crate::config::NetworkMode; use crate::connect_policy::TargetCheckedTcpConnector; use crate::mitm; @@ -39,7 +40,6 @@ use rama_core::error::ErrorExt as _; use rama_core::error::OpaqueError; use rama_core::extensions::ExtensionsMut; use rama_core::extensions::ExtensionsRef; -use rama_core::layer::AddInputExtensionLayer; use rama_core::service::service_fn; use rama_core::stream::Stream; use rama_http::Body; @@ -159,7 +159,11 @@ async fn run_http_proxy_with_listener( info!("HTTP proxy listening on {addr}"); listener - .serve(AddInputExtensionLayer::new(state).into_layer(http_service)) + .serve(BindConnectionAttribution::new( + http_service, + state, + environment_id, + )) .await; Ok(()) } diff --git a/codex-rs/network-proxy/src/lib.rs b/codex-rs/network-proxy/src/lib.rs index 82ab48ba89fc..6bf6d8eb3892 100644 --- a/codex-rs/network-proxy/src/lib.rs +++ b/codex-rs/network-proxy/src/lib.rs @@ -1,5 +1,6 @@ #![deny(clippy::print_stdout, clippy::print_stderr)] +mod attribution; mod certs; mod config; mod connect_policy; @@ -18,6 +19,8 @@ mod socks5; mod state; mod upstream; +pub use attribution::PROXY_ATTRIBUTION_TOKEN_ENV_KEY; +pub use attribution::write_attribution_frame; pub use certs::CUSTOM_CA_ENV_KEYS; pub use certs::is_managed_mitm_ca_trust_bundle_path; pub use config::NetworkDomainPermission; diff --git a/codex-rs/network-proxy/src/network_policy.rs b/codex-rs/network-proxy/src/network_policy.rs index e5792a07c0ee..23239e32b62c 100644 --- a/codex-rs/network-proxy/src/network_policy.rs +++ b/codex-rs/network-proxy/src/network_policy.rs @@ -84,6 +84,7 @@ pub struct NetworkPolicyRequest { pub method: Option, pub command: Option, pub exec_policy_hint: Option, + pub execution_id: Option, } pub struct NetworkPolicyRequestArgs { @@ -118,6 +119,7 @@ impl NetworkPolicyRequest { method, command, exec_policy_hint, + execution_id: None, } } } @@ -300,7 +302,14 @@ pub(crate) async fn evaluate_host_policy( HostBlockDecision::Allowed => (NetworkDecision::Allow, false), HostBlockDecision::Blocked(HostBlockReason::NotAllowed) => { if let Some(decider) = decider { - let decider_decision = map_decider_decision(decider.decide(request.clone()).await); + let mut request = request.clone(); + if request.environment_id.is_none() + && let Some(environment_id) = state.environment_id() + { + request.environment_id = Some(environment_id.to_string()); + } + request.execution_id = state.execution_id(); + let decider_decision = map_decider_decision(decider.decide(request).await); let policy_override = matches!(decider_decision, NetworkDecision::Allow); (decider_decision, policy_override) } else { diff --git a/codex-rs/network-proxy/src/proxy.rs b/codex-rs/network-proxy/src/proxy.rs index a338744b688d..fceb199bfe80 100644 --- a/codex-rs/network-proxy/src/proxy.rs +++ b/codex-rs/network-proxy/src/proxy.rs @@ -1,3 +1,6 @@ +mod execution_scope; + +use crate::attribution::PROXY_ATTRIBUTION_TOKEN_ENV_KEY; use crate::config; use crate::credential_broker::BROKERED_CREDENTIALS_ENV_KEY; use crate::credential_broker::CREDENTIAL_BROKER_ACTIVE_ENV_KEY; @@ -23,6 +26,8 @@ use std::sync::RwLock; use tokio::task::JoinHandle; use tracing::warn; +use self::execution_scope::ExecutionScope; + #[derive(Debug, Clone, Parser)] #[command(name = "codex-network-proxy", about = "Codex network sandbox proxy")] pub struct Args {} @@ -233,6 +238,7 @@ impl NetworkProxyBuilder { reserved_listeners, policy_decider: self.policy_decider, environment_proxies: Arc::new(Mutex::new(HashMap::new())), + execution_scope: None, }) } } @@ -370,6 +376,7 @@ pub struct NetworkProxy { reserved_listeners: Option>, policy_decider: Option>, environment_proxies: Arc>>, + execution_scope: Option>, } impl std::fmt::Debug for NetworkProxy { @@ -424,6 +431,7 @@ pub const PROXY_ENV_KEYS: &[&str] = &[ CREDENTIAL_BROKER_ACTIVE_ENV_KEY, BROKERED_CREDENTIALS_ENV_KEY, ALLOW_LOCAL_BINDING_ENV_KEY, + PROXY_ATTRIBUTION_TOKEN_ENV_KEY, ELECTRON_GET_USE_PROXY_ENV_KEY, NODE_USE_ENV_PROXY_ENV_KEY, "HTTP_PROXY", @@ -697,6 +705,12 @@ impl NetworkProxy { runtime_settings.mitm_ca_trust_bundle.as_ref(), ); self.state.virtualize_child_credentials(&mut env); + if let Some(execution_scope) = self.execution_scope.as_ref() { + env.insert( + PROXY_ATTRIBUTION_TOKEN_ENV_KEY.to_string(), + execution_scope.execution_id.clone(), + ); + } let mut loopback_ports = [ Some(addrs.http_addr), self.socks_enabled.then_some(addrs.socks_addr), @@ -778,6 +792,14 @@ impl NetworkProxy { } fn environment_proxy_addrs(&self, environment_id: &str) -> Result { + if let Some(execution_scope) = self.execution_scope.as_ref() { + anyhow::ensure!( + execution_scope.environment_id == environment_id, + "execution-scoped network proxy belongs to environment `{}`, not `{environment_id}`", + execution_scope.environment_id + ); + } + let mut proxies = self .environment_proxies .lock() @@ -895,6 +917,10 @@ impl NetworkProxy { } pub async fn run(&self) -> Result { + anyhow::ensure!( + self.execution_scope.is_none(), + "execution-scoped network proxy is already running" + ); let current_cfg = self.state.current_cfg().await?; if !current_cfg.network.enabled { warn!("network.enabled is false; skipping proxy listeners"); diff --git a/codex-rs/network-proxy/src/proxy/execution_scope.rs b/codex-rs/network-proxy/src/proxy/execution_scope.rs new file mode 100644 index 000000000000..4b9c01f86892 --- /dev/null +++ b/codex-rs/network-proxy/src/proxy/execution_scope.rs @@ -0,0 +1,41 @@ +use super::*; + +pub(super) struct ExecutionScope { + pub(super) environment_id: String, + pub(super) execution_id: String, + state: Arc, +} + +impl Drop for ExecutionScope { + fn drop(&mut self) { + self.state.unregister_execution(&self.execution_id); + } +} + +impl NetworkProxy { + /// Returns a proxy that attributes trusted bridge connections to one execution. + pub fn for_execution(&self, environment_id: &str, execution_id: String) -> Result { + #[cfg(not(target_os = "linux"))] + { + let _ = (environment_id, execution_id); + return Ok(self.clone()); + } + + #[cfg(target_os = "linux")] + { + anyhow::ensure!( + self.execution_scope.is_none(), + "cannot scope an execution-scoped network proxy" + ); + self.state.register_execution(environment_id, &execution_id); + + let mut proxy = self.clone(); + proxy.execution_scope = Some(Arc::new(ExecutionScope { + environment_id: environment_id.to_string(), + execution_id, + state: Arc::clone(&self.state), + })); + Ok(proxy) + } + } +} diff --git a/codex-rs/network-proxy/src/runtime.rs b/codex-rs/network-proxy/src/runtime.rs index b486ab057a85..f0e0e0dce604 100644 --- a/codex-rs/network-proxy/src/runtime.rs +++ b/codex-rs/network-proxy/src/runtime.rs @@ -33,6 +33,7 @@ use std::net::SocketAddr; use std::path::Path; use std::pin::Pin; use std::sync::Arc; +use std::sync::Mutex; use std::time::Duration; use time::OffsetDateTime; use tokio::net::lookup_host; @@ -96,6 +97,8 @@ pub struct BlockedRequest { pub method: Option, pub mode: Option, pub protocol: String, + #[serde(skip)] + pub execution_id: Option, #[serde(skip_serializing_if = "Option::is_none")] pub decision: Option, #[serde(skip_serializing_if = "Option::is_none")] @@ -137,6 +140,7 @@ impl BlockedRequest { method, mode, protocol, + execution_id: None, decision, source, port, @@ -211,6 +215,9 @@ pub struct NetworkProxyState { blocked_request_observer: Arc>>>, credential_broker: CredentialBroker, audit_metadata: NetworkProxyAuditMetadata, + execution_attributions: Arc>>, + environment_id: Option>, + execution_id: Option>, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -236,6 +243,9 @@ impl Clone for NetworkProxyState { blocked_request_observer: self.blocked_request_observer.clone(), credential_broker: self.credential_broker.clone(), audit_metadata: self.audit_metadata.clone(), + execution_attributions: self.execution_attributions.clone(), + environment_id: self.environment_id.clone(), + execution_id: self.execution_id.clone(), } } } @@ -287,9 +297,49 @@ impl NetworkProxyState { reloader, blocked_request_observer: Arc::new(RwLock::new(blocked_request_observer)), audit_metadata, + execution_attributions: Arc::new(Mutex::new(HashMap::new())), + environment_id: None, + execution_id: None, } } + #[cfg(any(target_os = "linux", test))] + pub(crate) fn register_execution(&self, environment_id: &str, execution_id: &str) { + self.execution_attributions + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(execution_id.to_string(), environment_id.to_string()); + } + + pub(crate) fn unregister_execution(&self, execution_id: &str) { + self.execution_attributions + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(execution_id); + } + + pub(crate) fn for_execution_token(&self, token: &str) -> Option { + let environment_id = self + .execution_attributions + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .get(token)? + .clone(); + Some(Self { + environment_id: Some(environment_id.into()), + execution_id: Some(token.to_string().into()), + ..self.clone() + }) + } + + pub(crate) fn environment_id(&self) -> Option<&str> { + self.environment_id.as_deref() + } + + pub(crate) fn execution_id(&self) -> Option { + self.execution_id.as_deref().map(str::to_string) + } + pub async fn set_blocked_request_observer( &self, blocked_request_observer: Option>, @@ -460,8 +510,9 @@ impl NetworkProxyState { } } - pub async fn record_blocked(&self, entry: BlockedRequest) -> Result<()> { + pub async fn record_blocked(&self, mut entry: BlockedRequest) -> Result<()> { self.reload_if_needed().await?; + entry.execution_id = self.execution_id(); let blocked_for_observer = entry.clone(); let blocked_request_observer = self.blocked_request_observer.read().await.clone(); let violation_line = blocked_request_violation_log_line(&entry); @@ -1306,6 +1357,7 @@ mod tests { method: Some("GET".to_string()), mode: Some(NetworkMode::Full), protocol: "http".to_string(), + execution_id: None, decision: Some("ask".to_string()), source: Some("decider".to_string()), port: Some(80), diff --git a/codex-rs/network-proxy/src/socks5.rs b/codex-rs/network-proxy/src/socks5.rs index 956187704309..d1c193198f59 100644 --- a/codex-rs/network-proxy/src/socks5.rs +++ b/codex-rs/network-proxy/src/socks5.rs @@ -1,3 +1,4 @@ +use crate::attribution::BindConnectionAttribution; use crate::config::NetworkMode; use crate::connect_policy::TargetCheckedTcpConnector; use crate::mitm; @@ -23,13 +24,11 @@ use crate::state::BlockedRequestArgs; use crate::state::NetworkProxyState; use anyhow::Context as _; use anyhow::Result; -use rama_core::Layer; use rama_core::Service; use rama_core::error::BoxError; use rama_core::extensions::Extensions; use rama_core::extensions::ExtensionsMut; use rama_core::extensions::ExtensionsRef; -use rama_core::layer::AddInputExtensionLayer; use rama_core::service::service_fn; use rama_net::address::HostWithPort; use rama_net::client::EstablishedClientConnection; @@ -165,11 +164,15 @@ async fn run_socks5_with_listener( })); let socks_acceptor = base.with_udp_associator(udp_relay); listener - .serve(AddInputExtensionLayer::new(state).into_layer(socks_acceptor)) + .serve(BindConnectionAttribution::new( + socks_acceptor, + state, + environment_id, + )) .await; } else { listener - .serve(AddInputExtensionLayer::new(state).into_layer(base)) + .serve(BindConnectionAttribution::new(base, state, environment_id)) .await; } Ok(()) From 53cf25d7960f398eb825384cc5e811de01aa6feb Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Tue, 23 Jun 2026 13:56:19 -0700 Subject: [PATCH 2/7] fix(network-proxy): attribute audit events to executions Co-authored-by: Codex noreply@openai.com --- codex-rs/network-proxy/src/network_policy.rs | 52 ++++++++++++++++++- .../src/proxy/execution_scope.rs | 2 +- 2 files changed, 52 insertions(+), 2 deletions(-) diff --git a/codex-rs/network-proxy/src/network_policy.rs b/codex-rs/network-proxy/src/network_policy.rs index 23239e32b62c..46e19ed8ed36 100644 --- a/codex-rs/network-proxy/src/network_policy.rs +++ b/codex-rs/network-proxy/src/network_policy.rs @@ -201,6 +201,7 @@ fn emit_non_domain_policy_decision_audit_event( args: BlockDecisionAuditEventArgs<'_>, decision: &'static str, ) { + let execution_id = state.execution_id(); emit_policy_audit_event( state, PolicyAuditEventArgs { @@ -213,6 +214,7 @@ fn emit_non_domain_policy_decision_audit_event( server_port: args.server_port, method: args.method, client_addr: args.client_addr, + execution_id: execution_id.as_deref(), policy_override: false, }, ); @@ -228,6 +230,7 @@ struct PolicyAuditEventArgs<'a> { server_port: u16, method: Option<&'a str>, client_addr: Option<&'a str>, + execution_id: Option<&'a str>, policy_override: bool, } @@ -256,6 +259,7 @@ fn emit_policy_audit_event(state: &NetworkProxyState, args: PolicyAuditEventArgs server.port = args.server_port, http.request.method = args.method.unwrap_or(DEFAULT_METHOD), client.address = args.client_addr.unwrap_or(DEFAULT_CLIENT_ADDRESS), + execution.id = args.execution_id, network.policy.override = args.policy_override, ); } @@ -297,6 +301,7 @@ pub(crate) async fn evaluate_host_policy( decider: Option<&Arc>, request: &NetworkPolicyRequest, ) -> Result { + let execution_id = state.execution_id(); let host_decision = state.host_blocked(&request.host, request.port).await?; let (decision, policy_override) = match host_decision { HostBlockDecision::Allowed => (NetworkDecision::Allow, false), @@ -308,7 +313,7 @@ pub(crate) async fn evaluate_host_policy( { request.environment_id = Some(environment_id.to_string()); } - request.execution_id = state.execution_id(); + request.execution_id = execution_id.clone(); let decider_decision = map_decider_decision(decider.decide(request).await); let policy_override = matches!(decider_decision, NetworkDecision::Allow); (decider_decision, policy_override) @@ -364,6 +369,7 @@ pub(crate) async fn evaluate_host_policy( server_port: request.port, method: request.method.as_deref(), client_addr: request.client_addr.as_deref(), + execution_id: execution_id.as_deref(), policy_override, }, ); @@ -688,6 +694,45 @@ mod tests { ); } + #[tokio::test(flavor = "current_thread")] + async fn evaluate_host_policy_emits_execution_id_for_baseline_allow() { + let state = network_proxy_state_for_policy({ + let mut network = NetworkProxySettings::default(); + network.set_allowed_domains(vec!["example.com".to_string()]); + network + }); + state.register_execution("local", "execution-baseline-allow"); + let state = state + .for_execution_token("execution-baseline-allow") + .expect("expected registered execution"); + let request = NetworkPolicyRequest::new(NetworkPolicyRequestArgs { + protocol: NetworkProtocol::Http, + host: "example.com".to_string(), + port: 80, + environment_id: None, + client_addr: None, + method: None, + command: None, + exec_policy_hint: None, + }); + + let (decision, events) = capture_events(|| async { + evaluate_host_policy(&state, /*decider*/ None, &request) + .await + .unwrap() + }) + .await; + assert_eq!(decision, NetworkDecision::Allow); + + let event = find_event_by_name(&events, POLICY_DECISION_EVENT_NAME) + .expect("expected policy decision audit event"); + assert_eq!(event.field("network.policy.decision"), Some("allow")); + assert_eq!( + event.field("execution.id"), + Some("execution-baseline-allow") + ); + } + #[tokio::test(flavor = "current_thread")] async fn evaluate_host_policy_emits_domain_event_for_baseline_deny() { let state = network_proxy_state_for_policy({ @@ -696,6 +741,10 @@ mod tests { network.set_denied_domains(vec!["blocked.com".to_string()]); network }); + state.register_execution("local", "execution-baseline-deny"); + let state = state + .for_execution_token("execution-baseline-deny") + .expect("expected registered execution"); let request = NetworkPolicyRequest::new(NetworkPolicyRequestArgs { protocol: NetworkProtocol::Http, host: "blocked.com".to_string(), @@ -733,6 +782,7 @@ mod tests { assert_eq!(event.field("network.policy.override"), Some("false")); assert_eq!(event.field("http.request.method"), Some("GET")); assert_eq!(event.field("client.address"), Some("127.0.0.1:1234")); + assert_eq!(event.field("execution.id"), Some("execution-baseline-deny")); } #[tokio::test(flavor = "current_thread")] diff --git a/codex-rs/network-proxy/src/proxy/execution_scope.rs b/codex-rs/network-proxy/src/proxy/execution_scope.rs index 4b9c01f86892..0c7e6b570e6a 100644 --- a/codex-rs/network-proxy/src/proxy/execution_scope.rs +++ b/codex-rs/network-proxy/src/proxy/execution_scope.rs @@ -18,7 +18,7 @@ impl NetworkProxy { #[cfg(not(target_os = "linux"))] { let _ = (environment_id, execution_id); - return Ok(self.clone()); + Ok(self.clone()) } #[cfg(target_os = "linux")] From f216c8b59a638c16185821dbe8b669e51b9cf25a Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Tue, 23 Jun 2026 14:25:59 -0700 Subject: [PATCH 3/7] fix(core): update network approval policy variant Co-authored-by: Codex noreply@openai.com --- codex-rs/core/tests/suite/network_approval.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/codex-rs/core/tests/suite/network_approval.rs b/codex-rs/core/tests/suite/network_approval.rs index 5c958c3f6901..2c80868c9102 100644 --- a/codex-rs/core/tests/suite/network_approval.rs +++ b/codex-rs/core/tests/suite/network_approval.rs @@ -212,7 +212,7 @@ async fn approved_network_host_for_one_environment_still_prompts_in_another() -> "fetch from the local environment", environments.clone(), ApprovalsReviewer::User, - AskForApproval::OnFailure, + AskForApproval::UnlessTrusted, ) .await?; let approval = expect_network_approval(&test, LOCAL_ENVIRONMENT_ID).await?; @@ -237,7 +237,7 @@ async fn approved_network_host_for_one_environment_still_prompts_in_another() -> "fetch from the remote environment", environments.clone(), ApprovalsReviewer::User, - AskForApproval::OnFailure, + AskForApproval::UnlessTrusted, ) .await?; let approval = expect_network_approval(&test, REMOTE_ENVIRONMENT_ID).await?; From 820de55627d23786205bc2509c89e524ee55b1a3 Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Tue, 23 Jun 2026 20:05:43 -0700 Subject: [PATCH 4/7] fix: address network attribution review comments Co-authored-by: Codex --- codex-rs/core/src/tools/network_approval.rs | 97 ++++++++----- codex-rs/core/tests/suite/network_approval.rs | 130 +++++++++++++----- codex-rs/linux-sandbox/src/proxy_routing.rs | 45 +++++- 3 files changed, 199 insertions(+), 73 deletions(-) diff --git a/codex-rs/core/src/tools/network_approval.rs b/codex-rs/core/src/tools/network_approval.rs index 3150d36968c9..bb2d723388b7 100644 --- a/codex-rs/core/src/tools/network_approval.rs +++ b/codex-rs/core/src/tools/network_approval.rs @@ -250,6 +250,11 @@ enum ActiveNetworkApprovalAttribution { Ambiguous, } +struct NetworkRequestAttribution { + owner_call: Option>, + environment_id: Option, +} + #[derive(Default)] struct NetworkApprovalCallState { active_calls: IndexMap>, @@ -347,6 +352,54 @@ impl NetworkApprovalService { } } + async fn resolve_request_attribution( + &self, + request: &NetworkPolicyRequest, + ) -> Option { + if let Some(execution_id) = request.execution_id.as_deref() { + let call = self + .resolve_active_call_by_execution_id(execution_id) + .await?; + let environment_id = request + .environment_id + .clone() + .unwrap_or_else(|| call.environment_id.clone()); + return (call.environment_id == environment_id).then_some(NetworkRequestAttribution { + owner_call: Some(call), + environment_id: Some(environment_id), + }); + } + + if let Some(environment_id) = request.environment_id.clone() { + let owner_call = match self.resolve_active_call_attribution().await { + ActiveNetworkApprovalAttribution::Single(call) => { + (call.environment_id == environment_id).then_some(call) + } + ActiveNetworkApprovalAttribution::None + | ActiveNetworkApprovalAttribution::Ambiguous => None, + }; + return Some(NetworkRequestAttribution { + owner_call, + environment_id: Some(environment_id), + }); + } + + match self.resolve_active_call_attribution().await { + ActiveNetworkApprovalAttribution::None => Some(NetworkRequestAttribution { + owner_call: None, + environment_id: None, + }), + ActiveNetworkApprovalAttribution::Single(call) => { + let environment_id = call.environment_id.clone(); + Some(NetworkRequestAttribution { + owner_call: Some(call), + environment_id: Some(environment_id), + }) + } + ActiveNetworkApprovalAttribution::Ambiguous => None, + } + } + async fn get_or_create_pending_approval( &self, key: HostApprovalKey, @@ -455,45 +508,13 @@ impl NetworkApprovalService { NetworkProtocol::Socks5Tcp => NetworkApprovalProtocol::Socks5Tcp, NetworkProtocol::Socks5Udp => NetworkApprovalProtocol::Socks5Udp, }; - let attributed_call = if let Some(execution_id) = request.execution_id.as_deref() { - let Some(call) = self.resolve_active_call_by_execution_id(execution_id).await else { - return NetworkDecision::deny(REASON_NOT_ALLOWED); - }; - Some(call) - } else { - None - }; - let (owner_call, active_environment_id) = if let Some(call) = attributed_call { - let environment_id = request - .environment_id - .clone() - .unwrap_or_else(|| call.environment_id.clone()); - let owner_call = (call.environment_id == environment_id).then_some(call); - (owner_call, Some(environment_id)) - } else if let Some(environment_id) = request.environment_id.clone() { - let owner_call = match self.resolve_active_call_attribution().await { - ActiveNetworkApprovalAttribution::Single(call) => { - (call.environment_id == environment_id).then_some(call) - } - ActiveNetworkApprovalAttribution::None - | ActiveNetworkApprovalAttribution::Ambiguous => None, - }; - (owner_call, Some(environment_id)) - } else { - match self.resolve_active_call_attribution().await { - ActiveNetworkApprovalAttribution::None => (None, None), - ActiveNetworkApprovalAttribution::Single(call) => { - let environment_id = call.environment_id.clone(); - (Some(call), Some(environment_id)) - } - ActiveNetworkApprovalAttribution::Ambiguous => { - return NetworkDecision::deny(REASON_NOT_ALLOWED); - } - } - }; - if request.execution_id.is_some() && owner_call.is_none() { + let Some(NetworkRequestAttribution { + owner_call, + environment_id: active_environment_id, + }) = self.resolve_request_attribution(&request).await + else { return NetworkDecision::deny(REASON_NOT_ALLOWED); - } + }; let turn_context = Self::active_turn_context(session.as_ref()).await; let Some(environment_id) = active_environment_id.or_else(|| { turn_context diff --git a/codex-rs/core/tests/suite/network_approval.rs b/codex-rs/core/tests/suite/network_approval.rs index 2c80868c9102..fb936512e8fe 100644 --- a/codex-rs/core/tests/suite/network_approval.rs +++ b/codex-rs/core/tests/suite/network_approval.rs @@ -127,35 +127,7 @@ async fn guardian_receives_exact_triggers_for_concurrent_network_requests() -> R .await?; wait_for_turn_complete(&test).await; - let mut actual_triggers = responses - .requests() - .into_iter() - .filter(|request| { - request.body_json()["client_metadata"]["x-openai-subagent"].as_str() == Some("guardian") - }) - .map(|request| { - let user_texts = request.message_input_texts("user"); - let action: Value = serde_json::from_str( - user_texts - .iter() - .find(|text| text.contains("\"tool\": \"network_access\"")) - .context("expected network access JSON in Guardian request")? - .trim(), - )?; - Ok(( - action - .pointer("/trigger/callId") - .and_then(Value::as_str) - .context("expected exact trigger call id")? - .to_string(), - action - .pointer("/trigger/command/2") - .and_then(Value::as_str) - .context("expected exact trigger command")? - .to_string(), - )) - }) - .collect::>>()?; + let mut actual_triggers = guardian_network_triggers(&responses)?; actual_triggers.sort_unstable(); assert_eq!( actual_triggers, @@ -168,6 +140,64 @@ async fn guardian_receives_exact_triggers_for_concurrent_network_requests() -> R Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +#[cfg_attr( + not(target_os = "linux"), + ignore = "requires the trusted Linux proxy bridge" +)] +async fn guardian_receives_exact_trigger_for_single_network_request() -> Result<()> { + skip_if_target_windows!(Ok(()), "uses the POSIX/Python network fixture"); + skip_if_host_windows!(Ok(())); + skip_if_no_network!(Ok(())); + skip_if_sandbox!(Ok(())); + + let server = start_mock_server().await; + let test = managed_network_unified_exec_test(&server).await?; + let command = network_python_fetch_command("1.1.1.1", /*timeout_secs*/ 10); + let responses = mount_sse_sequence( + &server, + vec![ + sse(vec![ + ev_response_created("resp-network-single"), + ev_function_call( + "exec-network-single", + "exec_command", + &serde_json::to_string(&network_exec_args(&command))?, + ), + ev_completed("resp-network-single"), + ]), + sse(vec![ + ev_response_created("resp-network-guardian"), + ev_assistant_message("msg-network-guardian", r#"{"outcome":"deny"}"#), + ev_completed("resp-network-guardian"), + ]), + sse(vec![ + ev_response_created("resp-network-done"), + ev_assistant_message("msg-network-done", "done"), + ev_completed("resp-network-done"), + ]), + ], + ) + .await; + + submit_managed_network_turn( + &test, + "run one network request", + vec![local(test.config.cwd.clone())], + ApprovalsReviewer::AutoReview, + AskForApproval::OnRequest, + ) + .await?; + wait_for_turn_complete(&test).await; + + assert_eq!( + guardian_network_triggers(&responses)?, + vec![("exec-network-single".to_string(), command)] + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn approved_network_host_for_one_environment_still_prompts_in_another() -> Result<()> { skip_if_target_windows!(Ok(()), "uses the POSIX/Python network fixture"); @@ -341,14 +371,18 @@ async fn mount_exec_network_turn( } fn network_fetch_args(environment_id: &str) -> Value { - let command = format!( - "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{NETWORK_TEST_HOST}', timeout=2).read().decode(errors='replace'))\"" - ); + let command = network_python_fetch_command(NETWORK_TEST_HOST, /*timeout_secs*/ 2); let mut args = network_exec_args(&command); args["environment_id"] = json!(environment_id); args } +fn network_python_fetch_command(host: &str, timeout_secs: u64) -> String { + format!( + "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{host}', timeout={timeout_secs}).read().decode(errors='replace'))\"" + ) +} + fn network_exec_args(command: &str) -> Value { json!({ "shell": "/bin/sh", @@ -407,6 +441,38 @@ async fn submit_managed_network_turn( Ok(()) } +fn guardian_network_triggers(responses: &ResponseMock) -> Result> { + responses + .requests() + .into_iter() + .filter(|request| { + request.body_json()["client_metadata"]["x-openai-subagent"].as_str() == Some("guardian") + }) + .map(|request| { + let user_texts = request.message_input_texts("user"); + let action: Value = serde_json::from_str( + user_texts + .iter() + .find(|text| text.contains("\"tool\": \"network_access\"")) + .context("expected network access JSON in Guardian request")? + .trim(), + )?; + Ok(( + action + .pointer("/trigger/callId") + .and_then(Value::as_str) + .context("expected exact trigger call id")? + .to_string(), + action + .pointer("/trigger/command/2") + .and_then(Value::as_str) + .context("expected exact trigger command")? + .to_string(), + )) + }) + .collect() +} + async fn expect_network_approval( test: &TestCodex, expected_environment_id: &str, diff --git a/codex-rs/linux-sandbox/src/proxy_routing.rs b/codex-rs/linux-sandbox/src/proxy_routing.rs index 32f949cb0fd5..707cbaafbe28 100644 --- a/codex-rs/linux-sandbox/src/proxy_routing.rs +++ b/codex-rs/linux-sandbox/src/proxy_routing.rs @@ -73,14 +73,12 @@ struct ProxyRoutePlan { } pub(crate) fn prepare_host_proxy_route_spec() -> io::Result { - let mut env: HashMap = std::env::vars().collect(); - let attribution_token = env.remove(PROXY_ATTRIBUTION_TOKEN_ENV_KEY); + let (attribution_token, plan) = extract_attribution_token_and_plan(std::env::vars().collect()); // SAFETY: the sandbox helper is single-threaded here, before it forks bridge workers or // executes the user command. unsafe { std::env::remove_var(PROXY_ATTRIBUTION_TOKEN_ENV_KEY); } - let plan = plan_proxy_routes(&env); if plan.routes.is_empty() { let message = if plan.has_proxy_config { @@ -133,6 +131,14 @@ pub(crate) fn prepare_host_proxy_route_spec() -> io::Result { serde_json::to_string(&ProxyRouteSpec { routes }).map_err(io::Error::other) } +fn extract_attribution_token_and_plan( + mut env: HashMap, +) -> (Option, ProxyRoutePlan) { + let attribution_token = env.remove(PROXY_ATTRIBUTION_TOKEN_ENV_KEY); + let plan = plan_proxy_routes(&env); + (attribution_token, plan) +} + pub(crate) fn activate_proxy_routes_in_netns(serialized_spec: &str) -> io::Result<()> { let spec: ProxyRouteSpec = serde_json::from_str(serialized_spec).map_err(io::Error::other)?; @@ -515,6 +521,8 @@ fn run_host_bridge( if let Some(attribution_token) = attribution_token && write_attribution_frame(&mut tcp_stream, &attribution_token).is_err() { + // The shared ingress must reject unauthenticated connections; do not forward + // application bytes if this bridge cannot prove the exec attribution first. return; } let _ = proxy_bidirectional(tcp_stream, unix_stream); @@ -701,12 +709,14 @@ fn close_fd(fd: libc::c_int) -> io::Result<()> { #[cfg(test)] mod tests { + use super::PROXY_ATTRIBUTION_TOKEN_ENV_KEY; use super::PROXY_SOCKET_DIR_PREFIX; use super::ProxyRouteEntry; use super::ProxyRouteSpec; use super::cleanup_proxy_socket_dir; use super::cleanup_stale_proxy_socket_dirs_in; use super::default_proxy_port; + use super::extract_attribution_token_and_plan; use super::is_proxy_env_key; use super::parse_loopback_proxy_endpoint; use super::parse_proxy_socket_dir_owner_pid; @@ -771,6 +781,35 @@ mod tests { ); } + #[test] + fn attribution_token_is_extracted_before_proxy_route_planning() { + let mut env = HashMap::new(); + env.insert( + "HTTP_PROXY".to_string(), + "http://127.0.0.1:43128".to_string(), + ); + env.insert( + PROXY_ATTRIBUTION_TOKEN_ENV_KEY.to_string(), + "exec-token".to_string(), + ); + + let (attribution_token, plan) = extract_attribution_token_and_plan(env); + + assert_eq!(attribution_token.as_deref(), Some("exec-token")); + assert_eq!( + plan, + super::ProxyRoutePlan { + routes: vec![super::PlannedProxyRoute { + env_key: "HTTP_PROXY".to_string(), + endpoint: "127.0.0.1:43128" + .parse::() + .expect("valid socket"), + }], + has_proxy_config: true, + } + ); + } + #[test] fn rewrites_proxy_url_to_local_loopback_port() { let rewritten = From 4b298f3dc64d9121f16bd761d97b432c672de52f Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Tue, 23 Jun 2026 20:09:45 -0700 Subject: [PATCH 5/7] test: inline single network request command Co-authored-by: Codex --- codex-rs/core/tests/suite/network_approval.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/codex-rs/core/tests/suite/network_approval.rs b/codex-rs/core/tests/suite/network_approval.rs index fb936512e8fe..6f087ac53044 100644 --- a/codex-rs/core/tests/suite/network_approval.rs +++ b/codex-rs/core/tests/suite/network_approval.rs @@ -153,7 +153,7 @@ async fn guardian_receives_exact_trigger_for_single_network_request() -> Result< let server = start_mock_server().await; let test = managed_network_unified_exec_test(&server).await?; - let command = network_python_fetch_command("1.1.1.1", /*timeout_secs*/ 10); + let command = "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://1.1.1.1', timeout=10).read().decode(errors='replace'))\"".to_string(); let responses = mount_sse_sequence( &server, vec![ @@ -371,18 +371,14 @@ async fn mount_exec_network_turn( } fn network_fetch_args(environment_id: &str) -> Value { - let command = network_python_fetch_command(NETWORK_TEST_HOST, /*timeout_secs*/ 2); + let command = format!( + "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{NETWORK_TEST_HOST}', timeout=2).read().decode(errors='replace'))\"" + ); let mut args = network_exec_args(&command); args["environment_id"] = json!(environment_id); args } -fn network_python_fetch_command(host: &str, timeout_secs: u64) -> String { - format!( - "python3 -c \"import urllib.request; opener = urllib.request.build_opener(urllib.request.ProxyHandler()); print('OK:' + opener.open('http://{host}', timeout={timeout_secs}).read().decode(errors='replace'))\"" - ) -} - fn network_exec_args(command: &str) -> Value { json!({ "shell": "/bin/sh", From 10f7126ff25890403c91f4afa079f4ebf4904d34 Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Wed, 24 Jun 2026 14:44:07 -0700 Subject: [PATCH 6/7] fix: separate proxy attribution token from execution id Co-authored-by: Codex --- codex-rs/core/src/tools/network_approval.rs | 21 ++++++--- codex-rs/core/src/tools/orchestrator.rs | 1 + .../network-proxy/src/attribution_tests.rs | 4 +- codex-rs/network-proxy/src/network_policy.rs | 10 +++-- codex-rs/network-proxy/src/proxy.rs | 2 +- .../src/proxy/execution_scope.rs | 45 +++++++++---------- codex-rs/network-proxy/src/runtime.rs | 34 ++++++++++---- 7 files changed, 70 insertions(+), 47 deletions(-) diff --git a/codex-rs/core/src/tools/network_approval.rs b/codex-rs/core/src/tools/network_approval.rs index bb2d723388b7..747ae80a6f12 100644 --- a/codex-rs/core/src/tools/network_approval.rs +++ b/codex-rs/core/src/tools/network_approval.rs @@ -27,6 +27,7 @@ use codex_protocol::protocol::Event; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::ReviewDecision; use codex_protocol::protocol::WarningEvent; +use codex_sandboxing::SandboxType; use indexmap::IndexMap; use std::collections::HashMap; use std::collections::HashSet; @@ -850,6 +851,7 @@ pub(crate) async fn begin_network_approval( session: &Session, turn_id: &str, managed_network_active: bool, + selected_sandbox: SandboxType, spec: Option, ) -> Result, ToolError> { let NetworkApprovalSpec { @@ -870,13 +872,18 @@ pub(crate) async fn begin_network_approval( } let registration_id = Uuid::new_v4().to_string(); - let execution_proxy = network - .for_execution(&environment_id, registration_id.clone()) - .map_err(|err| { - ToolError::Codex(codex_protocol::error::CodexErr::Io(io::Error::other( - format!("failed to create execution-scoped network proxy: {err}"), - ))) - })?; + let execution_proxy = if selected_sandbox == SandboxType::LinuxSeccomp { + let attribution_token = Uuid::new_v4().to_string(); + network + .for_execution(&environment_id, ®istration_id, attribution_token) + .map_err(|err| { + ToolError::Codex(codex_protocol::error::CodexErr::Io(io::Error::other( + format!("failed to create execution-scoped network proxy: {err}"), + ))) + })? + } else { + network.clone() + }; let cancellation_token = CancellationToken::new(); session .services diff --git a/codex-rs/core/src/tools/orchestrator.rs b/codex-rs/core/src/tools/orchestrator.rs index 8eaa44ce67b0..fc0a9c306dfe 100644 --- a/codex-rs/core/src/tools/orchestrator.rs +++ b/codex-rs/core/src/tools/orchestrator.rs @@ -72,6 +72,7 @@ impl ToolOrchestrator { &tool_ctx.session, &tool_ctx.turn.sub_id, managed_network_active, + attempt.sandbox, tool.network_approval_spec(req, tool_ctx), ) .await diff --git a/codex-rs/network-proxy/src/attribution_tests.rs b/codex-rs/network-proxy/src/attribution_tests.rs index 2021f8275963..e28110263c0b 100644 --- a/codex-rs/network-proxy/src/attribution_tests.rs +++ b/codex-rs/network-proxy/src/attribution_tests.rs @@ -31,7 +31,7 @@ async fn framed_connection_receives_registered_execution_state() -> Result<(), B let state = Arc::new(network_proxy_state_for_policy( NetworkProxySettings::default(), )); - state.register_execution("local", "token-1"); + state.register_execution("token-1", "local", "execution-1"); let listener = TcpListener::bind("127.0.0.1:0").await?; let addr = listener.local_addr()?; @@ -58,6 +58,6 @@ async fn framed_connection_receives_registered_execution_state() -> Result<(), B client.await??; assert_eq!(actual.environment_id(), Some("local")); - assert_eq!(actual.execution_id().as_deref(), Some("token-1")); + assert_eq!(actual.execution_id().as_deref(), Some("execution-1")); Ok(()) } diff --git a/codex-rs/network-proxy/src/network_policy.rs b/codex-rs/network-proxy/src/network_policy.rs index 46e19ed8ed36..956aa6f682df 100644 --- a/codex-rs/network-proxy/src/network_policy.rs +++ b/codex-rs/network-proxy/src/network_policy.rs @@ -701,9 +701,9 @@ mod tests { network.set_allowed_domains(vec!["example.com".to_string()]); network }); - state.register_execution("local", "execution-baseline-allow"); + state.register_execution("token-baseline-allow", "local", "execution-baseline-allow"); let state = state - .for_execution_token("execution-baseline-allow") + .for_execution_token("token-baseline-allow") .expect("expected registered execution"); let request = NetworkPolicyRequest::new(NetworkPolicyRequestArgs { protocol: NetworkProtocol::Http, @@ -731,6 +731,7 @@ mod tests { event.field("execution.id"), Some("execution-baseline-allow") ); + assert_ne!(event.field("execution.id"), Some("token-baseline-allow")); } #[tokio::test(flavor = "current_thread")] @@ -741,9 +742,9 @@ mod tests { network.set_denied_domains(vec!["blocked.com".to_string()]); network }); - state.register_execution("local", "execution-baseline-deny"); + state.register_execution("token-baseline-deny", "local", "execution-baseline-deny"); let state = state - .for_execution_token("execution-baseline-deny") + .for_execution_token("token-baseline-deny") .expect("expected registered execution"); let request = NetworkPolicyRequest::new(NetworkPolicyRequestArgs { protocol: NetworkProtocol::Http, @@ -783,6 +784,7 @@ mod tests { assert_eq!(event.field("http.request.method"), Some("GET")); assert_eq!(event.field("client.address"), Some("127.0.0.1:1234")); assert_eq!(event.field("execution.id"), Some("execution-baseline-deny")); + assert_ne!(event.field("execution.id"), Some("token-baseline-deny")); } #[tokio::test(flavor = "current_thread")] diff --git a/codex-rs/network-proxy/src/proxy.rs b/codex-rs/network-proxy/src/proxy.rs index fceb199bfe80..476c68aca8bb 100644 --- a/codex-rs/network-proxy/src/proxy.rs +++ b/codex-rs/network-proxy/src/proxy.rs @@ -708,7 +708,7 @@ impl NetworkProxy { if let Some(execution_scope) = self.execution_scope.as_ref() { env.insert( PROXY_ATTRIBUTION_TOKEN_ENV_KEY.to_string(), - execution_scope.execution_id.clone(), + execution_scope.attribution_token.clone(), ); } let mut loopback_ports = [ diff --git a/codex-rs/network-proxy/src/proxy/execution_scope.rs b/codex-rs/network-proxy/src/proxy/execution_scope.rs index 0c7e6b570e6a..37068625d491 100644 --- a/codex-rs/network-proxy/src/proxy/execution_scope.rs +++ b/codex-rs/network-proxy/src/proxy/execution_scope.rs @@ -2,40 +2,37 @@ use super::*; pub(super) struct ExecutionScope { pub(super) environment_id: String, - pub(super) execution_id: String, + pub(super) attribution_token: String, state: Arc, } impl Drop for ExecutionScope { fn drop(&mut self) { - self.state.unregister_execution(&self.execution_id); + self.state.unregister_execution(&self.attribution_token); } } impl NetworkProxy { /// Returns a proxy that attributes trusted bridge connections to one execution. - pub fn for_execution(&self, environment_id: &str, execution_id: String) -> Result { - #[cfg(not(target_os = "linux"))] - { - let _ = (environment_id, execution_id); - Ok(self.clone()) - } + pub fn for_execution( + &self, + environment_id: &str, + execution_id: &str, + attribution_token: String, + ) -> Result { + anyhow::ensure!( + self.execution_scope.is_none(), + "cannot scope an execution-scoped network proxy" + ); + self.state + .register_execution(&attribution_token, environment_id, execution_id); - #[cfg(target_os = "linux")] - { - anyhow::ensure!( - self.execution_scope.is_none(), - "cannot scope an execution-scoped network proxy" - ); - self.state.register_execution(environment_id, &execution_id); - - let mut proxy = self.clone(); - proxy.execution_scope = Some(Arc::new(ExecutionScope { - environment_id: environment_id.to_string(), - execution_id, - state: Arc::clone(&self.state), - })); - Ok(proxy) - } + let mut proxy = self.clone(); + proxy.execution_scope = Some(Arc::new(ExecutionScope { + environment_id: environment_id.to_string(), + attribution_token, + state: Arc::clone(&self.state), + })); + Ok(proxy) } } diff --git a/codex-rs/network-proxy/src/runtime.rs b/codex-rs/network-proxy/src/runtime.rs index f0e0e0dce604..a51f1d8f165b 100644 --- a/codex-rs/network-proxy/src/runtime.rs +++ b/codex-rs/network-proxy/src/runtime.rs @@ -215,7 +215,7 @@ pub struct NetworkProxyState { blocked_request_observer: Arc>>>, credential_broker: CredentialBroker, audit_metadata: NetworkProxyAuditMetadata, - execution_attributions: Arc>>, + execution_attributions: Arc>>, environment_id: Option>, execution_id: Option>, } @@ -227,6 +227,12 @@ pub(crate) enum HostMitmRequirement { Always, } +#[derive(Clone)] +struct ExecutionAttribution { + environment_id: String, + execution_id: String, +} + impl std::fmt::Debug for NetworkProxyState { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { // Avoid logging internal state (config contents, derived globsets, etc.) which can be noisy @@ -303,31 +309,41 @@ impl NetworkProxyState { } } - #[cfg(any(target_os = "linux", test))] - pub(crate) fn register_execution(&self, environment_id: &str, execution_id: &str) { + pub(crate) fn register_execution( + &self, + attribution_token: &str, + environment_id: &str, + execution_id: &str, + ) { self.execution_attributions .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert(execution_id.to_string(), environment_id.to_string()); + .insert( + attribution_token.to_string(), + ExecutionAttribution { + environment_id: environment_id.to_string(), + execution_id: execution_id.to_string(), + }, + ); } - pub(crate) fn unregister_execution(&self, execution_id: &str) { + pub(crate) fn unregister_execution(&self, attribution_token: &str) { self.execution_attributions .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(execution_id); + .remove(attribution_token); } pub(crate) fn for_execution_token(&self, token: &str) -> Option { - let environment_id = self + let attribution = self .execution_attributions .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .get(token)? .clone(); Some(Self { - environment_id: Some(environment_id.into()), - execution_id: Some(token.to_string().into()), + environment_id: Some(attribution.environment_id.into()), + execution_id: Some(attribution.execution_id.into()), ..self.clone() }) } From 139c6850ad42d429fef163ff66d0ae9d81b41f47 Mon Sep 17 00:00:00 2001 From: viyatb-oai Date: Thu, 25 Jun 2026 11:48:18 -0700 Subject: [PATCH 7/7] test: make concurrent network attribution deterministic Co-authored-by: Codex --- codex-rs/core/tests/suite/network_approval.rs | 165 ++++++++++++++---- 1 file changed, 127 insertions(+), 38 deletions(-) diff --git a/codex-rs/core/tests/suite/network_approval.rs b/codex-rs/core/tests/suite/network_approval.rs index 6f087ac53044..437155f32547 100644 --- a/codex-rs/core/tests/suite/network_approval.rs +++ b/codex-rs/core/tests/suite/network_approval.rs @@ -28,6 +28,7 @@ use core_test_support::responses::ev_assistant_message; use core_test_support::responses::ev_completed; use core_test_support::responses::ev_function_call; use core_test_support::responses::ev_response_created; +use core_test_support::responses::mount_sse_once_match; use core_test_support::responses::mount_sse_sequence; use core_test_support::responses::sse; use core_test_support::responses::start_mock_server; @@ -81,39 +82,61 @@ async fn guardian_receives_exact_triggers_for_concurrent_network_requests() -> R }; let first_command = network_command(&first_marker, &second_marker, "1.1.1.1"); let second_command = network_command(&second_marker, &first_marker, "8.8.8.8"); - let responses = mount_sse_sequence( + mount_sse_once_match( &server, - vec![ - sse(vec![ - ev_response_created("resp-network-concurrent"), - ev_function_call( - "exec-network-first", - "exec_command", - &serde_json::to_string(&network_exec_args(&first_command))?, - ), - ev_function_call( - "exec-network-second", - "exec_command", - &serde_json::to_string(&network_exec_args(&second_command))?, - ), - ev_completed("resp-network-concurrent"), - ]), - sse(vec![ - ev_response_created("resp-network-guardian-1"), - ev_assistant_message("msg-network-guardian-1", r#"{"outcome":"deny"}"#), - ev_completed("resp-network-guardian-1"), - ]), - sse(vec![ - ev_response_created("resp-network-guardian-2"), - ev_assistant_message("msg-network-guardian-2", r#"{"outcome":"deny"}"#), - ev_completed("resp-network-guardian-2"), - ]), - sse(vec![ - ev_response_created("resp-network-done"), - ev_assistant_message("msg-network-done", "done"), - ev_completed("resp-network-done"), - ]), - ], + |request: &wiremock::Request| { + !is_guardian_request(request) + && request_body_contains(request, "run both network requests") + && !request_body_contains(request, "exec-network-first") + }, + sse(vec![ + ev_response_created("resp-network-concurrent"), + ev_function_call( + "exec-network-first", + "exec_command", + &serde_json::to_string(&network_exec_args(&first_command))?, + ), + ev_function_call( + "exec-network-second", + "exec_command", + &serde_json::to_string(&network_exec_args(&second_command))?, + ), + ev_completed("resp-network-concurrent"), + ]), + ) + .await; + let first_guardian = mount_sse_once_match( + &server, + |request: &wiremock::Request| guardian_request_is_for(request, "exec-network-first"), + sse(vec![ + ev_response_created("resp-network-guardian-1"), + ev_assistant_message("msg-network-guardian-1", r#"{"outcome":"deny"}"#), + ev_completed("resp-network-guardian-1"), + ]), + ) + .await; + let second_guardian = mount_sse_once_match( + &server, + |request: &wiremock::Request| guardian_request_is_for(request, "exec-network-second"), + sse(vec![ + ev_response_created("resp-network-guardian-2"), + ev_assistant_message("msg-network-guardian-2", r#"{"outcome":"deny"}"#), + ev_completed("resp-network-guardian-2"), + ]), + ) + .await; + mount_sse_once_match( + &server, + |request: &wiremock::Request| { + !is_guardian_request(request) + && request_body_contains(request, "exec-network-first") + && request_body_contains(request, "exec-network-second") + }, + sse(vec![ + ev_response_created("resp-network-done"), + ev_assistant_message("msg-network-done", "done"), + ev_completed("resp-network-done"), + ]), ) .await; @@ -125,10 +148,21 @@ async fn guardian_receives_exact_triggers_for_concurrent_network_requests() -> R AskForApproval::OnRequest, ) .await?; + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let actual_triggers = loop { + let mut actual_triggers = guardian_network_triggers(&[&first_guardian, &second_guardian])?; + actual_triggers.sort_unstable(); + actual_triggers.dedup(); + if actual_triggers.len() == 2 { + break actual_triggers; + } + if tokio::time::Instant::now() >= deadline { + anyhow::bail!("timed out waiting for both Guardian network reviews"); + } + tokio::time::sleep(Duration::from_millis(10)).await; + }; wait_for_turn_complete(&test).await; - let mut actual_triggers = guardian_network_triggers(&responses)?; - actual_triggers.sort_unstable(); assert_eq!( actual_triggers, vec![ @@ -191,7 +225,7 @@ async fn guardian_receives_exact_trigger_for_single_network_request() -> Result< wait_for_turn_complete(&test).await; assert_eq!( - guardian_network_triggers(&responses)?, + guardian_network_triggers(&[&responses])?, vec![("exec-network-single".to_string(), command)] ); @@ -437,10 +471,65 @@ async fn submit_managed_network_turn( Ok(()) } -fn guardian_network_triggers(responses: &ResponseMock) -> Result> { +fn decoded_request_body(request: &wiremock::Request) -> Option> { + let is_zstd = request + .headers + .get("content-encoding") + .and_then(|value| value.to_str().ok()) + .is_some_and(|value| { + value + .split(',') + .any(|entry| entry.trim().eq_ignore_ascii_case("zstd")) + }); + if is_zstd { + zstd::stream::decode_all(std::io::Cursor::new(&request.body)).ok() + } else { + Some(request.body.clone()) + } +} + +fn request_body_contains(request: &wiremock::Request, text: &str) -> bool { + decoded_request_body(request) + .and_then(|body| String::from_utf8(body).ok()) + .is_some_and(|body| body.contains(text)) +} + +fn is_guardian_request(request: &wiremock::Request) -> bool { + decoded_request_body(request) + .and_then(|body| serde_json::from_slice::(&body).ok()) + .is_some_and(|body| { + body.pointer("/client_metadata/x-openai-subagent") + .and_then(Value::as_str) + == Some("guardian") + }) +} + +fn guardian_request_is_for(request: &wiremock::Request, call_id: &str) -> bool { + decoded_request_body(request) + .and_then(|body| serde_json::from_slice::(&body).ok()) + .filter(|body| { + body.pointer("/client_metadata/x-openai-subagent") + .and_then(Value::as_str) + == Some("guardian") + }) + .and_then(|body| { + body.get("input") + .and_then(Value::as_array) + .and_then(|input| { + input + .iter() + .rev() + .find(|item| item.get("role").and_then(Value::as_str) == Some("user")) + }) + .cloned() + }) + .is_some_and(|latest_user_message| latest_user_message.to_string().contains(call_id)) +} + +fn guardian_network_triggers(responses: &[&ResponseMock]) -> Result> { responses - .requests() - .into_iter() + .iter() + .flat_map(|responses| responses.requests()) .filter(|request| { request.body_json()["client_metadata"]["x-openai-subagent"].as_str() == Some("guardian") })