diff --git a/architecture/sandbox.md b/architecture/sandbox.md index 2f9d88a1fb..366e0b12fc 100644 --- a/architecture/sandbox.md +++ b/architecture/sandbox.md @@ -277,6 +277,12 @@ the shared raw byte relay after the existing adapter gates. Forward HTTP retains its guarded single-request relay while sharing authorization, request context, policy-pinning, and destination boundaries. Adapter-specific response and OCSF event shapes remain at the protocol boundary. +HTTP response framing and connection persistence are separate decisions. After +forwarding a complete closing response (explicit `Connection: close` or HTTP/1.0 +without keep-alive), the relay flushes and shuts down downstream writes before +ending the exchange, including TLS close notification. Response middleware +preserves this lifetime rule; persistent responses remain eligible for reuse. + An explicit `protocol: tcp` endpoint with a valid DNS hostname opts into native DNS and transparent TCP when the selected runtime advertises that substrate. Hostless `allowed_ips` and literal-IP selectors remain available only to the diff --git a/crates/openshell-supervisor-network/src/l7/rest.rs b/crates/openshell-supervisor-network/src/l7/rest.rs index 399db42301..4592600231 100644 --- a/crates/openshell-supervisor-network/src/l7/rest.rs +++ b/crates/openshell-supervisor-network/src/l7/rest.rs @@ -3663,7 +3663,7 @@ mod tests { const VALID_WS_ACCEPT: &str = "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="; const TEXT_OPCODE: u8 = 0x1; - #[derive(Clone, Copy)] + #[derive(Debug, Clone, Copy)] enum ResponseRelayScript { HeadersOnly, WholeBody, @@ -7154,7 +7154,7 @@ mod tests { ResponseRelayScript::HeadersOnly, ) .await; - assert!(matches!(outcome.unwrap(), RelayOutcome::Reusable)); + assert!(matches!(outcome.unwrap(), RelayOutcome::Consumed)); let (outcome, delivered) = run_response_middleware_relay( b"HTTP/1.1 200 OK\r\n\r\n", @@ -7931,6 +7931,117 @@ mod tests { assert!(received_str.contains("hello")); } + #[tokio::test] + async fn response_persistence_is_independent_of_framing_and_middleware() { + for script in [ + None, + Some(ResponseRelayScript::HeadersOnly), + Some(ResponseRelayScript::Stream), + Some(ResponseRelayScript::WholeBody), + ] { + for (method, response, closes) in [ + ( + "GET", + "HTTP/1.0 200 OK\r\nContent-Length: 5\r\n\r\nhello", + true, + ), + ( + "GET", + "HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Length: 5\r\n\r\nhello", + true, + ), + ( + "GET", + "HTTP/1.1 200 OK\r\nConnection: close\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhello\r\n0\r\n\r\n", + true, + ), + ("HEAD", "HTTP/1.0 200 OK\r\nContent-Length: 5\r\n\r\n", true), + ("GET", "HTTP/1.0 204 No Content\r\n\r\n", true), + ( + "GET", + "HTTP/1.1 304 Not Modified\r\nConnection: close\r\n\r\n", + true, + ), + ( + "GET", + "HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhello", + false, + ), + ( + "GET", + "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhello\r\n0\r\n\r\n", + false, + ), + ( + "HEAD", + "HTTP/1.0 200 OK\r\nConnection: keep-alive\r\nContent-Length: 5\r\n\r\n", + false, + ), + ] { + // Body-inspecting scripts reject bodiless responses at preflight. + if (method == "HEAD" + || response.contains("204 No Content") + || response.contains("304 Not Modified")) + && matches!( + script, + Some(ResponseRelayScript::Stream | ResponseRelayScript::WholeBody) + ) + { + continue; + } + let fixture = script.map(response_middleware_fixture); + let middleware = fixture + .as_ref() + .map(|(runner, chain)| response_middleware_context(runner, chain, method)); + let mut upstream = response.as_bytes(); + let (mut reader, mut writer) = tokio::io::duplex(4096); + let outcome = relay_response( + method, + &mut upstream, + &mut writer, + RelayResponseOptions::default(), + middleware, + ) + .await + .unwrap(); + assert_eq!( + matches!(outcome, RelayOutcome::Consumed), + closes, + "{response}" + ); + let mut delivered = Vec::new(); + if closes { + tokio::time::timeout( + std::time::Duration::from_secs(2), + reader.read_to_end(&mut delivered), + ) + .await + .unwrap_or_else(|error| panic!("{script:?} {method} {response:?}: {error}")) + .unwrap(); + } else { + // A second write demonstrates that persistent output is still open. + writer.write_all(b"next response").await.unwrap(); + writer.shutdown().await.unwrap(); + reader.read_to_end(&mut delivered).await.unwrap(); + assert!(delivered.ends_with(b"next response")); + } + if let Some(script) = script { + let expected = match script { + ResponseRelayScript::HeadersOnly => "hello", + ResponseRelayScript::Stream => "HELLO", + ResponseRelayScript::WholeBody => "whole:hello", + _ => unreachable!(), + }; + if response.ends_with("hello") { + assert!(String::from_utf8_lossy(&delivered).contains(expected)); + } + } else { + assert!(delivered.starts_with(response.as_bytes())); + } + } + } + } + #[tokio::test] async fn relay_response_connection_close_with_content_length() { let response = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\nhello"; @@ -7956,18 +8067,19 @@ mod tests { .expect("relay must not deadlock"); let outcome = result.expect("relay_response should succeed"); - // With explicit framing, Connection: close is still reported as reusable - // so the relay loop continues. The *next* upstream write will fail and - // exit the loop via the normal error path. - assert!( - matches!(outcome, RelayOutcome::Reusable), - "explicit framing keeps loop alive despite Connection: close" - ); + assert!(matches!(outcome, RelayOutcome::Consumed)); - client_write.shutdown().await.unwrap(); + // Keep the relay-side stream alive: EOF must come from shutdown, + // not from dropping the stream or the test closing it manually. let mut received = Vec::new(); - client_read.read_to_end(&mut received).await.unwrap(); - assert!(String::from_utf8_lossy(&received).contains("hello")); + tokio::time::timeout( + std::time::Duration::from_secs(2), + client_read.read_to_end(&mut received), + ) + .await + .expect("closing response must deliver EOF") + .unwrap(); + assert_eq!(received, response); } #[tokio::test] diff --git a/crates/openshell-supervisor-network/src/l7/rest/http_response.rs b/crates/openshell-supervisor-network/src/l7/rest/http_response.rs index 9a04244282..8b637ef263 100644 --- a/crates/openshell-supervisor-network/src/l7/rest/http_response.rs +++ b/crates/openshell-supervisor-network/src/l7/rest/http_response.rs @@ -178,7 +178,11 @@ where )) .await? { - return Ok(outcome); + return if matches!(outcome, RelayOutcome::Consumed) { + finish_response(client, true).await + } else { + Ok(outcome) + }; } // Bodiless responses (HEAD, 1xx, 204, 304): forward headers only, skip body @@ -187,12 +191,7 @@ where .write_all(&buf[..header_end]) .await .into_diagnostic()?; - client.flush().await.into_diagnostic()?; - return if server_wants_close { - Ok(RelayOutcome::Consumed) - } else { - Ok(RelayOutcome::Reusable) - }; + return finish_response(client, server_wants_close || http_10_closes_by_default).await; } // No explicit framing (no Content-Length, no Transfer-Encoding). @@ -261,13 +260,24 @@ where "relay_response complete (explicit framing)" ); - // When body framing is explicit (Content-Length / Chunked), always report - // the connection as reusable so the relay loop continues. If the server - // sent `Connection: close`, the *next* upstream write will fail and the - // loop will exit via the normal error path. Exiting early here would - // tear down the CONNECT tunnel before the client can detect the close, - // causing ~30 s retry delays in clients like `gh`. - Ok(RelayOutcome::Reusable) + finish_response(client, server_wants_close || http_10_closes_by_default).await +} + +/// Body framing determines when delivery finishes, not whether another +/// request is permitted. Signal EOF (including TLS `close_notify`) before the +/// caller tears down a closing CONNECT tunnel; waiting for another request +/// deadlocks clients that are themselves waiting for EOF. +async fn finish_response(client: &mut C, close: bool) -> Result +where + C: AsyncWrite + Unpin, +{ + client.flush().await.into_diagnostic()?; + if close { + client.shutdown().await.into_diagnostic()?; + Ok(RelayOutcome::Consumed) + } else { + Ok(RelayOutcome::Reusable) + } } #[allow(clippy::too_many_arguments)] @@ -758,8 +768,9 @@ where } client.flush().await.into_diagnostic()?; Ok(Some( - if (committed && close_delimited_output) - || (matches!(body_length, BodyLength::None) && (server_wants_close || event_stream)) + if server_wants_close + || (committed && close_delimited_output) + || (matches!(body_length, BodyLength::None) && event_stream) { RelayOutcome::Consumed } else { @@ -827,7 +838,11 @@ where BodyLength::None => {} } client.flush().await.into_diagnostic()?; - Ok(RelayOutcome::Reusable) + Ok(if server_wants_close { + RelayOutcome::Consumed + } else { + RelayOutcome::Reusable + }) } fn emit_http_response_middleware_invocations( diff --git a/crates/openshell-supervisor-network/src/proxy/relay.rs b/crates/openshell-supervisor-network/src/proxy/relay.rs index 4c89e5e382..77055245a5 100644 --- a/crates/openshell-supervisor-network/src/proxy/relay.rs +++ b/crates/openshell-supervisor-network/src/proxy/relay.rs @@ -374,6 +374,101 @@ mod tests { } } + async fn assert_response_lifecycle(mut caller: C, mut client: P) + where + C: AsyncRead + AsyncWrite + Unpin, + P: AsyncRead + AsyncWrite + Unpin + Send, + { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let (mut upstream, mut server) = tokio::io::duplex(1024); + let engine = OpaEngine::from_strings(POLICY_REGO, EMPTY_POLICY_DATA).unwrap(); + let decision = decision(engine.current_generation()); + let request = request_context(); + let context = prepare_http_relay(None, &engine, &decision, &request).unwrap(); + let persistent = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhello"; + // Larger than either transport buffer: closing must drain the body. + let body = vec![b'x'; 64 * 1024]; + let mut closing = format!( + "HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + body.len() + ) + .into_bytes(); + closing.extend_from_slice(&body); + let relay = Box::pin(relay_http_stream(&mut client, &mut upstream, context)); + let serve = async { + for response in [persistent.as_slice(), closing.as_slice()] { + let mut request = Vec::new(); + while !request.ends_with(b"\r\n\r\n") { + request.push(server.read_u8().await.unwrap()); + assert!(request.len() < 4096); + } + assert!(request.starts_with(b"GET / HTTP/1.1\r\n")); + server.write_all(response).await.unwrap(); + server.flush().await.unwrap(); + } + // Retain the upstream socket: response headers decide persistence. + }; + let receive = async { + let request = b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n"; + caller.write_all(request).await.unwrap(); + caller.flush().await.unwrap(); + let mut first = vec![0; persistent.len()]; + caller.read_exact(&mut first).await.unwrap(); + assert_eq!(first, persistent); + // A real second exchange verifies reuse, without a timing assertion. + caller.write_all(request).await.unwrap(); + caller.flush().await.unwrap(); + let mut second = Vec::new(); + // TLS must return clean EOF, not UnexpectedEof from a dropped socket. + caller.read_to_end(&mut second).await.unwrap(); + assert_eq!(second, closing); + }; + tokio::time::timeout(std::time::Duration::from_secs(5), async { + let (result, (), ()) = tokio::join!(relay, serve, receive); + result.unwrap(); + }) + .await + .expect("response delivery and EOF must not await another request"); + } + + #[tokio::test] + async fn http_relay_reuses_then_closes_after_complete_response() { + let (caller, client) = tokio::io::duplex(1024); + assert_response_lifecycle(caller, client).await; + } + + #[tokio::test] + async fn tls_http_relay_reuses_then_sends_close_notify_after_complete_response() { + use crate::l7::tls::{CertCache, ProxyTlsState, SandboxCa, tls_terminate_client}; + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + let ca = SandboxCa::generate().unwrap(); + let mut roots = rustls::RootCertStore::empty(); + for cert in rustls_pemfile::certs(&mut ca.cert_pem().as_bytes()) { + roots.add(cert.unwrap()).unwrap(); + } + let config = Arc::new( + rustls::ClientConfig::builder() + .with_root_certificates(roots) + .with_no_client_auth(), + ); + let state = ProxyTlsState::new(CertCache::new(ca), config.clone()); + let connector = tokio_rustls::TlsConnector::from(config); + let (caller, client) = tokio::io::duplex(1024); + let (caller, client) = tokio::time::timeout(std::time::Duration::from_secs(5), async { + tokio::join!( + connector.connect( + rustls::pki_types::ServerName::try_from("example.com").unwrap(), + caller + ), + tls_terminate_client(client, &state, "example.com"), + ) + }) + .await + .expect("TLS handshake must complete"); + assert_response_lifecycle(caller.unwrap(), client.unwrap()).await; + } + #[test] fn relay_without_route_pins_l4_decision_generation() { let engine = OpaEngine::from_strings(POLICY_REGO, EMPTY_POLICY_DATA).unwrap(); diff --git a/docs/reference/policy-schema.mdx b/docs/reference/policy-schema.mdx index 3e1539b7b0..a72a21bb2c 100644 --- a/docs/reference/policy-schema.mdx +++ b/docs/reference/policy-schema.mdx @@ -267,6 +267,11 @@ access selects no preset. Unknown enum numbers are rejected before activation. - `deny_rules: []` (empty list) is rejected; remove it if no denials are needed. - `credential_signing` requires a resolvable AWS credential source before a sandbox policy can activate. Use an attached endpoint-bearing profile that declares `AWS_ACCESS_KEY_ID` and `AWS_SECRET_ACCESS_KEY` and covers the signed endpoint, or bind an attached endpointless profile that declares those keys with `credential_binding.provider`. +For proxied HTTP responses, OpenShell completes delivery and signals EOF when +the server sends `Connection: close` or an HTTP/1.0 response without keep-alive. +This also applies to HTTPS and responses processed by middleware. A +`Content-Length` or chunked body does not override the server's close decision. + Credential rewrite recognizes the canonical `openshell:resolve:env:KEY` placeholder form and whole-token provider-shaped aliases such as `provider-OPENSHELL-RESOLVE-ENV-API_TOKEN` when the referenced environment key exists in the configured provider credentials. Static provider placeholders also require the request host, port, and path to diff --git a/e2e/rust/tests/host_gateway_alias.rs b/e2e/rust/tests/host_gateway_alias.rs index e4ae811472..8d355f97fb 100644 --- a/e2e/rust/tests/host_gateway_alias.rs +++ b/e2e/rust/tests/host_gateway_alias.rs @@ -340,6 +340,78 @@ async fn sandbox_reaches_host_openshell_internal_via_host_gateway_alias() { ); } +#[tokio::test] +async fn sandbox_receives_eof_after_closing_http_response() { + for response in [ + "HTTP/1.0 200 OK\r\nContent-Length: 3\r\n\r\nOK\n", + "HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Length: 3\r\n\r\nOK\n", + "HTTP/1.1 200 OK\r\nConnection: close\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nOK\n\r\n0\r\n\r\n", + ] { + let listener = TcpListener::bind(("0.0.0.0", 0)).await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let server = HostServer { + port, + task: tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.unwrap(); + stream.set_nodelay(true).expect("disable Nagle on fixture"); + let mut request = Vec::new(); + while !request.ends_with(b"\r\n\r\n") { + request.push(stream.read_u8().await.unwrap()); + assert!(request.len() < 4096); + } + stream.write_all(response.as_bytes()).await.unwrap(); + stream.shutdown().await.unwrap(); + }), + }; + let policy = write_policy(server.port).unwrap(); + let command = format!( + r#"set -eu +exec 3<>/dev/tcp/host.openshell.internal/{port} +printf 'GET / HTTP/1.1\r\nHost: host.openshell.internal:{port}\r\n\r\n' >&3 +while true; do + line= + if IFS= read -r -t 5 line <&3; then + printf '%s\n' "$line" + else + status=$? + [ "$status" -eq 1 ] || {{ echo EOF_TIMEOUT; exit 1; }} + [ -z "$line" ] || printf '%s\n' "$line" + break + fi +done +printf 'RESPONSE_EOF\n' +"#, + ); + let mut sandbox = SandboxGuard::create(&[ + "--policy", + policy.path().to_str().unwrap(), + "--no-auto-providers", + "--", + "/usr/bin/bash", + "-c", + &command, + ]) + .await + .expect("closing response must finish without a client timeout"); + assert!( + sandbox.create_output.lines().any(|line| line == "OK"), + "{}", + sandbox.create_output + ); + assert!( + sandbox.create_output.contains("RESPONSE_EOF"), + "{}", + sandbox.create_output + ); + assert!( + !sandbox.create_output.contains("EOF_TIMEOUT"), + "{}", + sandbox.create_output + ); + sandbox.cleanup().await; + } +} + #[tokio::test] async fn static_provider_credentials_are_bound_to_profile_endpoints() { let server = HostServer::start_with_auth_check("", Some("Bearer e2e-bound-secret")) diff --git a/e2e/rust/tests/proxy_egress_pipeline.rs b/e2e/rust/tests/proxy_egress_pipeline.rs index 9556c16429..09eda47615 100644 --- a/e2e/rust/tests/proxy_egress_pipeline.rs +++ b/e2e/rust/tests/proxy_egress_pipeline.rs @@ -614,10 +614,9 @@ impl PipelineProbeServer { } } observed.lock().unwrap().extend_from_slice(&request); + // Keep the first response reusable so the queued request reaches policy evaluation. let _ = stream - .write_all( - b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok", - ) + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok") .await; }); }