From 1ae4c1d0ad217f0b2956e1c3ff1933e2898a1df6 Mon Sep 17 00:00:00 2001 From: Drew Newberry Date: Sat, 26 Sep 2026 17:11:53 -0700 Subject: [PATCH 1/3] fix(sandbox): charge control slots only after bearer authentication The sandbox control listener requires no TLS client certificate; peers prove themselves with the EdDSA session bearer on each RPC. A connection that completed TLS was nevertheless charged one of the 128 control slots and held it indefinitely, so a same-UID workload process reaching the Unix socket could exhaust the capacity the supervisor needs to reconnect. - Acquire the control slot on the first authenticated RPC of a connection and return RESOURCE_EXHAUSTED when none is available. - Arm a 10-second deadline on every connection until a bearer is accepted, and make the gRPC transport report EOF on shutdown so a peer that never sent an HTTP/2 preface is actually torn down. - Admit same-namespace Unix peers only when they are the sandbox or one of its ancestors, so an orphan reparented to PID 1 or a subreaper is rejected even when the sandbox is not PID 1. - Correct comments that described the channel as mutual TLS and document the trust model in the sandbox architecture doc. Signed-off-by: Drew Newberry --- architecture/sandbox.md | 10 +- crates/openshell-driver-docker/src/lib.rs | 7 +- .../openshell-sandbox/src/boundary_server.rs | 341 ++++++++++++++++-- 3 files changed, 327 insertions(+), 31 deletions(-) diff --git a/architecture/sandbox.md b/architecture/sandbox.md index b1169ca64e..1696f43ffb 100644 --- a/architecture/sandbox.md +++ b/architecture/sandbox.md @@ -23,7 +23,15 @@ or backend-admin authority. The compute driver provisions separate protected configurations and one mutually authenticated gRPC connection over a private Unix socket, Kubernetes -TCP Service, or VM vsock channel. Independent bidirectional `Exchange` RPCs +TCP Service, or VM vsock channel. The sandbox listener presents a server +certificate only; the supervisor authenticates every RPC with an EdDSA session +bearer bound to the sandbox, runtime generation, and credential epoch. The +Unix socket is reachable by same-UID workload processes because Landlock does +not govern `connect()` on a filesystem socket, so the sandbox rejects peers +inside its own PID namespace that are not the sandbox or one of its ancestors +before TLS, charges a bounded control-connection slot only after the first +accepted bearer, and closes connections that present no valid bearer within a +short deadline. Independent bidirectional `Exchange` RPCs carry lifecycle, exec, TCP, and forwarding traffic, while one persistent bidirectional `Mediate` RPC carries multiplexed DNS traffic. General application UDP is unsupported; UDP DNS remains mediated by the supervisor. diff --git a/crates/openshell-driver-docker/src/lib.rs b/crates/openshell-driver-docker/src/lib.rs index 1f4c91f5e4..cab27cfa77 100644 --- a/crates/openshell-driver-docker/src/lib.rs +++ b/crates/openshell-driver-docker/src/lib.rs @@ -4389,8 +4389,11 @@ fn docker_sandbox_bundle_archive( ".openshell/channel/sandbox", // The sandbox owns this directory so it can consume bootstrap files // and create the control socket. The separate non-root supervisor - // needs execute-only traversal to that known socket path; mutual TLS - // authenticates the endpoint and the files beneath remain 0600. + // needs execute-only traversal to that known socket path. The socket + // presents a server certificate only; the supervisor proves itself + // with the EdDSA session bearer on each RPC, and the sandbox rejects + // same-namespace workload peers before TLS. Files beneath stay 0600 + // and are consumed at startup. 0o711, identity.uid, identity.gid, diff --git a/crates/openshell-sandbox/src/boundary_server.rs b/crates/openshell-sandbox/src/boundary_server.rs index 7a2a73f731..2d224ec304 100644 --- a/crates/openshell-sandbox/src/boundary_server.rs +++ b/crates/openshell-sandbox/src/boundary_server.rs @@ -80,6 +80,12 @@ mod linux { const FORCE_KILL_REAP_TIMEOUT: Duration = Duration::from_secs(2); const MAX_PENDING_HANDSHAKES: usize = 32; const MAX_CONTROL_CONNECTIONS: usize = 128; + /// A completed TLS handshake proves nothing about the peer because the + /// listener requires no client certificate. A connection that has not + /// presented a valid bearer holds no control slot and closes at this + /// deadline, so a same-UID workload that reaches the Unix socket cannot + /// pin the resources the supervisor needs to reconnect. + const CONTROL_UNAUTHENTICATED_DEADLINE: Duration = Duration::from_secs(10); const MAX_REPLAY_LEDGER_ENTRIES: usize = 4096; const MAX_RETAINED_EXEC_PROCESSES: usize = 64; @@ -533,6 +539,7 @@ mod linux { runtime, pending, active_connections, + CONTROL_UNAUTHENTICATED_DEADLINE, ) .await { @@ -564,26 +571,36 @@ mod linux { runtime: Arc, pending: tokio::sync::OwnedSemaphorePermit, active_connections: Arc, + unauthenticated_deadline: Duration, ) -> Result<(), String> { let stream = stream .establish_async(&runtime.process_runtime) .await - .map_err(|error| format!("authenticate boundary transport: {error}"))?; + .map_err(|error| format!("establish boundary transport: {error}"))?; drop(pending); - let Some(_slot) = acquire_control_connection_slot(&active_connections) else { - return Err("authenticated control connection limit reached".to_string()); - }; - serve_grpc(stream.into_tokio()?, runtime, SandboxConnectionId::new()).await + // TLS completion is not authentication. The control slot is acquired + // only after an RPC on this connection presents a valid bearer. + serve_grpc( + stream.into_tokio()?, + runtime, + SandboxConnectionId::new(), + active_connections, + unauthenticated_deadline, + ) + .await } async fn serve_grpc( stream: openshell_isolation_interface::contract::BoundaryDuplexStream, runtime: Arc, connection_id: SandboxConnectionId, + active_connections: Arc, + unauthenticated_deadline: Duration, ) -> Result<(), String> { let (connection_shutdown, connection_closed) = tokio::sync::watch::channel(()); runtime.register_connection(connection_id, connection_shutdown.clone()); let connection_expiry = Arc::new(ConnectionExpiry::new(connection_shutdown.clone())); + connection_expiry.update_deadline(tokio::time::Instant::now() + unauthenticated_deadline); let incoming = tokio_stream::StreamExt::chain( tokio_stream::iter([Ok::<_, io::Error>(GrpcServerIo { stream, @@ -592,6 +609,13 @@ mod linux { runtime: Arc::downgrade(&runtime), connection_id, }, + closed: Box::pin({ + let mut closed = connection_closed.clone(); + async move { + let _ = closed.changed().await; + } + }), + is_closed: false, })]), tokio_stream::pending(), ); @@ -608,6 +632,8 @@ mod linux { connection_id, connection_expiry, connection_closed, + active_connections, + control_slot: Arc::new(Mutex::new(None)), }) .max_decoding_message_size(64 * 1024) .max_encoding_message_size(64 * 1024), @@ -626,6 +652,21 @@ mod linux { // bridges, including on keepalive failure or task cancellation. _connection_alive: tokio::sync::watch::Sender<()>, _disconnect: TransportDisconnectGuard, + /// Resolves once the connection deadline or an explicit shutdown + /// fires. Graceful HTTP/2 shutdown alone does not tear down a peer + /// that completed TLS but never sent a preface, so the transport + /// reports EOF itself. + closed: Pin + Send>>, + is_closed: bool, + } + + impl GrpcServerIo { + fn poll_closed(&mut self, context: &mut Context<'_>) -> bool { + if !self.is_closed && self.closed.as_mut().poll(context).is_ready() { + self.is_closed = true; + } + self.is_closed + } } struct TransportDisconnectGuard { @@ -647,6 +688,9 @@ mod linux { context: &mut Context<'_>, buffer: &mut tokio::io::ReadBuf<'_>, ) -> Poll> { + if self.poll_closed(context) { + return Poll::Ready(Ok(())); + } Pin::new(&mut self.stream).poll_read(context, buffer) } } @@ -657,6 +701,12 @@ mod linux { context: &mut Context<'_>, buffer: &[u8], ) -> Poll> { + if self.poll_closed(context) { + return Poll::Ready(Err(io::Error::new( + io::ErrorKind::ConnectionAborted, + "boundary control connection closed", + ))); + } Pin::new(&mut self.stream).poll_write(context, buffer) } @@ -684,6 +734,45 @@ mod linux { connection_id: SandboxConnectionId, connection_expiry: Arc, connection_closed: tokio::sync::watch::Receiver<()>, + active_connections: Arc, + /// Held from the first authenticated RPC until the connection ends. + /// Shared across tonic's per-request service clones. + control_slot: Arc>>, + } + + impl GrpcBoundaryService { + /// Authenticate one RPC and, on the first success for this + /// connection, charge it against the bounded control slots. + /// + /// Unauthenticated connections never hold a slot, so a peer that can + /// only complete TLS cannot exhaust the supervisor's reconnect + /// capacity. The unauthenticated deadline stays armed until a bearer + /// is accepted here. + fn authenticate_and_admit( + &self, + metadata: &tonic::metadata::MetadataMap, + ) -> Result { + let principal = self + .runtime + .authenticate_request(self.connection_id, metadata)?; + { + let mut slot = lock(&self.control_slot); + if slot.is_none() { + *slot = Some( + acquire_control_connection_slot(&self.active_connections).ok_or_else( + || { + tonic::Status::resource_exhausted( + "authenticated control connection limit reached", + ) + }, + )?, + ); + } + } + self.connection_expiry + .update(principal.session().expires_at); + Ok(principal) + } } struct ConnectionExpiry { @@ -772,11 +861,7 @@ mod linux { &self, request: tonic::Request>, ) -> Result, tonic::Status> { - let principal = self - .runtime - .authenticate_request(self.connection_id, request.metadata())?; - self.connection_expiry - .update(principal.session().expires_at); + let principal = self.authenticate_and_admit(request.metadata())?; let (stream, response) = bridge_grpc_server_stream(request.into_inner(), self.connection_closed.clone()); let runtime = self.runtime.clone(); @@ -796,11 +881,7 @@ mod linux { &self, request: tonic::Request>, ) -> Result, tonic::Status> { - let principal = self - .runtime - .authenticate_request(self.connection_id, request.metadata())?; - self.connection_expiry - .update(principal.session().expires_at); + let principal = self.authenticate_and_admit(request.metadata())?; let (stream, response) = bridge_grpc_server_stream(request.into_inner(), self.connection_closed.clone()); let runtime = self.runtime.clone(); @@ -3150,8 +3231,14 @@ mod linux { BoundaryListenerConfig::Unix { socket_path, tls } => { remove_owned_stale_control_socket(socket_path)?; let listener = std::os::unix::net::UnixListener::bind(socket_path)?; - // Mutual TLS makes a same-UID pathname replacement a - // detectable denial of service rather than impersonation. + // The workload shares this UID and Landlock does not + // govern connect() on a filesystem socket, so any workload + // process can reach this listener. The listener presents + // a server certificate only; peers prove themselves with + // the EdDSA session bearer on each RPC, and the peer + // credential check in accept() rejects workload processes + // before TLS. A same-UID pathname replacement is therefore + // a detectable denial of service, not impersonation. std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o666))?; listener.set_nonblocking(true)?; let server_config = Arc::new(load_tls_server_config(tls)?); @@ -3300,11 +3387,17 @@ mod linux { } let peer = u32::try_from(credentials.pid) .map_err(|_| io::Error::from_raw_os_error(libc::EACCES))?; - // Linux reports PID zero for a peer outside our PID namespace. Such a - // peer still must authenticate with the per-sandbox mTLS certificate. + // Linux reports PID zero for a peer outside our PID namespace, which + // is where every driver places the supervisor. Such a peer still must + // present the EdDSA session bearer on each RPC. Inside our namespace, + // only this process and its ancestors are trusted: workload processes + // are descendants, and an orphan that double-forked away from the + // sandbox reparents to PID 1 or a subreaper, which is never an + // ancestor of the sandbox unless the sandbox is PID 1 itself. Read + // only kernel-owned ancestry, never workload-supplied data. if peer != 0 && peer != std::process::id() - && is_process_descendant(peer, std::process::id()) + && !is_process_descendant(std::process::id(), peer) .map_err(|_| io::Error::from_raw_os_error(libc::EACCES))? { return Err(io::Error::from_raw_os_error(libc::EACCES)); @@ -3313,10 +3406,7 @@ mod linux { } fn is_process_descendant(mut process: u32, ancestor: u32) -> io::Result { - // Drivers run the sandbox as workload PID 1, so orphaned descendants - // reparent to it and cannot escape this check by double-forking. - // Read kernel-owned ancestry, never workload-supplied paths or UIDs. - // If a peer exits during inspection, fail closed for that connection. + // If a process exits during inspection, fail closed for that connection. for _ in 0..1024 { if process == ancestor { return Ok(true); @@ -3956,8 +4046,8 @@ mod linux { drop(stream); assert!(child.wait().unwrap().success()); assert!(!is_process_descendant(std::process::id(), child.id()).unwrap()); - // Trusted same-process connections and external ancestors remain - // eligible for mTLS; we do not equate same UID with workload trust. + // Same-process connections and ancestors of the sandbox remain + // eligible to present a bearer; same UID alone is not workload trust. let client = std::os::unix::net::UnixStream::connect(&path).unwrap(); let (stream, _) = listener.accept().unwrap(); reject_workload_unix_peer(&stream).unwrap(); @@ -4068,6 +4158,7 @@ mod linux { runtime.clone(), permit, active.clone(), + CONTROL_UNAUTHENTICATED_DEADLINE, )); } assert!(pending.clone().try_acquire_owned().is_err()); @@ -4084,6 +4175,186 @@ mod linux { drop(clients); } + #[tokio::test(flavor = "multi_thread")] + async fn unauthenticated_tls_connection_holds_no_slot_and_closes_at_deadline() { + let directory = tempfile::tempdir().unwrap(); + let (server_tls, client_tls) = stage_test_tls(directory.path(), "unauthenticated"); + let server_config = Arc::new(load_tls_server_config(&server_tls).unwrap()); + let (runtime, _) = availability_test_runtime(); + let active = Arc::new(AtomicUsize::new(0)); + let pending = Arc::new(tokio::sync::Semaphore::new(MAX_PENDING_HANDSHAKES)); + let (server, client) = std::os::unix::net::UnixStream::pair().unwrap(); + client.set_nonblocking(true).unwrap(); + let client = tokio::net::UnixStream::from_std(client).unwrap(); + let permit = pending.clone().try_acquire_owned().unwrap(); + let deadline = Duration::from_secs(2); + let started = std::time::Instant::now(); + let server = tokio::spawn(serve_control_connection( + ControlStream::PendingTls { + stream: PlainControlStream::Unix(server), + server_config, + }, + runtime, + permit, + active.clone(), + deadline, + )); + + // A same-UID workload peer can complete TLS: the listener asks + // for no client certificate and never sees a bearer here. + let client_config = test_client_config(&client_tls); + let server_name = rustls::pki_types::ServerName::try_from(client_tls.server_name) + .expect("valid server name"); + let mut stream = tokio_rustls::TlsConnector::from(Arc::new(client_config)) + .connect(server_name, client) + .await + .expect("TLS completes without any bearer"); + tokio::time::timeout(Duration::from_secs(2), async { + while pending.available_permits() != MAX_PENDING_HANDSHAKES { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("handshake permit is released once TLS completes"); + assert_eq!( + active.load(Ordering::Acquire), + 0, + "TLS completion must not consume a control slot" + ); + + // The server speaks HTTP/2 first (SETTINGS), so drain until the + // peer closes rather than treating the first bytes as the end. + tokio::time::timeout(deadline + Duration::from_secs(3), async { + let mut buffer = [0_u8; 256]; + loop { + match stream.read(&mut buffer).await { + Ok(0) | Err(_) => break, + Ok(_) => {} + } + } + }) + .await + .expect("server must close the idle unauthenticated connection at the deadline"); + assert!( + started.elapsed() >= deadline, + "connection must survive until the unauthenticated deadline" + ); + assert_eq!(active.load(Ordering::Acquire), 0); + let _ = tokio::time::timeout(Duration::from_secs(3), server) + .await + .expect("server task must finish after the deadline"); + assert_eq!(active.load(Ordering::Acquire), 0); + } + + fn discover_policy_stream() -> tokio_stream::Iter> { + let request = + RequestEnvelope::new(Request::DiscoverPolicy).expect("encode discover request"); + tokio_stream::iter([BoundaryChunk { + data: encode_frame(&request).expect("encode logical request"), + }]) + } + + #[tokio::test(flavor = "multi_thread")] + async fn control_slot_is_charged_only_after_a_valid_bearer() { + let (runtime, token) = availability_test_runtime(); + let active = Arc::new(AtomicUsize::new(0)); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server_active = active.clone(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + serve_grpc( + Box::new(stream), + runtime, + SandboxConnectionId::new(), + server_active, + Duration::from_secs(30), + ) + .await + }); + let channel = tonic::transport::Endpoint::from_shared(format!("http://{address}")) + .unwrap() + .connect() + .await + .unwrap(); + let mut client = IsolationBoundaryClient::new(channel); + + let status = client + .exchange(tonic::Request::new(discover_policy_stream())) + .await + .expect_err("missing bearer must be rejected"); + assert_eq!(status.code(), tonic::Code::Unauthenticated); + assert_eq!(active.load(Ordering::Acquire), 0); + + let status = client + .exchange(bearer_request( + discover_policy_stream(), + "not-a-session-token", + )) + .await + .expect_err("malformed bearer must be rejected"); + assert_eq!(status.code(), tonic::Code::Unauthenticated); + assert_eq!(active.load(Ordering::Acquire), 0); + + for _ in 0..2 { + let mut body = client + .exchange(bearer_request(discover_policy_stream(), &token)) + .await + .expect("valid bearer is admitted") + .into_inner(); + while body.message().await.unwrap().is_some() {} + assert_eq!( + active.load(Ordering::Acquire), + 1, + "one authenticated connection charges exactly one slot" + ); + } + + drop(client); + let _ = tokio::time::timeout(Duration::from_secs(5), server) + .await + .expect("server ends when the authenticated client disconnects"); + assert_eq!( + active.load(Ordering::Acquire), + 0, + "the slot is released with the connection" + ); + } + + #[tokio::test(flavor = "multi_thread")] + async fn authenticated_rpc_is_refused_when_control_slots_are_exhausted() { + let (runtime, token) = availability_test_runtime(); + let active = Arc::new(AtomicUsize::new(MAX_CONTROL_CONNECTIONS)); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server_active = active.clone(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + serve_grpc( + Box::new(stream), + runtime, + SandboxConnectionId::new(), + server_active, + Duration::from_secs(30), + ) + .await + }); + let channel = tonic::transport::Endpoint::from_shared(format!("http://{address}")) + .unwrap() + .connect() + .await + .unwrap(); + let mut client = IsolationBoundaryClient::new(channel); + let status = client + .exchange(bearer_request(discover_policy_stream(), &token)) + .await + .expect_err("exhausted slots must refuse even a valid bearer"); + assert_eq!(status.code(), tonic::Code::ResourceExhausted); + assert_eq!(active.load(Ordering::Acquire), MAX_CONTROL_CONNECTIONS); + drop(client); + server.abort(); + } + #[tokio::test(flavor = "multi_thread")] async fn disconnected_session_reconfirms_before_becoming_active() { let (runtime, token) = availability_test_runtime(); @@ -4238,7 +4509,14 @@ mod linux { let server_runtime = runtime.clone(); let server = tokio::spawn(async move { let (stream, _) = server_listener.accept().await.unwrap(); - serve_grpc(Box::new(stream), server_runtime, connection_id).await + serve_grpc( + Box::new(stream), + server_runtime, + connection_id, + Arc::new(AtomicUsize::new(0)), + CONTROL_UNAUTHENTICATED_DEADLINE, + ) + .await }); let proxy_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let proxy_address = proxy_listener.local_addr().unwrap(); @@ -4673,7 +4951,14 @@ mod linux { let address = listener.local_addr().expect("gRPC test address"); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("accept gRPC client"); - serve_grpc(Box::new(stream), boundary, SandboxConnectionId::new()).await + serve_grpc( + Box::new(stream), + boundary, + SandboxConnectionId::new(), + Arc::new(AtomicUsize::new(0)), + CONTROL_UNAUTHENTICATED_DEADLINE, + ) + .await }); let channel = tonic::transport::Endpoint::from_shared(format!("http://{address}")) .expect("valid gRPC endpoint") From db5cc767648d1ba99ccbe5fca5919aa41d442eda Mon Sep 17 00:00:00 2001 From: Drew Newberry Date: Sat, 26 Sep 2026 17:42:28 -0700 Subject: [PATCH 2/3] fix(sandbox): bound unauthenticated control connections Signed-off-by: Drew Newberry --- architecture/sandbox.md | 4 +- .../openshell-sandbox/src/boundary_server.rs | 105 +++++++++++++++++- 2 files changed, 107 insertions(+), 2 deletions(-) diff --git a/architecture/sandbox.md b/architecture/sandbox.md index 1696f43ffb..439b169f41 100644 --- a/architecture/sandbox.md +++ b/architecture/sandbox.md @@ -31,7 +31,9 @@ not govern `connect()` on a filesystem socket, so the sandbox rejects peers inside its own PID namespace that are not the sandbox or one of its ancestors before TLS, charges a bounded control-connection slot only after the first accepted bearer, and closes connections that present no valid bearer within a -short deadline. Independent bidirectional `Exchange` RPCs +short deadline. TLS-complete connections awaiting a bearer have a separate +bounded pool; admitting a new connection closes the oldest waiting peer when +that pool is full. Independent bidirectional `Exchange` RPCs carry lifecycle, exec, TCP, and forwarding traffic, while one persistent bidirectional `Mediate` RPC carries multiplexed DNS traffic. General application UDP is unsupported; UDP DNS remains mediated by the supervisor. diff --git a/crates/openshell-sandbox/src/boundary_server.rs b/crates/openshell-sandbox/src/boundary_server.rs index 2d224ec304..0ee92dc689 100644 --- a/crates/openshell-sandbox/src/boundary_server.rs +++ b/crates/openshell-sandbox/src/boundary_server.rs @@ -13,6 +13,7 @@ use std::path::Path; #[cfg(target_os = "linux")] mod linux { + use std::collections::VecDeque; use std::fs::File; use std::io::{self, Read, Write}; use std::mem::size_of; @@ -80,6 +81,7 @@ mod linux { const FORCE_KILL_REAP_TIMEOUT: Duration = Duration::from_secs(2); const MAX_PENDING_HANDSHAKES: usize = 32; const MAX_CONTROL_CONNECTIONS: usize = 128; + const MAX_UNAUTHENTICATED_CONTROL_CONNECTIONS: usize = 128; /// A completed TLS handshake proves nothing about the peer because the /// listener requires no client certificate. A connection that has not /// presented a valid bearer holds no control slot and closes at this @@ -190,6 +192,53 @@ mod linux { .ok() .map(|_| ControlConnectionSlot(active.clone())) } + + struct UnauthenticatedControlConnections { + capacity: usize, + entries: Mutex)>>, + } + + impl UnauthenticatedControlConnections { + fn new(capacity: usize) -> Self { + assert!(capacity > 0); + Self { + capacity, + entries: Mutex::new(VecDeque::new()), + } + } + + fn register( + &self, + connection_id: SandboxConnectionId, + shutdown: tokio::sync::watch::Sender<()>, + ) { + let mut entries = lock(&self.entries); + if entries.len() == self.capacity + && let Some((_, oldest)) = entries.pop_front() + { + // Make room for a new bearer attempt even when every older + // connection is idle. Its transport observes this shutdown. + let _ = oldest.send(()); + } + entries.push_back((connection_id, shutdown)); + } + + fn remove(&self, connection_id: SandboxConnectionId) { + lock(&self.entries).retain(|(id, _)| *id != connection_id); + } + } + + struct UnauthenticatedControlRegistration { + connections: Arc, + connection_id: SandboxConnectionId, + } + + impl Drop for UnauthenticatedControlRegistration { + fn drop(&mut self) { + self.connections.remove(self.connection_id); + } + } + static BOUNDARY_TERMINATION_REQUESTED: AtomicBool = AtomicBool::new(false); static BOUNDARY_TERMINATION_SIGNAL: AtomicI32 = AtomicI32::new(0); @@ -500,6 +549,9 @@ mod linux { let listener = ControlListener::bind(config) .map_err(|error| format!("bind boundary control listener: {error}"))?; let active_connections = Arc::new(AtomicUsize::new(0)); + let unauthenticated_connections = Arc::new(UnauthenticatedControlConnections::new( + MAX_UNAUTHENTICATED_CONTROL_CONNECTIONS, + )); let pending_handshakes = Arc::new(tokio::sync::Semaphore::new(MAX_PENDING_HANDSHAKES)); tracing::info!(?config, "Boundary control listener ready"); loop { @@ -530,6 +582,7 @@ mod linux { continue; }; let active_connections = active_connections.clone(); + let unauthenticated_connections = unauthenticated_connections.clone(); let runtime = runtime.clone(); runtime.process_runtime.spawn({ let runtime = runtime.clone(); @@ -539,6 +592,7 @@ mod linux { runtime, pending, active_connections, + unauthenticated_connections, CONTROL_UNAUTHENTICATED_DEADLINE, ) .await @@ -571,6 +625,7 @@ mod linux { runtime: Arc, pending: tokio::sync::OwnedSemaphorePermit, active_connections: Arc, + unauthenticated_connections: Arc, unauthenticated_deadline: Duration, ) -> Result<(), String> { let stream = stream @@ -585,6 +640,7 @@ mod linux { runtime, SandboxConnectionId::new(), active_connections, + unauthenticated_connections, unauthenticated_deadline, ) .await @@ -595,9 +651,15 @@ mod linux { runtime: Arc, connection_id: SandboxConnectionId, active_connections: Arc, + unauthenticated_connections: Arc, unauthenticated_deadline: Duration, ) -> Result<(), String> { let (connection_shutdown, connection_closed) = tokio::sync::watch::channel(()); + unauthenticated_connections.register(connection_id, connection_shutdown.clone()); + let _unauthenticated_registration = UnauthenticatedControlRegistration { + connections: unauthenticated_connections.clone(), + connection_id, + }; runtime.register_connection(connection_id, connection_shutdown.clone()); let connection_expiry = Arc::new(ConnectionExpiry::new(connection_shutdown.clone())); connection_expiry.update_deadline(tokio::time::Instant::now() + unauthenticated_deadline); @@ -633,6 +695,7 @@ mod linux { connection_expiry, connection_closed, active_connections, + unauthenticated_connections, control_slot: Arc::new(Mutex::new(None)), }) .max_decoding_message_size(64 * 1024) @@ -735,6 +798,7 @@ mod linux { connection_expiry: Arc, connection_closed: tokio::sync::watch::Receiver<()>, active_connections: Arc, + unauthenticated_connections: Arc, /// Held from the first authenticated RPC until the connection ends. /// Shared across tonic's per-request service clones. control_slot: Arc>>, @@ -771,6 +835,7 @@ mod linux { } self.connection_expiry .update(principal.session().expires_at); + self.unauthenticated_connections.remove(self.connection_id); Ok(principal) } } @@ -1468,7 +1533,7 @@ mod linux { #[derive(Default)] struct ReplayLedger { entries: std::collections::HashMap, - order: std::collections::VecDeque, + order: VecDeque, } impl ReplayLedger { @@ -4019,6 +4084,28 @@ mod linux { assert_eq!(active.load(Ordering::Acquire), MAX_CONTROL_CONNECTIONS - 1); } + #[tokio::test] + async fn unauthenticated_connection_pool_evicts_oldest_to_admit_new_peer() { + let connections = UnauthenticatedControlConnections::new(2); + let first_id = SandboxConnectionId::new(); + let second_id = SandboxConnectionId::new(); + let third_id = SandboxConnectionId::new(); + let (first_shutdown, mut first_closed) = tokio::sync::watch::channel(()); + let (second_shutdown, second_closed) = tokio::sync::watch::channel(()); + let (third_shutdown, _third_closed) = tokio::sync::watch::channel(()); + + connections.register(first_id, first_shutdown); + connections.register(second_id, second_shutdown); + connections.register(third_id, third_shutdown); + + first_closed.changed().await.expect("oldest peer is closed"); + assert!(!second_closed.has_changed().unwrap()); + let entries = lock(&connections.entries); + assert_eq!(entries.len(), 2); + assert_eq!(entries[0].0, second_id); + assert_eq!(entries[1].0, third_id); + } + #[test] fn unix_control_rejects_workload_descendants_before_admission() { const CHILD_SOCKET: &str = "OPENSHELL_TEST_CONTROL_PEER_SOCKET"; @@ -4158,6 +4245,9 @@ mod linux { runtime.clone(), permit, active.clone(), + Arc::new(UnauthenticatedControlConnections::new( + MAX_UNAUTHENTICATED_CONTROL_CONNECTIONS, + )), CONTROL_UNAUTHENTICATED_DEADLINE, )); } @@ -4197,6 +4287,9 @@ mod linux { runtime, permit, active.clone(), + Arc::new(UnauthenticatedControlConnections::new( + MAX_UNAUTHENTICATED_CONTROL_CONNECTIONS, + )), deadline, )); @@ -4258,9 +4351,11 @@ mod linux { async fn control_slot_is_charged_only_after_a_valid_bearer() { let (runtime, token) = availability_test_runtime(); let active = Arc::new(AtomicUsize::new(0)); + let unauthenticated = Arc::new(UnauthenticatedControlConnections::new(2)); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server_active = active.clone(); + let server_unauthenticated = unauthenticated.clone(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.unwrap(); serve_grpc( @@ -4268,6 +4363,7 @@ mod linux { runtime, SandboxConnectionId::new(), server_active, + server_unauthenticated, Duration::from_secs(30), ) .await @@ -4285,6 +4381,7 @@ mod linux { .expect_err("missing bearer must be rejected"); assert_eq!(status.code(), tonic::Code::Unauthenticated); assert_eq!(active.load(Ordering::Acquire), 0); + assert_eq!(lock(&unauthenticated.entries).len(), 1); let status = client .exchange(bearer_request( @@ -4295,6 +4392,7 @@ mod linux { .expect_err("malformed bearer must be rejected"); assert_eq!(status.code(), tonic::Code::Unauthenticated); assert_eq!(active.load(Ordering::Acquire), 0); + assert_eq!(lock(&unauthenticated.entries).len(), 1); for _ in 0..2 { let mut body = client @@ -4308,6 +4406,7 @@ mod linux { 1, "one authenticated connection charges exactly one slot" ); + assert!(lock(&unauthenticated.entries).is_empty()); } drop(client); @@ -4328,6 +4427,7 @@ mod linux { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server_active = active.clone(); + let unauthenticated = Arc::new(UnauthenticatedControlConnections::new(2)); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.unwrap(); serve_grpc( @@ -4335,6 +4435,7 @@ mod linux { runtime, SandboxConnectionId::new(), server_active, + unauthenticated, Duration::from_secs(30), ) .await @@ -4514,6 +4615,7 @@ mod linux { server_runtime, connection_id, Arc::new(AtomicUsize::new(0)), + Arc::new(UnauthenticatedControlConnections::new(2)), CONTROL_UNAUTHENTICATED_DEADLINE, ) .await @@ -4956,6 +5058,7 @@ mod linux { boundary, SandboxConnectionId::new(), Arc::new(AtomicUsize::new(0)), + Arc::new(UnauthenticatedControlConnections::new(2)), CONTROL_UNAUTHENTICATED_DEADLINE, ) .await From 1e1af9810a8e4a6fd099f7b0007de94c9f51b84f Mon Sep 17 00:00:00 2001 From: Drew Newberry Date: Sat, 26 Sep 2026 18:08:13 -0700 Subject: [PATCH 3/3] fix(sandbox): hold permits until unauthenticated peers close Signed-off-by: Drew Newberry --- architecture/sandbox.md | 8 +- .../openshell-sandbox/src/boundary_server.rs | 144 ++++++++++++++---- 2 files changed, 119 insertions(+), 33 deletions(-) diff --git a/architecture/sandbox.md b/architecture/sandbox.md index 439b169f41..75dd53a322 100644 --- a/architecture/sandbox.md +++ b/architecture/sandbox.md @@ -31,9 +31,11 @@ not govern `connect()` on a filesystem socket, so the sandbox rejects peers inside its own PID namespace that are not the sandbox or one of its ancestors before TLS, charges a bounded control-connection slot only after the first accepted bearer, and closes connections that present no valid bearer within a -short deadline. TLS-complete connections awaiting a bearer have a separate -bounded pool; admitting a new connection closes the oldest waiting peer when -that pool is full. Independent bidirectional `Exchange` RPCs +short deadline. Unauthenticated gRPC transports have a separate bounded pool. +When it is full, a new connection closes the oldest waiting peer and holds a +handshake permit until the closed peer releases its pool permit. The same +unauthenticated deadline covers that wait. Independent +bidirectional `Exchange` RPCs carry lifecycle, exec, TCP, and forwarding traffic, while one persistent bidirectional `Mediate` RPC carries multiplexed DNS traffic. General application UDP is unsupported; UDP DNS remains mediated by the supervisor. diff --git a/crates/openshell-sandbox/src/boundary_server.rs b/crates/openshell-sandbox/src/boundary_server.rs index 0ee92dc689..9a6a8b2d1a 100644 --- a/crates/openshell-sandbox/src/boundary_server.rs +++ b/crates/openshell-sandbox/src/boundary_server.rs @@ -194,43 +194,82 @@ mod linux { } struct UnauthenticatedControlConnections { - capacity: usize, - entries: Mutex)>>, + permits: Arc, + entries: Mutex>, + } + + struct UnauthenticatedControlEntry { + connection_id: SandboxConnectionId, + shutdown: tokio::sync::watch::Sender<()>, + eviction_requested: bool, } impl UnauthenticatedControlConnections { fn new(capacity: usize) -> Self { assert!(capacity > 0); Self { - capacity, + permits: Arc::new(tokio::sync::Semaphore::new(capacity)), entries: Mutex::new(VecDeque::new()), } } - fn register( - &self, + async fn register( + self: &Arc, connection_id: SandboxConnectionId, shutdown: tokio::sync::watch::Sender<()>, - ) { + ) -> Arc { + let permit = match self.permits.clone().try_acquire_owned() { + Ok(permit) => permit, + Err(tokio::sync::TryAcquireError::NoPermits) => { + self.evict_oldest(); + self.permits + .clone() + .acquire_owned() + .await + .expect("unauthenticated connection pool stays open") + } + Err(tokio::sync::TryAcquireError::Closed) => { + unreachable!("unauthenticated connection pool stays open") + } + }; let mut entries = lock(&self.entries); - if entries.len() == self.capacity - && let Some((_, oldest)) = entries.pop_front() - { - // Make room for a new bearer attempt even when every older - // connection is idle. Its transport observes this shutdown. - let _ = oldest.send(()); + entries.push_back(UnauthenticatedControlEntry { + connection_id, + shutdown, + eviction_requested: false, + }); + Arc::new(UnauthenticatedControlRegistration { + connections: self.clone(), + connection_id, + permit: Mutex::new(Some(permit)), + }) + } + + fn evict_oldest(&self) { + let mut entries = lock(&self.entries); + if let Some(oldest) = entries.iter_mut().find(|entry| !entry.eviction_requested) { + oldest.eviction_requested = true; + // Keep its permit charged until the transport actually closes. + let _ = oldest.shutdown.send(()); } - entries.push_back((connection_id, shutdown)); } fn remove(&self, connection_id: SandboxConnectionId) { - lock(&self.entries).retain(|(id, _)| *id != connection_id); + lock(&self.entries).retain(|entry| entry.connection_id != connection_id); } } struct UnauthenticatedControlRegistration { connections: Arc, connection_id: SandboxConnectionId, + permit: Mutex>, + } + + impl UnauthenticatedControlRegistration { + fn release(&self) { + self.connections.remove(self.connection_id); + lock(&self.permit).take(); + } } impl Drop for UnauthenticatedControlRegistration { @@ -632,7 +671,6 @@ mod linux { .establish_async(&runtime.process_runtime) .await .map_err(|error| format!("establish boundary transport: {error}"))?; - drop(pending); // TLS completion is not authentication. The control slot is acquired // only after an RPC on this connection presents a valid bearer. serve_grpc( @@ -642,6 +680,7 @@ mod linux { active_connections, unauthenticated_connections, unauthenticated_deadline, + Some(pending), ) .await } @@ -653,16 +692,20 @@ mod linux { active_connections: Arc, unauthenticated_connections: Arc, unauthenticated_deadline: Duration, + pending_handshake: Option, ) -> Result<(), String> { + let unauthenticated_deadline = tokio::time::Instant::now() + unauthenticated_deadline; let (connection_shutdown, connection_closed) = tokio::sync::watch::channel(()); - unauthenticated_connections.register(connection_id, connection_shutdown.clone()); - let _unauthenticated_registration = UnauthenticatedControlRegistration { - connections: unauthenticated_connections.clone(), - connection_id, - }; + let unauthenticated_registration = tokio::time::timeout_at( + unauthenticated_deadline, + unauthenticated_connections.register(connection_id, connection_shutdown.clone()), + ) + .await + .map_err(|_| "unauthenticated control connection deadline reached".to_string())?; + drop(pending_handshake); runtime.register_connection(connection_id, connection_shutdown.clone()); let connection_expiry = Arc::new(ConnectionExpiry::new(connection_shutdown.clone())); - connection_expiry.update_deadline(tokio::time::Instant::now() + unauthenticated_deadline); + connection_expiry.update_deadline(unauthenticated_deadline); let incoming = tokio_stream::StreamExt::chain( tokio_stream::iter([Ok::<_, io::Error>(GrpcServerIo { stream, @@ -695,7 +738,7 @@ mod linux { connection_expiry, connection_closed, active_connections, - unauthenticated_connections, + unauthenticated_registration, control_slot: Arc::new(Mutex::new(None)), }) .max_decoding_message_size(64 * 1024) @@ -798,7 +841,7 @@ mod linux { connection_expiry: Arc, connection_closed: tokio::sync::watch::Receiver<()>, active_connections: Arc, - unauthenticated_connections: Arc, + unauthenticated_registration: Arc, /// Held from the first authenticated RPC until the connection ends. /// Shared across tonic's per-request service clones. control_slot: Arc>>, @@ -835,7 +878,7 @@ mod linux { } self.connection_expiry .update(principal.session().expires_at); - self.unauthenticated_connections.remove(self.connection_id); + self.unauthenticated_registration.release(); Ok(principal) } } @@ -4086,7 +4129,7 @@ mod linux { #[tokio::test] async fn unauthenticated_connection_pool_evicts_oldest_to_admit_new_peer() { - let connections = UnauthenticatedControlConnections::new(2); + let connections = Arc::new(UnauthenticatedControlConnections::new(2)); let first_id = SandboxConnectionId::new(); let second_id = SandboxConnectionId::new(); let third_id = SandboxConnectionId::new(); @@ -4094,16 +4137,53 @@ mod linux { let (second_shutdown, second_closed) = tokio::sync::watch::channel(()); let (third_shutdown, _third_closed) = tokio::sync::watch::channel(()); - connections.register(first_id, first_shutdown); - connections.register(second_id, second_shutdown); - connections.register(third_id, third_shutdown); + let first = connections.register(first_id, first_shutdown).await; + let _second = connections.register(second_id, second_shutdown).await; + let third = tokio::spawn({ + let connections = connections.clone(); + async move { connections.register(third_id, third_shutdown).await } + }); first_closed.changed().await.expect("oldest peer is closed"); + assert!(!third.is_finished(), "new peer waits for the old permit"); + assert_eq!(connections.permits.available_permits(), 0); assert!(!second_closed.has_changed().unwrap()); + first.release(); + let _third = tokio::time::timeout(Duration::from_secs(1), third) + .await + .expect("new peer is admitted once the old one releases its permit") + .unwrap(); let entries = lock(&connections.entries); assert_eq!(entries.len(), 2); - assert_eq!(entries[0].0, second_id); - assert_eq!(entries[1].0, third_id); + assert_eq!(entries[0].connection_id, second_id); + assert_eq!(entries[1].connection_id, third_id); + } + + #[tokio::test] + async fn unauthenticated_connection_wait_keeps_the_original_deadline() { + let (runtime, _) = availability_test_runtime(); + let connections = Arc::new(UnauthenticatedControlConnections::new(1)); + let (shutdown, _closed) = tokio::sync::watch::channel(()); + let held = connections + .register(SandboxConnectionId::new(), shutdown) + .await; + let (stream, _client) = tokio::io::duplex(64); + let result = serve_grpc( + Box::new(stream), + runtime, + SandboxConnectionId::new(), + Arc::new(AtomicUsize::new(0)), + connections.clone(), + Duration::from_millis(40), + None, + ) + .await; + assert!(result.unwrap_err().contains("deadline reached")); + assert_eq!(lock(&connections.entries).len(), 1); + assert_eq!(connections.permits.available_permits(), 0); + drop(held); + assert!(lock(&connections.entries).is_empty()); + assert_eq!(connections.permits.available_permits(), 1); } #[test] @@ -4365,6 +4445,7 @@ mod linux { server_active, server_unauthenticated, Duration::from_secs(30), + None, ) .await }); @@ -4437,6 +4518,7 @@ mod linux { server_active, unauthenticated, Duration::from_secs(30), + None, ) .await }); @@ -4617,6 +4699,7 @@ mod linux { Arc::new(AtomicUsize::new(0)), Arc::new(UnauthenticatedControlConnections::new(2)), CONTROL_UNAUTHENTICATED_DEADLINE, + None, ) .await }); @@ -5060,6 +5143,7 @@ mod linux { Arc::new(AtomicUsize::new(0)), Arc::new(UnauthenticatedControlConnections::new(2)), CONTROL_UNAUTHENTICATED_DEADLINE, + None, ) .await });