Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions architecture/sandbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
136 changes: 124 additions & 12 deletions crates/openshell-supervisor-network/src/l7/rest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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";
Expand All @@ -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]
Expand Down
49 changes: 32 additions & 17 deletions crates/openshell-supervisor-network/src/l7/rest/http_response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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).
Expand Down Expand Up @@ -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<C>(client: &mut C, close: bool) -> Result<RelayOutcome>
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)]
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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(
Expand Down
95 changes: 95 additions & 0 deletions crates/openshell-supervisor-network/src/proxy/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -374,6 +374,101 @@ mod tests {
}
}

async fn assert_response_lifecycle<C, P>(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();
Expand Down
5 changes: 5 additions & 0 deletions docs/reference/policy-schema.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading