From ed94456aa037d6087866fa52ec729d305e7c26d9 Mon Sep 17 00:00:00 2001 From: Prekshi Vyas Date: Thu, 17 Sep 2026 14:54:12 -0700 Subject: [PATCH 1/2] fix(ocsf): attribute MXC proxy events to sandboxes NVBug 6783086 Signed-off-by: Prekshi Vyas --- architecture/security-policy.md | 7 + .../openshell-supervisor-network/src/host.rs | 77 +++++- .../openshell-supervisor-network/src/proxy.rs | 236 ++++++++++++++---- .../src/proxy/tests/compatibility.rs | 6 + 4 files changed, 270 insertions(+), 56 deletions(-) diff --git a/architecture/security-policy.md b/architecture/security-policy.md index 655886d87c..7503532ed8 100644 --- a/architecture/security-policy.md +++ b/architecture/security-policy.md @@ -418,6 +418,13 @@ middleware, token grants, credential rewriting, policy-generation checks, and the HTTP relay have succeeded so a later denial cannot coexist with an allowed record for the same request. +Each Windows MXC sandbox owns a separate host proxy. That proxy carries an +immutable per-sandbox OCSF context so proxy lifecycle events and top-level +CONNECT/forward decisions use the correct `container.uid` and `container.name` +even when one gateway serves multiple sandboxes concurrently. The process-wide +sandbox context is only suitable for the one-supervisor-per-sandbox runtime +model. + Never log secrets, credentials, bearer tokens, or query parameters in OCSF messages. OCSF JSONL output may be shipped to external systems. The gateway-local OCSF JSONL file sink is restricted to the Windows/MXC path diff --git a/crates/openshell-supervisor-network/src/host.rs b/crates/openshell-supervisor-network/src/host.rs index 9d6f5cfffa..ea4082b9fa 100644 --- a/crates/openshell-supervisor-network/src/host.rs +++ b/crates/openshell-supervisor-network/src/host.rs @@ -22,7 +22,7 @@ use openshell_core::proposals::AgentProposals; use openshell_core::proto::SandboxPolicy as ProtoSandboxPolicy; use openshell_core::provider_credentials::ProviderCredentialState; use openshell_ocsf::{ - ConfigStateChangeBuilder, SeverityId, StateId, StatusId, ctx::ctx as ocsf_ctx, ocsf_emit, + ConfigStateChangeBuilder, EventContext, SeverityId, StateId, StatusId, ocsf_emit, }; use tokio::sync::mpsc::UnboundedSender; @@ -68,7 +68,11 @@ pub struct HostProxyConfig { /// Per-sandbox client authentication. Host-side MXC proxies must set this /// so another sandbox cannot borrow this proxy's identity and policy. pub client_auth: HostProxyClientAuth, + /// Stable sandbox identifier used to attribute host-proxy OCSF events. + /// Required and non-empty for every host-side proxy. pub sandbox_id: Option, + /// Sandbox display name used to attribute host-proxy OCSF events. + /// Required and non-empty for every host-side proxy. pub sandbox_name: Option, pub openshell_endpoint: Option, pub provider_credentials: Option, @@ -99,6 +103,36 @@ impl HostProxyHandle { } } +fn host_proxy_event_context(config: &HostProxyConfig) -> Result { + let sandbox_id = config + .sandbox_id + .as_deref() + .map(str::trim) + .filter(|value| !value.is_empty()) + .ok_or_else(|| miette::miette!("host proxy requires a non-empty sandbox_id"))?; + let sandbox_name = config + .sandbox_name + .as_deref() + .map(str::trim) + .filter(|value| !value.is_empty()) + .ok_or_else(|| miette::miette!("host proxy requires a non-empty sandbox_name"))?; + + Ok(EventContext { + sandbox_id: sandbox_id.to_string(), + sandbox_name: sandbox_name.to_string(), + container_image: String::new(), + hostname: std::env::var("COMPUTERNAME") + .or_else(|_| std::env::var("HOSTNAME")) + .ok() + .map(|hostname| hostname.trim().to_string()) + .filter(|hostname| !hostname.is_empty()) + .unwrap_or_else(|| "openshell-gateway".to_string()), + product_version: env!("CARGO_PKG_VERSION").to_string(), + proxy_ip: config.bind_addr.ip(), + proxy_port: config.bind_addr.port(), + }) +} + /// Start a host-side proxy for one sandbox. /// /// Linux supervisor mode should continue to use `run::run_networking`; this API @@ -117,6 +151,7 @@ pub async fn start_host_proxy(config: HostProxyConfig) -> Result Result Result { ocsf_emit!( - ConfigStateChangeBuilder::new(ocsf_ctx()) + ConfigStateChangeBuilder::new(&event_context) .severity(SeverityId::High) .status(StatusId::Failure) .state(StateId::Disabled, "disabled") @@ -174,7 +209,7 @@ pub async fn start_host_proxy(config: HostProxyConfig) -> Result { ocsf_emit!( - ConfigStateChangeBuilder::new(ocsf_ctx()) + ConfigStateChangeBuilder::new(&event_context) .severity(SeverityId::High) .status(StatusId::Failure) .state(StateId::Disabled, "disabled") @@ -189,7 +224,7 @@ pub async fn start_host_proxy(config: HostProxyConfig) -> Result { ocsf_emit!( - ConfigStateChangeBuilder::new(ocsf_ctx()) + ConfigStateChangeBuilder::new(&event_context) .severity(SeverityId::High) .status(StatusId::Failure) .state(StateId::Disabled, "disabled") @@ -203,7 +238,7 @@ pub async fn start_host_proxy(config: HostProxyConfig) -> Result { ocsf_emit!( - ConfigStateChangeBuilder::new(ocsf_ctx()) + ConfigStateChangeBuilder::new(&event_context) .severity(SeverityId::High) .status(StatusId::Failure) .state(StateId::Disabled, "disabled") @@ -218,7 +253,8 @@ pub async fn start_host_proxy(config: HostProxyConfig) -> Result Result String { let mut request = String::from( "GET http://policy.local/v1/policy/current HTTP/1.1\r\nHost: policy.local\r\n", diff --git a/crates/openshell-supervisor-network/src/proxy.rs b/crates/openshell-supervisor-network/src/proxy.rs index 22bb8191c3..64a6fb4cf7 100644 --- a/crates/openshell-supervisor-network/src/proxy.rs +++ b/crates/openshell-supervisor-network/src/proxy.rs @@ -31,8 +31,8 @@ use openshell_core::policy::ProxyPolicy; use openshell_core::provider_credentials::{ProviderCredentialSnapshot, ProviderCredentialState}; use openshell_core::secrets::{self, SecretResolver, rewrite_header_line_checked}; use openshell_ocsf::{ - ActionId, ActivityId, DispositionId, Endpoint, HttpActivityBuilder, HttpRequest, HttpResponse, - NetworkActivityBuilder, Process, SeverityId, StatusId, Url as OcsfUrl, ocsf_emit, + ActionId, ActivityId, DispositionId, Endpoint, EventContext, HttpActivityBuilder, HttpRequest, + HttpResponse, NetworkActivityBuilder, Process, SeverityId, StatusId, Url as OcsfUrl, ocsf_emit, }; #[cfg(target_os = "linux")] use std::mem::size_of; @@ -175,6 +175,9 @@ pub(crate) enum ProxyIdentityMode { binary_path: PathBuf, binary_sha256: String, required_proxy_authorization: Option>, + /// Per-sandbox context for host-side proxies. The process-wide OCSF + /// context cannot identify one sandbox when a gateway hosts many. + event_context: Option>, }, } @@ -206,9 +209,33 @@ impl ProxyIdentityMode { binary_path, binary_sha256, required_proxy_authorization, + event_context: None, }) } + #[cfg(target_os = "windows")] + pub(super) fn with_event_context(mut self, context: EventContext) -> Self { + match &mut self { + #[cfg(target_os = "linux")] + Self::Procfs { .. } => {} + Self::Static { event_context, .. } => { + *event_context = Some(Arc::new(context)); + } + } + self + } + + fn event_context(&self) -> &EventContext { + match self { + #[cfg(any(not(target_os = "linux"), test))] + Self::Static { + event_context: Some(context), + .. + } => context, + _ => openshell_ocsf::ctx::ctx(), + } + } + fn required_proxy_authorization(&self) -> Option<&str> { match self { #[cfg(target_os = "linux")] @@ -269,8 +296,9 @@ impl ProxyHandle { let listener = TcpListener::bind(http_addr).await.into_diagnostic()?; let local_addr = listener.local_addr().into_diagnostic()?; + let event_context = Arc::new(identity_mode.event_context().clone()); { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Listen) .severity(SeverityId::Informational) .status(StatusId::Success) @@ -307,22 +335,21 @@ impl ProxyHandle { // access. let upstream_proxy: Arc> = Arc::new( UpstreamProxyConfig::from_args(upstream_proxy_args).map_err(|err| { - let event = - openshell_ocsf::ConfigStateChangeBuilder::new(openshell_ocsf::ctx::ctx()) - .severity(SeverityId::High) - .status(StatusId::Failure) - .state(openshell_ocsf::StateId::Disabled, "invalid") - .message(format!( - "Upstream corporate proxy configuration invalid; \ + let event = openshell_ocsf::ConfigStateChangeBuilder::new(&event_context) + .severity(SeverityId::High) + .status(StatusId::Failure) + .state(openshell_ocsf::StateId::Disabled, "invalid") + .message(format!( + "Upstream corporate proxy configuration invalid; \ refusing to start: {err}" - )) - .build(); + )) + .build(); ocsf_emit!(event); miette::miette!("invalid upstream corporate proxy configuration: {err}") })?, ); if let Some(cfg) = upstream_proxy.as_ref() { - let event = openshell_ocsf::ConfigStateChangeBuilder::new(openshell_ocsf::ctx::ctx()) + let event = openshell_ocsf::ConfigStateChangeBuilder::new(&event_context) .severity(SeverityId::Informational) .status(StatusId::Success) .state(openshell_ocsf::StateId::Enabled, "enabled") @@ -392,6 +419,7 @@ impl ProxyHandle { let dtx = denial_tx.clone(); let atx = activity_tx.clone(); let endpoint_observations = endpoint_observation_tx.clone(); + let event_context = event_context.clone(); tokio::spawn(async move { #[allow(clippy::large_futures)] if let Err(err) = handle_tcp_connection( @@ -412,7 +440,7 @@ impl ProxyHandle { ) .await { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -429,7 +457,7 @@ impl ProxyHandle { &mut consecutive_unknown_errors, ) { AcceptAction::Terminal => { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::High) .status(StatusId::Failure) @@ -441,7 +469,7 @@ impl ProxyHandle { break; } AcceptAction::Retry { backoff, severity } => { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(severity) .status(StatusId::Failure) @@ -1366,6 +1394,7 @@ fn emit_denial_simple( #[allow(clippy::too_many_arguments)] fn build_connect_allow_ocsf_event( + event_context: &EventContext, peer_addr: SocketAddr, host: &str, port: u16, @@ -1381,7 +1410,7 @@ fn build_connect_allow_ocsf_event( } else { "CONNECT" }; - NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + NetworkActivityBuilder::new(event_context) .activity(ActivityId::Open) .action(ActionId::Allowed) .disposition(DispositionId::Allowed) @@ -1397,6 +1426,7 @@ fn build_connect_allow_ocsf_event( #[allow(clippy::too_many_arguments)] fn build_forward_allow_ocsf_event( + event_context: &EventContext, peer_addr: SocketAddr, method: &str, host: &str, @@ -1408,7 +1438,7 @@ fn build_forward_allow_ocsf_event( cmdline: &str, policy: &str, ) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::for_http_method(method)) .action(ActionId::Allowed) .disposition(DispositionId::Allowed) @@ -1426,8 +1456,11 @@ fn build_forward_allow_ocsf_event( .build() } -fn build_forward_parse_error_ocsf_event(path: &str) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) +fn build_forward_parse_error_ocsf_event( + event_context: &EventContext, + path: &str, +) -> openshell_ocsf::OcsfEvent { + HttpActivityBuilder::new(event_context) .activity(ActivityId::Other) .http_response(HttpResponse { code: StatusCode::BAD_REQUEST.as_u16(), @@ -1443,12 +1476,13 @@ fn build_forward_parse_error_ocsf_event(path: &str) -> openshell_ocsf::OcsfEvent /// contain credentials; the method and generated response provide the HTTP /// context required by OCSF 1.8. fn build_forward_unsupported_scheme_ocsf_event( + event_context: &EventContext, method: &str, scheme: &str, host: &str, port: u16, ) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::for_http_method(method)) .http_request(HttpRequest { http_method: method.parse().expect("HTTP method parsing is infallible"), @@ -1470,6 +1504,7 @@ fn build_forward_unsupported_scheme_ocsf_event( #[allow(clippy::too_many_arguments)] fn build_forward_l7_parse_rejection_ocsf_event( + event_context: &EventContext, peer_addr: SocketAddr, method: &str, host: &str, @@ -1482,7 +1517,7 @@ fn build_forward_l7_parse_rejection_ocsf_event( policy: &str, detail: &str, ) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::for_http_method(method)) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -1505,6 +1540,7 @@ fn build_forward_l7_parse_rejection_ocsf_event( #[allow(clippy::too_many_arguments)] fn build_forward_policy_deny_ocsf_event( + event_context: &EventContext, peer_addr: SocketAddr, method: &str, host: &str, @@ -1516,7 +1552,7 @@ fn build_forward_policy_deny_ocsf_event( cmdline: &str, reason: &str, ) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::Other) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -1556,6 +1592,7 @@ fn endpoint_result_for_destination_failure(kind: DestinationDenialKind) -> Endpo #[allow(clippy::too_many_arguments)] fn build_connect_destination_deny_ocsf_event( + event_context: &EventContext, denial: &DestinationDenial, peer_addr: SocketAddr, host: &str, @@ -1572,7 +1609,7 @@ fn build_connect_destination_deny_ocsf_event( format!("CONNECT blocked: {detail} for {host}:{port}") }; - NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + NetworkActivityBuilder::new(event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -1589,6 +1626,7 @@ fn build_connect_destination_deny_ocsf_event( #[allow(clippy::too_many_arguments)] fn build_forward_destination_deny_ocsf_event( + event_context: &EventContext, denial: &DestinationDenial, peer_addr: SocketAddr, method: &str, @@ -1608,7 +1646,7 @@ fn build_forward_destination_deny_ocsf_event( detail }; - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::for_http_method(method)) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -1629,6 +1667,7 @@ fn build_forward_destination_deny_ocsf_event( #[allow(clippy::too_many_arguments)] async fn deny_connect_destination( + event_context: &EventContext, client: &mut TcpStream, denial: &DestinationDenial, peer_addr: SocketAddr, @@ -1644,7 +1683,15 @@ async fn deny_connect_destination( ) -> Result<()> { let detail = destination_denial_detail(denial.kind); ocsf_emit!(build_connect_destination_deny_ocsf_event( - denial, peer_addr, host, port, binary, pid, ancestors, cmdline, + event_context, + denial, + peer_addr, + host, + port, + binary, + pid, + ancestors, + cmdline, )); emit_denial( @@ -1675,6 +1722,7 @@ async fn deny_connect_destination( #[allow(clippy::too_many_arguments)] async fn deny_forward_destination( + event_context: &EventContext, client: &mut TcpStream, denial: &DestinationDenial, peer_addr: SocketAddr, @@ -1693,7 +1741,18 @@ async fn deny_forward_destination( ) -> Result<()> { let detail = destination_denial_detail(denial.kind); ocsf_emit!(build_forward_destination_deny_ocsf_event( - denial, peer_addr, method, host, port, path, binary, pid, ancestors, cmdline, policy, + event_context, + denial, + peer_addr, + method, + host, + port, + path, + binary, + pid, + ancestors, + cmdline, + policy, )); emit_denial_simple( @@ -1782,6 +1841,7 @@ async fn handle_tcp_connection( activity_tx: Option, endpoint_observation_tx: Option, ) -> Result<()> { + let event_context = identity_mode.event_context().clone(); // Capture authority before request parsing or policy selection can yield. // A later inventory installation cannot acquire this connection's result. let endpoint_observation_context = endpoint_observation_tx @@ -1935,7 +1995,7 @@ async fn handle_tcp_connection( // Allowed connections are logged after the L7 config check (below) // so we can distinguish CONNECT (L4-only) from CONNECT_L7 (L7 follows). if matches!(decision.action, NetworkAction::Deny { .. }) { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -2021,6 +2081,7 @@ async fn handle_tcp_connection( observer.observe(endpoint_result_for_destination_failure(denial.kind)); } deny_connect_destination( + &event_context, &mut client, &denial, workload_addr, @@ -2062,6 +2123,7 @@ async fn handle_tcp_connection( observer.observe(endpoint_result_for_destination_failure(denial.kind)); } deny_connect_destination( + &event_context, &mut client, &denial, workload_addr, @@ -2093,7 +2155,7 @@ async fn handle_tcp_connection( if let Some(observer) = connect_endpoint_observer.as_ref() { observer.observe(EndpointResult::TlsFailed); } - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -2131,7 +2193,7 @@ async fn handle_tcp_connection( if let Some(observer) = connect_endpoint_observer.as_ref() { observer.observe(EndpointResult::PolicyDenied); } - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -2231,6 +2293,7 @@ async fn handle_tcp_connection( // Log the allowed CONNECT — use CONNECT_L7 when L7 inspection follows, // so log consumers can distinguish L4-only decisions from tunnel lifecycle events. ocsf_emit!(build_connect_allow_ocsf_event( + &event_context, workload_addr, &host_lc, port, @@ -2339,7 +2402,7 @@ async fn handle_tcp_connection( if let Some(observer) = connect_endpoint_observer.as_ref() { observer.observe(EndpointResult::TlsFailed); } - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -2368,7 +2431,7 @@ async fn handle_tcp_connection( "TLS connection closed" ); } else { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -2391,7 +2454,7 @@ async fn handle_tcp_connection( // placeholder verbatim). const DETAIL: &str = "TLS termination unavailable after tunnel establishment; \ closing connection - credential rewrite would be bypassed"; - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -2443,7 +2506,7 @@ async fn handle_tcp_connection( } else { format!("HTTP relay error: {e}") }; - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -2462,7 +2525,7 @@ async fn handle_tcp_connection( if requirement == InspectionRequirement::RequiredMiddleware { crate::l7::middleware::emit_middleware_uninspectable(&ctx, protocol_detail, true); } - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -4491,6 +4554,7 @@ async fn handle_forward_proxy( activity_tx: Option<&ActivitySender>, endpoint_observation_tx: Option, ) -> Result<()> { + let event_context = identity_mode.event_context().clone(); // Capture authority before asynchronous authorization or credential selection. // Policy and provider snapshots must belong to this same installation. let endpoint_observation_context = endpoint_observation_tx @@ -4501,7 +4565,10 @@ async fn handle_forward_proxy( // canonicalized below before credential binding, policy-path evaluation, // upstream bytes, or telemetry consume it. let Ok((scheme, host, port, mut path)) = parse_proxy_uri(target_uri) else { - ocsf_emit!(build_forward_parse_error_ocsf_event(&telemetry_path)); + ocsf_emit!(build_forward_parse_error_ocsf_event( + &event_context, + &telemetry_path + )); respond(client, b"HTTP/1.1 400 Bad Request\r\n\r\n").await?; return Ok(()); }; @@ -4543,7 +4610,13 @@ async fn handle_forward_proxy( } if scheme != "http" { - let event = build_forward_unsupported_scheme_ocsf_event(method, &scheme, &host_lc, port); + let event = build_forward_unsupported_scheme_ocsf_event( + &event_context, + method, + &scheme, + &host_lc, + port, + ); ocsf_emit!(event); if scheme == "https" { respond( @@ -4624,6 +4697,7 @@ async fn handle_forward_proxy( NetworkAction::Allow { matched_policy } => matched_policy.clone(), NetworkAction::Deny { reason } => { ocsf_emit!(build_forward_policy_deny_ocsf_event( + &event_context, workload_addr, method, &host_lc, @@ -4735,7 +4809,7 @@ async fn handle_forward_proxy( let prepared_target = match prepare_forward_target(&path, canonicalize_options) { Ok(prepared) => prepared, Err(error) => { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Medium) .status(StatusId::Failure) @@ -4906,6 +4980,7 @@ async fn handle_forward_proxy( observer.observe(EndpointResult::PolicyDenied); } ocsf_emit!(build_forward_l7_parse_rejection_ocsf_event( + &event_context, workload_addr, method, &host_lc, @@ -4935,7 +5010,7 @@ async fn handle_forward_proxy( if let Some(observer) = endpoint_observer.as_ref() { observer.observe(EndpointResult::PolicyDenied); } - let event = HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = HttpActivityBuilder::new(&event_context) .activity(ActivityId::Other) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -5015,7 +5090,7 @@ async fn handle_forward_proxy( { Ok(info) => info, Err(e) => { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Medium) .status(StatusId::Failure) @@ -5069,7 +5144,7 @@ async fn handle_forward_proxy( { Ok(body) => body, Err(e) => { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Medium) .status(StatusId::Failure) @@ -5132,7 +5207,7 @@ async fn handle_forward_proxy( || { crate::l7::relay::evaluate_l7_request(&tunnel_engine, &l7_ctx, &request_info) .unwrap_or_else(|e| { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = NetworkActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -5197,7 +5272,7 @@ async fn handle_forward_proxy( ) }, ); - let event = HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = HttpActivityBuilder::new(&event_context) .activity(ActivityId::Other) .action(action_id) .disposition(disposition_id) @@ -5271,6 +5346,7 @@ async fn handle_forward_proxy( observer.observe(EndpointResult::PolicyDenied); } deny_forward_destination( + &event_context, client, &denial, workload_addr, @@ -5311,6 +5387,7 @@ async fn handle_forward_proxy( observer.observe(endpoint_result_for_destination_failure(denial.kind)); } deny_forward_destination( + &event_context, client, &denial, workload_addr, @@ -5702,7 +5779,7 @@ async fn handle_forward_proxy( if let Some(observer) = endpoint_observer.as_ref() { observer.observe(EndpointResult::TransportFailed); } - let event = HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + let event = HttpActivityBuilder::new(&event_context) .activity(ActivityId::Fail) .severity(SeverityId::Low) .status(StatusId::Failure) @@ -5860,6 +5937,7 @@ async fn handle_forward_proxy( // rewriting, generation checks, and the HTTP relay. Only now record the // final allowed outcome. ocsf_emit!(build_forward_allow_ocsf_event( + &event_context, workload_addr, method, &host_lc, @@ -6154,6 +6232,18 @@ mod tests { use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::{TcpListener, TcpStream}; + fn proxy_event_context(sandbox_id: &str, sandbox_name: &str) -> EventContext { + EventContext { + sandbox_id: sandbox_id.to_string(), + sandbox_name: sandbox_name.to_string(), + container_image: String::new(), + hostname: "windows-gateway".to_string(), + product_version: "test".to_string(), + proxy_ip: Ipv4Addr::LOCALHOST.into(), + proxy_port: 3128, + } + } + #[test] fn endpoint_result_distinguishes_resolution_from_policy() { assert_eq!( @@ -7169,6 +7259,7 @@ network_policies: fn forward_policy_denial_ocsf_includes_validation_rationale() { let reason = "policy validation failed; fail-closed quarantine is active; candidate version 7 rejected: conflicting tls metadata"; let event = build_forward_policy_deny_ocsf_event( + openshell_ocsf::ctx::ctx(), "127.0.0.1:45123".parse().unwrap(), "GET", "api.example.com", @@ -7187,9 +7278,45 @@ network_policies: assert_eq!(json["disposition"], "Blocked"); } + #[test] + fn forward_policy_denial_ocsf_keeps_per_proxy_sandbox_attribution() { + use openshell_ocsf::validation::{load_class_schema, validate_required_fields}; + + let peer = "127.0.0.1:45123".parse().unwrap(); + let build_event = |context: &EventContext| { + build_forward_policy_deny_ocsf_event( + context, + peer, + "GET", + "api.example.com", + 80, + "/v1/models", + r"C:\agent.exe", + "-", + "-", + r"C:\agent.exe", + "endpoint is not allowed by any policy", + ) + .to_json() + .unwrap() + }; + + let sandbox_a = build_event(&proxy_event_context("sandbox-a-id", "sandbox-a")); + let sandbox_b = build_event(&proxy_event_context("sandbox-b-id", "sandbox-b")); + + assert_eq!(sandbox_a["container"]["uid"], "sandbox-a-id"); + assert_eq!(sandbox_a["container"]["name"], "sandbox-a"); + assert_eq!(sandbox_b["container"]["uid"], "sandbox-b-id"); + assert_eq!(sandbox_b["container"]["name"], "sandbox-b"); + let schema = load_class_schema("http_activity"); + validate_required_fields(&sandbox_a, &schema); + validate_required_fields(&sandbox_b, &schema); + } + #[test] fn forward_l7_parse_rejection_ocsf_includes_denial_context() { let event = build_forward_l7_parse_rejection_ocsf_event( + openshell_ocsf::ctx::ctx(), "127.0.0.1:45123".parse().unwrap(), "GET", "api.example.com", @@ -7364,6 +7491,7 @@ network_policies: assert_eq!(path, "/v1/[CREDENTIAL]"); let allowed = build_forward_allow_ocsf_event( + openshell_ocsf::ctx::ctx(), peer, "GET", "api.example.com", @@ -7378,6 +7506,7 @@ network_policies: .to_json() .unwrap(); let denied = build_forward_policy_deny_ocsf_event( + openshell_ocsf::ctx::ctx(), peer, "GET", "api.example.com", @@ -7404,6 +7533,7 @@ network_policies: assert_eq!(host, "api.example.com"); assert_eq!(path, "/?token=real-secret"); let no_path_query = build_forward_allow_ocsf_event( + openshell_ocsf::ctx::ctx(), peer, "GET", &host, @@ -7423,9 +7553,12 @@ network_policies: assert!(!serialized.contains("real-secret"), "{serialized}"); assert!(!serialized.contains("?token="), "{serialized}"); - let malformed = build_forward_parse_error_ocsf_event(&forward_telemetry_path( - "not-a-uri?token=real-secret&key=openshell:resolve:env:API_TOKEN", - )) + let malformed = build_forward_parse_error_ocsf_event( + openshell_ocsf::ctx::ctx(), + &forward_telemetry_path( + "not-a-uri?token=real-secret&key=openshell:resolve:env:API_TOKEN", + ), + ) .to_json() .unwrap(); assert_eq!( @@ -10387,8 +10520,13 @@ network_policies: fn unsupported_forward_scheme_event_omits_request_url() { use openshell_ocsf::validation::{load_class_schema, validate_required_fields}; - let event = - build_forward_unsupported_scheme_ocsf_event("GET", "https", "api.example.com", 443); + let event = build_forward_unsupported_scheme_ocsf_event( + openshell_ocsf::ctx::ctx(), + "GET", + "https", + "api.example.com", + 443, + ); let json = event.to_json().unwrap(); assert_eq!(json["http_request"]["http_method"], "GET"); diff --git a/crates/openshell-supervisor-network/src/proxy/tests/compatibility.rs b/crates/openshell-supervisor-network/src/proxy/tests/compatibility.rs index 5948cd5224..759f888b20 100644 --- a/crates/openshell-supervisor-network/src/proxy/tests/compatibility.rs +++ b/crates/openshell-supervisor-network/src/proxy/tests/compatibility.rs @@ -92,6 +92,7 @@ async fn destination_denials_preserve_adapter_specific_wire_contracts() { let (mut app, mut proxy) = tcp_pair().await; deny_connect_destination( + openshell_ocsf::ctx::ctx(), &mut proxy, &denial, peer, @@ -119,6 +120,7 @@ async fn destination_denials_preserve_adapter_specific_wire_contracts() { let (mut app, mut proxy) = tcp_pair().await; deny_forward_destination( + openshell_ocsf::ctx::ctx(), &mut proxy, &denial, peer, @@ -165,6 +167,7 @@ fn representative_adapter_denials_preserve_ocsf_fields() { // global tracing pipeline. Its callsite-interest cache is process-global, // so parallel tests can otherwise make captured-event assertions flaky. let connect = serde_json::to_value(build_connect_destination_deny_ocsf_event( + openshell_ocsf::ctx::ctx(), &denial, peer, "target.example", @@ -194,6 +197,7 @@ fn representative_adapter_denials_preserve_ocsf_fields() { assert_eq!(connect["status_detail"], denial_reason); let forward = serde_json::to_value(build_forward_destination_deny_ocsf_event( + openshell_ocsf::ctx::ctx(), &denial, peer, "POST", @@ -229,6 +233,7 @@ fn representative_adapter_denials_preserve_ocsf_fields() { fn representative_adapter_allows_preserve_ocsf_fields() { let peer: SocketAddr = "127.0.0.1:41000".parse().unwrap(); let connect = serde_json::to_value(build_connect_allow_ocsf_event( + openshell_ocsf::ctx::ctx(), peer, "target.example", 8443, @@ -254,6 +259,7 @@ fn representative_adapter_allows_preserve_ocsf_fields() { assert_eq!(connect["message"], "CONNECT_L7 allowed target.example:8443"); let forward = serde_json::to_value(build_forward_allow_ocsf_event( + openshell_ocsf::ctx::ctx(), peer, "GET", "target.example", From 1e5cd656d8bff5efc6e51d067a65233c67ffc46d Mon Sep 17 00:00:00 2001 From: Prekshi Vyas Date: Mon, 21 Sep 2026 23:35:29 -0700 Subject: [PATCH 2/2] fix(ocsf): preserve sandbox context across proxy denials Signed-off-by: Prekshi Vyas --- .../src/l7/mod.rs | 29 ++- .../src/l7/relay.rs | 1 + .../src/l7/websocket.rs | 8 +- .../openshell-supervisor-network/src/proxy.rs | 208 +++++++++++++++--- .../src/proxy/relay.rs | 86 +++++++- 5 files changed, 286 insertions(+), 46 deletions(-) diff --git a/crates/openshell-supervisor-network/src/l7/mod.rs b/crates/openshell-supervisor-network/src/l7/mod.rs index 2085b46681..d0ecb1e18a 100644 --- a/crates/openshell-supervisor-network/src/l7/mod.rs +++ b/crates/openshell-supervisor-network/src/l7/mod.rs @@ -35,6 +35,7 @@ use std::sync::{Arc, Mutex}; use crate::opa::PolicyGenerationGuard; pub(crate) fn build_credential_endpoint_mismatch_finding( + event_context: &openshell_ocsf::EventContext, policy_name: &str, host: &str, protocol: Option<&str>, @@ -50,7 +51,7 @@ pub(crate) fn build_credential_endpoint_mismatch_finding( } evidence.push(("disposition", "denied")); - DetectionFindingBuilder::new(openshell_ocsf::ctx::ctx()) + DetectionFindingBuilder::new(event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -413,8 +414,13 @@ pub fn parse_l7_config(val: ®orus::Value) -> Option { }) } -pub(crate) fn emit_uninspected_credential_finding(host: &str, policy_name: &str, surface: &str) { - let event = openshell_ocsf::DetectionFindingBuilder::new(openshell_ocsf::ctx::ctx()) +pub(crate) fn build_uninspected_credential_finding( + event_context: &openshell_ocsf::EventContext, + host: &str, + policy_name: &str, + surface: &str, +) -> openshell_ocsf::OcsfEvent { + openshell_ocsf::DetectionFindingBuilder::new(event_context) .severity(openshell_ocsf::SeverityId::High) .finding_info(openshell_ocsf::FindingInfo::new( "openshell.credentials.traffic_uninspectable", @@ -427,8 +433,21 @@ pub(crate) fn emit_uninspected_credential_finding(host: &str, policy_name: &str, ("disposition", "denied"), ]) .message("Uninspected credential-bearing traffic denied") - .build(); - openshell_ocsf::ocsf_emit!(event); + .build() +} + +pub(crate) fn emit_uninspected_credential_finding( + event_context: &openshell_ocsf::EventContext, + host: &str, + policy_name: &str, + surface: &str, +) { + openshell_ocsf::ocsf_emit!(build_uninspected_credential_finding( + event_context, + host, + policy_name, + surface, + )); } impl L7EndpointConfig { diff --git a/crates/openshell-supervisor-network/src/l7/relay.rs b/crates/openshell-supervisor-network/src/l7/relay.rs index 31acd3a7fd..b319c54112 100644 --- a/crates/openshell-supervisor-network/src/l7/relay.rs +++ b/crates/openshell-supervisor-network/src/l7/relay.rs @@ -371,6 +371,7 @@ fn build_credential_resolution_event( fn build_credential_endpoint_mismatch_finding(ctx: &L7EvalContext) -> openshell_ocsf::OcsfEvent { crate::l7::build_credential_endpoint_mismatch_finding( + openshell_ocsf::ctx::ctx(), &ctx.policy_name, &ctx.host, None, diff --git a/crates/openshell-supervisor-network/src/l7/websocket.rs b/crates/openshell-supervisor-network/src/l7/websocket.rs index dab11a1112..60f2b1c535 100644 --- a/crates/openshell-supervisor-network/src/l7/websocket.rs +++ b/crates/openshell-supervisor-network/src/l7/websocket.rs @@ -636,6 +636,7 @@ fn emit_credential_endpoint_mismatch(host: &str, port: u16, policy_name: &str) { .build() ); ocsf_emit!(crate::l7::build_credential_endpoint_mismatch_finding( + openshell_ocsf::ctx::ctx(), policy_name, host, Some("websocket"), @@ -1387,7 +1388,12 @@ fn emit_uninspected_credential_denial(host: &str, port: u16, policy_name: &str, )) .build(); ocsf_emit!(event); - crate::l7::emit_uninspected_credential_finding(host, policy_name, surface); + crate::l7::emit_uninspected_credential_finding( + openshell_ocsf::ctx::ctx(), + host, + policy_name, + surface, + ); } fn inspect_websocket_text_message( diff --git a/crates/openshell-supervisor-network/src/proxy.rs b/crates/openshell-supervisor-network/src/proxy.rs index 1f3285e1e6..0bfea5731d 100644 --- a/crates/openshell-supervisor-network/src/proxy.rs +++ b/crates/openshell-supervisor-network/src/proxy.rs @@ -74,12 +74,13 @@ const FORWARD_ENCODED_SLASH_REJECTION_DETAIL: &str = const SIDECAR_SUPERVISOR_TOPOLOGY: &str = "sidecar"; fn build_credential_endpoint_mismatch_event( + event_context: &EventContext, method: &str, host: &str, port: u16, policy_name: &str, ) -> openshell_ocsf::OcsfEvent { - HttpActivityBuilder::new(openshell_ocsf::ctx::ctx()) + HttpActivityBuilder::new(event_context) .activity(ActivityId::for_http_method(method)) .http_request(HttpRequest { http_method: method.parse().expect("HTTP method parsing is infallible"), @@ -99,10 +100,18 @@ fn build_credential_endpoint_mismatch_event( .build() } -fn emit_credential_endpoint_mismatch(method: &str, host: &str, port: u16, policy_name: &str) { - let event = build_credential_endpoint_mismatch_event(method, host, port, policy_name); +fn emit_credential_endpoint_mismatch( + event_context: &EventContext, + method: &str, + host: &str, + port: u16, + policy_name: &str, +) { + let event = + build_credential_endpoint_mismatch_event(event_context, method, host, port, policy_name); ocsf_emit!(event); let finding = crate::l7::build_credential_endpoint_mismatch_finding( + event_context, policy_name, host, None, @@ -818,7 +827,14 @@ async fn handle_transparent_tcp_connection( } )); emit_activity(&activity_tx, false, "transparent_tcp"); - relay::relay_tcp(&mut client, &mut upstream, &generation_guard, &ctx).await + relay::relay_tcp( + &mut client, + &mut upstream, + &generation_guard, + &ctx, + openshell_ocsf::ctx::ctx(), + ) + .await } #[cfg(any(target_os = "linux", test))] @@ -2041,6 +2057,7 @@ async fn handle_tcp_connection( Ok(guard) => guard, Err(error) => { reject_stale_connect_policy( + &event_context, &mut client, &host_lc, port, @@ -2213,6 +2230,7 @@ async fn handle_tcp_connection( .build(); ocsf_emit!(event); crate::l7::emit_uninspected_credential_finding( + &event_context, &host_lc, policy_str, if effective_tls_skip { "tls-skip" } else { "l4" }, @@ -2241,8 +2259,15 @@ async fn handle_tcp_connection( if let Err(error) = relay::validate_route_generation(l7_route, connect_generation_guard.captured_generation()) { - reject_stale_connect_policy(&mut client, &host_lc, port, activity_tx.as_ref(), error) - .await?; + reject_stale_connect_policy( + &event_context, + &mut client, + &host_lc, + port, + activity_tx.as_ref(), + error, + ) + .await?; return Ok(()); } @@ -2252,6 +2277,7 @@ async fn handle_tcp_connection( }; let Some(upstream_result) = upstream_result else { reject_stale_connect_policy( + &event_context, &mut client, &host_lc, port, @@ -2276,8 +2302,15 @@ async fn handle_tcp_connection( } }; if let Err(error) = connect_generation_guard.ensure_current() { - reject_stale_connect_policy(&mut client, &host_lc, port, activity_tx.as_ref(), error) - .await?; + reject_stale_connect_policy( + &event_context, + &mut client, + &host_lc, + port, + activity_tx.as_ref(), + error, + ) + .await?; return Ok(()); } @@ -2360,11 +2393,19 @@ async fn handle_tcp_connection( port = port, "tls: skip — bypassing TLS auto-detection, raw tunnel" ); - let Some(generation_guard) = relay::prepare_raw_relay(l7_route, &opa_engine, &decision) + let Some(generation_guard) = + relay::prepare_raw_relay(l7_route, &opa_engine, &decision, &event_context) else { return Ok(()); }; - relay::relay_tcp(&mut client, &mut upstream, &generation_guard, &ctx).await?; + relay::relay_tcp( + &mut client, + &mut upstream, + &generation_guard, + &ctx, + &event_context, + ) + .await?; return Ok(()); } @@ -2414,9 +2455,13 @@ async fn handle_tcp_connection( } }; let tls_result = async { - let Some(relay_context) = - relay::prepare_http_relay(l7_route, &opa_engine, &decision, &ctx) - else { + let Some(relay_context) = relay::prepare_http_relay( + l7_route, + &opa_engine, + &decision, + &ctx, + &event_context, + ) else { return Ok(()); }; @@ -2489,7 +2534,8 @@ async fn handle_tcp_connection( // Plaintext HTTP detected. ctx.request_default_port = Some(80); let is_l7_relay = l7_route.is_some_and(|route| !route.configs.is_empty()); - let Some(relay_context) = relay::prepare_http_relay(l7_route, &opa_engine, &decision, &ctx) + let Some(relay_context) = + relay::prepare_http_relay(l7_route, &opa_engine, &decision, &ctx, &event_context) else { return Ok(()); }; @@ -2576,11 +2622,19 @@ async fn handle_tcp_connection( port = port, "Non-TLS non-HTTP traffic detected, raw tunnel" ); - let Some(generation_guard) = relay::prepare_raw_relay(l7_route, &opa_engine, &decision) + let Some(generation_guard) = + relay::prepare_raw_relay(l7_route, &opa_engine, &decision, &event_context) else { return Ok(()); }; - relay::relay_tcp(&mut client, &mut upstream, &generation_guard, &ctx).await?; + relay::relay_tcp( + &mut client, + &mut upstream, + &generation_guard, + &ctx, + &event_context, + ) + .await?; } Ok(()) @@ -3044,8 +3098,13 @@ fn authorize_egress_intent( } } -fn emit_l7_tunnel_close_after_policy_change(host: &str, port: u16, error: miette::Report) { - let event = NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) +fn build_l7_tunnel_close_after_policy_change_event( + event_context: &EventContext, + host: &str, + port: u16, + error: &miette::Report, +) -> openshell_ocsf::OcsfEvent { + NetworkActivityBuilder::new(event_context) .activity(ActivityId::Open) .action(ActionId::Denied) .disposition(DispositionId::Blocked) @@ -3055,11 +3114,21 @@ fn emit_l7_tunnel_close_after_policy_change(host: &str, port: u16, error: miette .message(format!( "L7 tunnel closed before inspection because policy changed: {error}" )) - .build(); + .build() +} + +fn emit_l7_tunnel_close_after_policy_change( + event_context: &EventContext, + host: &str, + port: u16, + error: miette::Report, +) { + let event = build_l7_tunnel_close_after_policy_change_event(event_context, host, port, &error); ocsf_emit!(event); } async fn reject_stale_connect_policy( + event_context: &EventContext, client: &mut TcpStream, host: &str, port: u16, @@ -3072,7 +3141,7 @@ async fn reject_stale_connect_policy( error = %error, "CONNECT rejected because policy changed after L4 authorization" ); - emit_l7_tunnel_close_after_policy_change(host, port, error); + emit_l7_tunnel_close_after_policy_change(event_context, host, port, error); emit_activity_simple(activity_tx, true, "policy_stale"); respond( client, @@ -4759,7 +4828,7 @@ async fn handle_forward_proxy( error = %e, "Forward proxy rejected request because policy generation changed after L4 decision" ); - emit_l7_tunnel_close_after_policy_change(&host_lc, port, e); + emit_l7_tunnel_close_after_policy_change(&event_context, &host_lc, port, e); emit_activity_simple(activity_tx, true, "policy_stale"); respond( client, @@ -4880,6 +4949,7 @@ async fn handle_forward_proxy( "Forward proxy rejected request because L7 route lookup used a different policy generation" ); emit_l7_tunnel_close_after_policy_change( + &event_context, &host_lc, port, miette::miette!( @@ -4912,7 +4982,7 @@ async fn handle_forward_proxy( error = %e, "Forward proxy rejected request because L7 tunnel engine could not be cloned" ); - emit_l7_tunnel_close_after_policy_change(&host_lc, port, e); + emit_l7_tunnel_close_after_policy_change(&event_context, &host_lc, port, e); emit_activity_simple(activity_tx, true, "policy_stale"); respond( client, @@ -5421,7 +5491,7 @@ async fn handle_forward_proxy( error = %e, "Forward proxy rejected request because policy changed before upstream connect" ); - emit_l7_tunnel_close_after_policy_change(&host_lc, port, e); + emit_l7_tunnel_close_after_policy_change(&event_context, &host_lc, port, e); emit_activity_simple(activity_tx, true, "policy_stale"); respond( client, @@ -5452,6 +5522,7 @@ async fn handle_forward_proxy( observer.observe(EndpointResult::PolicyDenied); } emit_l7_tunnel_close_after_policy_change( + &event_context, &host_lc, port, miette::miette!( @@ -5674,7 +5745,13 @@ async fn handle_forward_proxy( if let Some(observer) = endpoint_observer.as_ref() { observer.observe_credential_failure(true); } - emit_credential_endpoint_mismatch(method, &host_lc, port, policy_str); + emit_credential_endpoint_mismatch( + &event_context, + method, + &host_lc, + port, + policy_str, + ); respond( client, &build_json_error_response( @@ -5749,7 +5826,7 @@ async fn handle_forward_proxy( error = %e, "Forward proxy rejected request because policy changed before relay" ); - emit_l7_tunnel_close_after_policy_change(&host_lc, port, e); + emit_l7_tunnel_close_after_policy_change(&event_context, &host_lc, port, e); if let Some(session) = middleware_session.take() { session .end(openshell_core::proto::MiddlewareSessionEndReason::PolicyReload) @@ -5829,7 +5906,7 @@ async fn handle_forward_proxy( error = %e, "Forward proxy rejected request because policy changed during upstream connect" ); - emit_l7_tunnel_close_after_policy_change(&host_lc, port, e); + emit_l7_tunnel_close_after_policy_change(&event_context, &host_lc, port, e); if let Some(session) = middleware_session.take() { session .end(openshell_core::proto::MiddlewareSessionEndReason::PolicyReload) @@ -10542,18 +10619,93 @@ network_policies: fn credential_endpoint_mismatch_event_includes_method_and_response() { use openshell_ocsf::validation::{load_class_schema, validate_required_fields}; - let event = - build_credential_endpoint_mismatch_event("POST", "api.example.com", 443, "bound"); + let event_context = proxy_event_context("sandbox-incident-id", "sandbox-incident"); + let event = build_credential_endpoint_mismatch_event( + &event_context, + "POST", + "api.example.com", + 443, + "bound", + ); let json = event.to_json().unwrap(); + let finding = crate::l7::build_credential_endpoint_mismatch_finding( + &event_context, + "bound", + "api.example.com", + None, + "Provider credential endpoint binding mismatch; request denied", + ) + .to_json() + .unwrap(); assert_eq!(json["class_uid"], 4002); assert_eq!(json["activity_name"], "Post"); assert_eq!(json["http_request"]["http_method"], "POST"); assert!(json["http_request"].get("url").is_none()); assert_eq!(json["http_response"]["code"], 403); + for incident_event in [&json, &finding] { + assert_eq!(incident_event["container"]["uid"], "sandbox-incident-id"); + assert_eq!(incident_event["container"]["name"], "sandbox-incident"); + } validate_required_fields(&json, &load_class_schema("http_activity")); } + #[test] + fn uninspected_credential_finding_keeps_per_proxy_sandbox_attribution() { + let event_context = proxy_event_context("sandbox-uninspected-id", "sandbox-uninspected"); + let finding = crate::l7::build_uninspected_credential_finding( + &event_context, + "api.example.com", + "bound", + "tls-skip", + ) + .to_json() + .unwrap(); + + assert_eq!(finding["container"]["uid"], "sandbox-uninspected-id"); + assert_eq!(finding["container"]["name"], "sandbox-uninspected"); + } + + #[tokio::test] + async fn stale_connect_rejection_keeps_context_and_returns_policy_denial() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let client = TcpStream::connect(address); + let accepted = listener.accept(); + let (client, accepted) = tokio::join!(client, accepted); + let mut client = client.unwrap(); + let (mut proxy, _) = accepted.unwrap(); + let event_context = proxy_event_context("sandbox-stale-id", "sandbox-stale"); + + reject_stale_connect_policy( + &event_context, + &mut proxy, + "api.example.com", + 443, + None, + miette::miette!("policy generation is stale"), + ) + .await + .unwrap(); + + let mut response = vec![0; 512]; + let bytes_read = client.read(&mut response).await.unwrap(); + let response = String::from_utf8_lossy(&response[..bytes_read]); + assert!(response.starts_with("HTTP/1.1 403 Forbidden"), "{response}"); + assert!(response.contains("policy_denied"), "{response}"); + + let event = build_l7_tunnel_close_after_policy_change_event( + &event_context, + "api.example.com", + 443, + &miette::miette!("policy generation is stale"), + ) + .to_json() + .unwrap(); + assert_eq!(event["container"]["uid"], "sandbox-stale-id"); + assert_eq!(event["container"]["name"], "sandbox-stale"); + } + #[test] fn forward_credentials_capture_endpoint_resolver_and_revision_together() { use openshell_core::proto::{StaticCredentialBinding, StaticCredentialEndpointBinding}; diff --git a/crates/openshell-supervisor-network/src/proxy/relay.rs b/crates/openshell-supervisor-network/src/proxy/relay.rs index 4c89e5e382..485b5eb04d 100644 --- a/crates/openshell-supervisor-network/src/proxy/relay.rs +++ b/crates/openshell-supervisor-network/src/proxy/relay.rs @@ -36,6 +36,7 @@ pub(super) struct RelayContext<'a> { request: &'a L7EvalContext, policy: PreparedHttpPolicy, middleware_engine: &'a OpaEngine, + event_context: &'a openshell_ocsf::EventContext, } /// Non-blocking observation channels attached to an authorized HTTP relay. @@ -146,9 +147,11 @@ pub(super) fn prepare_http_relay<'a>( opa_engine: &'a OpaEngine, decision: &EgressDecision, request: &'a L7EvalContext, + event_context: &'a openshell_ocsf::EventContext, ) -> Option> { if let Err(error) = validate_route_generation(route, decision.policy_generation) { emit_l7_tunnel_close_after_policy_change( + event_context, &decision.intent.destination.host, decision.intent.destination.port, error, @@ -161,6 +164,7 @@ pub(super) fn prepare_http_relay<'a>( Ok(evaluator) => evaluator, Err(error) => { emit_l7_tunnel_close_after_policy_change( + event_context, &decision.intent.destination.host, decision.intent.destination.port, error, @@ -182,6 +186,7 @@ pub(super) fn prepare_http_relay<'a>( Ok(guard) => guard, Err(error) => { emit_l7_tunnel_close_after_policy_change( + event_context, &decision.intent.destination.host, decision.intent.destination.port, error, @@ -196,6 +201,7 @@ pub(super) fn prepare_http_relay<'a>( request, policy, middleware_engine: opa_engine, + event_context, }) } @@ -206,9 +212,11 @@ pub(super) fn prepare_raw_relay( route: Option<&L7RouteSnapshot>, opa_engine: &OpaEngine, decision: &EgressDecision, + event_context: &openshell_ocsf::EventContext, ) -> Option { if let Err(error) = validate_route_generation(route, decision.policy_generation) { emit_l7_tunnel_close_after_policy_change( + event_context, &decision.intent.destination.host, decision.intent.destination.port, error, @@ -220,6 +228,7 @@ pub(super) fn prepare_raw_relay( Ok(guard) => Some(guard), Err(error) => { emit_l7_tunnel_close_after_policy_change( + event_context, &decision.intent.destination.host, decision.intent.destination.port, error, @@ -255,7 +264,11 @@ where context.request, ) => result, () = generation_guard.wait_until_stale() => { - emit_stale_relay_close(context.request, &generation_guard); + emit_stale_relay_close( + context.request, + &generation_guard, + context.event_context, + ); Ok(()) } } @@ -271,7 +284,11 @@ where context.request, ) => result, () = generation_guard.wait_until_stale() => { - emit_stale_relay_close(context.request, &generation_guard); + emit_stale_relay_close( + context.request, + &generation_guard, + context.event_context, + ); Ok(()) } } @@ -286,7 +303,11 @@ where Some(context.middleware_engine), ) => result, () = generation_guard.wait_until_stale() => { - emit_stale_relay_close(context.request, &generation_guard); + emit_stale_relay_close( + context.request, + &generation_guard, + context.event_context, + ); Ok(()) } } @@ -300,6 +321,7 @@ pub(super) async fn relay_tcp( upstream: &mut U, generation_guard: &PolicyGenerationGuard, request: &L7EvalContext, + event_context: &openshell_ocsf::EventContext, ) -> Result<()> where C: AsyncRead + AsyncWrite + Unpin, @@ -310,14 +332,19 @@ where result.into_diagnostic()?; } () = generation_guard.wait_until_stale() => { - emit_stale_relay_close(request, generation_guard); + emit_stale_relay_close(request, generation_guard, event_context); } } Ok(()) } -fn emit_stale_relay_close(request: &L7EvalContext, guard: &PolicyGenerationGuard) { +fn emit_stale_relay_close( + request: &L7EvalContext, + guard: &PolicyGenerationGuard, + event_context: &openshell_ocsf::EventContext, +) { emit_l7_tunnel_close_after_policy_change( + event_context, &request.host, request.port, miette::miette!( @@ -380,8 +407,14 @@ mod tests { let decision = decision(engine.current_generation()); let request = request_context(); - let context = prepare_http_relay(None, &engine, &decision, &request) - .expect("current L4 generation should prepare a relay"); + let context = prepare_http_relay( + None, + &engine, + &decision, + &request, + openshell_ocsf::ctx::ctx(), + ) + .expect("current L4 generation should prepare a relay"); let PreparedHttpPolicy::Passthrough { generation_guard } = context.policy else { panic!("route-less relay should use a generation guard"); }; @@ -403,7 +436,14 @@ mod tests { let request = request_context(); assert!( - prepare_http_relay(Some(&route), &engine, &decision, &request).is_none(), + prepare_http_relay( + Some(&route), + &engine, + &decision, + &request, + openshell_ocsf::ctx::ctx(), + ) + .is_none(), "a current L7 lookup must not freshen a stale L4 allow" ); } @@ -441,7 +481,14 @@ mod tests { let request = request_context(); assert!( - prepare_http_relay(Some(&route), &engine, &decision, &request).is_none(), + prepare_http_relay( + Some(&route), + &engine, + &decision, + &request, + openshell_ocsf::ctx::ctx(), + ) + .is_none(), "an inspected route must use the generation that authorized CONNECT" ); } @@ -456,7 +503,8 @@ mod tests { }; assert!( - prepare_raw_relay(Some(&route), &engine, &decision).is_none(), + prepare_raw_relay(Some(&route), &engine, &decision, openshell_ocsf::ctx::ctx(),) + .is_none(), "a raw relay must not freshen a stale L4 allow" ); } @@ -469,7 +517,14 @@ mod tests { engine.reload(POLICY_REGO, EMPTY_POLICY_DATA).unwrap(); assert!( - prepare_http_relay(None, &engine, &decision, &request).is_none(), + prepare_http_relay( + None, + &engine, + &decision, + &request, + openshell_ocsf::ctx::ctx(), + ) + .is_none(), "policy reload must prevent a stale relay from starting" ); } @@ -485,7 +540,14 @@ mod tests { let (_upstream_peer, mut proxy_upstream) = tokio::io::duplex(64); let relay = tokio::spawn(async move { - relay_tcp(&mut proxy_client, &mut proxy_upstream, &guard, &request).await + relay_tcp( + &mut proxy_client, + &mut proxy_upstream, + &guard, + &request, + openshell_ocsf::ctx::ctx(), + ) + .await }); tokio::task::yield_now().await;