use std::collections::{BTreeMap, BTreeSet};
use std::fmt;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use futures::{SinkExt, StreamExt};
use supercode::{
CoordinatedRuntime, RuntimeAuthorization, RuntimeClientId, SdkRuntime,
DEFAULT_RUNTIME_LEASE_TTL_MS,
};
use tokio::net::{TcpListener, TcpStream};
#[cfg(unix)]
use tokio::net::{UnixListener, UnixStream};
use tokio::sync::{watch, Notify};
use tokio_tungstenite::tungstenite::handshake::server::{ErrorResponse, Request, Response};
use tokio_tungstenite::tungstenite::http::StatusCode;
use tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
use tokio_tungstenite::tungstenite::Message;
use zeroize::Zeroize;
use zeroize::Zeroizing;
use crate::codex_app_server_v0_144::{CodexAppServerAdapter, CodexCompatibilityMode};
const TOKEN_BYTES: usize = 32;
const MAX_MESSAGE_BYTES: usize = 16 * 1024 * 1024;
const MAX_FRAME_BYTES: usize = 1024 * 1024;
const MAX_WRITE_BUFFER_BYTES: usize = 2 * MAX_MESSAGE_BYTES;
const MAX_HANDSHAKE_BYTES: usize = 64 * 1024;
const MAX_CONCURRENT_CONNECTIONS: usize = 32;
#[cfg(not(test))]
const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
#[cfg(test)]
const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
#[derive(Debug, thiserror::Error)]
pub enum CredentialError {
#[error("operating system random source failed")]
RandomSource,
#[error("credential generation collided repeatedly")]
Collision,
#[error("credential is already revoked")]
Revoked,
#[error("bootstrap credential is invalid for this runtime generation")]
InvalidBootstrap,
#[error("invalid runtime client identity: {0}")]
InvalidClient(String),
#[error("private credential file operation failed: {0}")]
Io(#[from] std::io::Error),
#[error("endpoint configuration is already shared")]
ConfigurationLocked,
}
struct CredentialGrant {
client_id: RuntimeClientId,
authorization: RuntimeAuthorization,
revoked: watch::Receiver<bool>,
runtime_id: String,
generation: [u8; 16],
active_channel: watch::Sender<bool>,
}
impl Drop for CredentialGrant {
fn drop(&mut self) {
self.active_channel.send_replace(false);
}
}
struct CredentialRecord {
client_id: RuntimeClientId,
authorization: RuntimeAuthorization,
revoked: watch::Sender<bool>,
active_channel: watch::Sender<bool>,
}
struct CodexCredentialRegistry {
records: Mutex<BTreeMap<[u8; 32], CredentialRecord>>,
runtime_id: String,
generation: [u8; 16],
}
pub struct CodexBootstrapCredential {
secret: [u8; TOKEN_BYTES],
digest: [u8; 32],
generation: [u8; 16],
revoked: bool,
}
impl fmt::Debug for CodexBootstrapCredential {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("CodexBootstrapCredential([REDACTED])")
}
}
impl Drop for CodexBootstrapCredential {
fn drop(&mut self) {
zero(&mut self.secret);
}
}
pub struct CodexBootstrapFile {
path: PathBuf,
generation: [u8; 16],
}
impl Drop for CodexBootstrapFile {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
impl CodexCredentialRegistry {
fn new(runtime_id: String, generation: [u8; 16]) -> Arc<Self> {
Arc::new(Self {
records: Mutex::new(BTreeMap::new()),
runtime_id,
generation,
})
}
fn issue(
self: &Arc<Self>,
client_id: impl Into<String>,
authorization: RuntimeAuthorization,
) -> Result<CodexClientCredential, CredentialError> {
self.issue_with_fill(client_id, authorization, |secret| {
getrandom::getrandom(secret).map_err(|_| CredentialError::RandomSource)
})
}
fn issue_with_fill(
self: &Arc<Self>,
client_id: impl Into<String>,
authorization: RuntimeAuthorization,
mut fill: impl FnMut(&mut [u8; TOKEN_BYTES]) -> Result<(), CredentialError>,
) -> Result<CodexClientCredential, CredentialError> {
let client_id = RuntimeClientId::parse(client_id.into())
.map_err(|error| CredentialError::InvalidClient(error.to_string()))?;
for _ in 0..3 {
let mut secret = [0u8; TOKEN_BYTES];
fill(&mut secret)?;
let digest = credential_digest(&secret);
let mut records = self
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if records.contains_key(&digest) {
zero(&mut secret);
continue;
}
let (revoked, _) = watch::channel(false);
let (active_channel, _) = watch::channel(false);
records.insert(
digest,
CredentialRecord {
client_id: client_id.clone(),
authorization,
revoked,
active_channel,
},
);
return Ok(CodexClientCredential {
secret,
digest,
client_id,
revoked: false,
});
}
Err(CredentialError::Collision)
}
async fn revoke(&self, credential: &mut CodexClientCredential) -> Result<(), CredentialError> {
if credential.revoked {
return Err(CredentialError::Revoked);
}
let record = self
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&credential.digest)
.ok_or(CredentialError::Revoked)?;
let mut active_channel = record.active_channel.subscribe();
let _ = record.revoked.send(true);
while *active_channel.borrow() {
if active_channel.changed().await.is_err() {
break;
}
}
credential.revoked = true;
zero(&mut credential.secret);
Ok(())
}
async fn rotate(
self: &Arc<Self>,
credential: &mut CodexClientCredential,
client_id: impl Into<String>,
authorization: RuntimeAuthorization,
) -> Result<CodexClientCredential, CredentialError> {
let mut replacement = self.issue(client_id, authorization)?;
if let Err(error) = self.revoke(credential).await {
let _ = self.revoke(&mut replacement).await;
return Err(error);
}
Ok(replacement)
}
async fn revoke_all(&self) {
let records = {
let mut records = self
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
std::mem::take(&mut *records)
};
let mut channels = Vec::with_capacity(records.len());
for (_, record) in records {
let mut active = record.active_channel.subscribe();
let _ = record.revoked.send(true);
channels.push(async move {
while *active.borrow() {
if active.changed().await.is_err() {
break;
}
}
});
}
futures::future::join_all(channels).await;
}
fn authenticate(&self, token: &str) -> Result<CredentialGrant, AuthenticationError> {
let (presented_client_id, secret_text) = token
.split_once('.')
.ok_or(AuthenticationError::Unauthenticated)?;
let mut secret =
decode_hex_secret(secret_text).ok_or(AuthenticationError::Unauthenticated)?;
let digest = credential_digest(&secret);
zero(&mut secret);
let mut records = self
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let record = records
.get_mut(&digest)
.ok_or(AuthenticationError::Unauthenticated)?;
if record.client_id.as_str() != presented_client_id {
return Err(AuthenticationError::Unauthorized);
}
if *record.active_channel.borrow() {
return Err(AuthenticationError::Busy);
}
record.active_channel.send_replace(true);
Ok(CredentialGrant {
client_id: record.client_id.clone(),
authorization: record.authorization.clone(),
revoked: record.revoked.subscribe(),
runtime_id: self.runtime_id.clone(),
generation: self.generation,
active_channel: record.active_channel.clone(),
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AuthenticationError {
Unauthenticated,
Unauthorized,
Busy,
}
pub struct CodexClientCredential {
secret: [u8; TOKEN_BYTES],
digest: [u8; 32],
client_id: RuntimeClientId,
revoked: bool,
}
impl CodexClientCredential {
pub fn spawn_tokio_child(
&self,
command: &mut tokio::process::Command,
variable: &str,
) -> std::io::Result<tokio::process::Child> {
let bearer = Zeroizing::new(self.bearer_value());
command.env(variable, bearer.as_str());
let child = command.spawn();
command.env_remove(variable);
child
}
#[cfg(test)]
fn bearer(&self) -> String {
format!("Bearer {}", self.bearer_value())
}
fn bearer_value(&self) -> String {
format!("{}.{}", self.client_id.as_str(), encode_hex(&self.secret))
}
}
impl fmt::Debug for CodexClientCredential {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("CodexClientCredential([REDACTED])")
}
}
impl Drop for CodexClientCredential {
fn drop(&mut self) {
zero(&mut self.secret);
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CodexEndpointKind {
WebSocket(SocketAddr),
#[cfg(unix)]
Unix(PathBuf),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CodexEndpointHealth {
Ready,
ShuttingDown,
Stopped,
}
pub struct CodexEndpoint {
coordinator: Arc<CoordinatedRuntime>,
credentials: Arc<CodexCredentialRegistry>,
bootstrap_digest: Mutex<[u8; 32]>,
cwd: PathBuf,
client_home: Option<PathBuf>,
mode: CodexCompatibilityMode,
forwarded_ports: BTreeSet<u16>,
}
impl CodexEndpoint {
pub fn new(
runtime: Arc<dyn SdkRuntime>,
runtime_id: impl Into<String>,
cwd: impl Into<PathBuf>,
) -> Result<(Arc<Self>, CodexBootstrapCredential), CredentialError> {
let generation = random_generation()?;
let bootstrap = new_bootstrap(generation)?;
let endpoint = Arc::new(Self {
coordinator: CoordinatedRuntime::with_lease_ttl(runtime, DEFAULT_RUNTIME_LEASE_TTL_MS),
credentials: CodexCredentialRegistry::new(runtime_id.into(), generation),
bootstrap_digest: Mutex::new(bootstrap.digest),
cwd: cwd.into(),
client_home: None,
mode: CodexCompatibilityMode::TracedOnly,
forwarded_ports: BTreeSet::new(),
});
Ok((endpoint, bootstrap))
}
pub fn with_mode(
runtime: Arc<dyn SdkRuntime>,
runtime_id: impl Into<String>,
cwd: impl Into<PathBuf>,
mode: CodexCompatibilityMode,
) -> Result<(Arc<Self>, CodexBootstrapCredential), CredentialError> {
let generation = random_generation()?;
let bootstrap = new_bootstrap(generation)?;
let endpoint = Arc::new(Self {
coordinator: CoordinatedRuntime::with_lease_ttl(runtime, DEFAULT_RUNTIME_LEASE_TTL_MS),
credentials: CodexCredentialRegistry::new(runtime_id.into(), generation),
bootstrap_digest: Mutex::new(bootstrap.digest),
cwd: cwd.into(),
client_home: None,
mode,
forwarded_ports: BTreeSet::new(),
});
Ok((endpoint, bootstrap))
}
pub fn with_forwarded_ports(
mut self: Arc<Self>,
ports: impl IntoIterator<Item = u16>,
) -> Result<Arc<Self>, CredentialError> {
Arc::get_mut(&mut self)
.ok_or(CredentialError::ConfigurationLocked)?
.forwarded_ports
.extend(ports.into_iter().filter(|port| *port != 0));
Ok(self)
}
pub fn with_client_home(
mut self: Arc<Self>,
path: impl Into<PathBuf>,
) -> Result<Arc<Self>, CredentialError> {
Arc::get_mut(&mut self)
.ok_or(CredentialError::ConfigurationLocked)?
.client_home = Some(path.into());
Ok(self)
}
pub fn issue_interactive(
self: &Arc<Self>,
bootstrap: &CodexBootstrapCredential,
client_id: impl Into<String>,
) -> Result<CodexClientCredential, CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.credentials.issue(
client_id,
RuntimeAuthorization::new([
supercode::RuntimePermission::Observe,
supercode::RuntimePermission::Interact,
supercode::RuntimePermission::Approve,
]),
)
}
pub fn issue_observer(
self: &Arc<Self>,
bootstrap: &CodexBootstrapCredential,
client_id: impl Into<String>,
) -> Result<CodexClientCredential, CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.credentials
.issue(client_id, RuntimeAuthorization::observer())
}
pub async fn rotate_interactive(
self: &Arc<Self>,
bootstrap: &CodexBootstrapCredential,
credential: &mut CodexClientCredential,
client_id: impl Into<String>,
) -> Result<CodexClientCredential, CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.credentials
.rotate(
credential,
client_id,
RuntimeAuthorization::new([
supercode::RuntimePermission::Observe,
supercode::RuntimePermission::Interact,
supercode::RuntimePermission::Approve,
]),
)
.await
}
pub fn rotate_bootstrap(
&self,
bootstrap: &mut CodexBootstrapCredential,
) -> Result<CodexBootstrapCredential, CredentialError> {
let replacement = new_bootstrap(self.credentials.generation)?;
let mut expected = self
.bootstrap_digest
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if bootstrap.revoked
|| bootstrap.generation != self.credentials.generation
|| *expected != bootstrap.digest
|| credential_digest(&bootstrap.secret) != *expected
{
return Err(CredentialError::InvalidBootstrap);
}
*expected = replacement.digest;
bootstrap.revoked = true;
zero(&mut bootstrap.secret);
Ok(replacement)
}
#[cfg(unix)]
pub fn rotate_bootstrap_file(
&self,
bootstrap: &mut CodexBootstrapCredential,
bootstrap_file: &mut CodexBootstrapFile,
runtime_dir: impl Into<PathBuf>,
) -> Result<(CodexBootstrapCredential, CodexBootstrapFile), CredentialError> {
self.rotate_bootstrap_file_with_remove(
bootstrap,
bootstrap_file,
runtime_dir.into(),
|path| std::fs::remove_file(path),
)
}
#[cfg(unix)]
fn rotate_bootstrap_file_with_remove(
&self,
bootstrap: &mut CodexBootstrapCredential,
bootstrap_file: &mut CodexBootstrapFile,
runtime_dir: PathBuf,
remove_old: impl FnOnce(&std::path::Path) -> std::io::Result<()>,
) -> Result<(CodexBootstrapCredential, CodexBootstrapFile), CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.authenticate_bootstrap_file(bootstrap_file)?;
let replacement = new_bootstrap(self.credentials.generation)?;
let replacement_file =
write_bootstrap_secret(runtime_dir, &replacement.secret, replacement.generation)?;
let mut expected = self
.bootstrap_digest
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if *expected != bootstrap.digest || credential_digest(&bootstrap.secret) != *expected {
drop(replacement_file);
return Err(CredentialError::InvalidBootstrap);
}
remove_old(&bootstrap_file.path)?;
*expected = replacement.digest;
bootstrap.revoked = true;
zero(&mut bootstrap.secret);
bootstrap_file.generation = [0; 16];
Ok((replacement, replacement_file))
}
pub async fn revoke_frontend(
&self,
bootstrap: &CodexBootstrapCredential,
credential: &mut CodexClientCredential,
) -> Result<(), CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.credentials.revoke(credential).await
}
#[cfg(unix)]
pub fn write_bootstrap_file(
&self,
bootstrap: &CodexBootstrapCredential,
runtime_dir: impl Into<PathBuf>,
) -> Result<CodexBootstrapFile, CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
write_bootstrap_secret(runtime_dir.into(), &bootstrap.secret, bootstrap.generation)
}
#[cfg(unix)]
pub fn issue_interactive_from_file(
self: &Arc<Self>,
bootstrap_file: &CodexBootstrapFile,
client_id: impl Into<String>,
) -> Result<CodexClientCredential, CredentialError> {
self.authenticate_bootstrap_file(bootstrap_file)?;
self.credentials.issue(
client_id,
RuntimeAuthorization::new([
supercode::RuntimePermission::Observe,
supercode::RuntimePermission::Interact,
supercode::RuntimePermission::Approve,
]),
)
}
#[cfg(unix)]
pub async fn revoke_frontend_from_file(
&self,
bootstrap_file: &CodexBootstrapFile,
credential: &mut CodexClientCredential,
) -> Result<(), CredentialError> {
self.authenticate_bootstrap_file(bootstrap_file)?;
self.credentials.revoke(credential).await
}
pub async fn terminate(
self: &Arc<Self>,
bootstrap: &CodexBootstrapCredential,
) -> Result<(), CredentialError> {
self.authenticate_bootstrap(bootstrap)?;
self.credentials.revoke_all().await;
let client_id = RuntimeClientId::parse("supercode-bootstrap-operator")
.map_err(|error| CredentialError::InvalidClient(error.to_string()))?;
self.coordinator
.client(client_id, RuntimeAuthorization::owner())
.close()
.await
.map_err(|_| CredentialError::InvalidBootstrap)
}
fn authenticate_bootstrap(
&self,
bootstrap: &CodexBootstrapCredential,
) -> Result<(), CredentialError> {
if bootstrap.revoked || bootstrap.generation != self.credentials.generation {
return Err(CredentialError::InvalidBootstrap);
}
let expected = self
.bootstrap_digest
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if *expected != bootstrap.digest || credential_digest(&bootstrap.secret) != *expected {
return Err(CredentialError::InvalidBootstrap);
}
Ok(())
}
#[cfg(unix)]
fn authenticate_bootstrap_file(
&self,
bootstrap_file: &CodexBootstrapFile,
) -> Result<(), CredentialError> {
if bootstrap_file.generation != self.credentials.generation {
return Err(CredentialError::InvalidBootstrap);
}
let mut encoded = std::fs::read(&bootstrap_file.path)?;
let decoded = std::str::from_utf8(&encoded)
.ok()
.and_then(decode_hex_secret)
.ok_or(CredentialError::InvalidBootstrap);
encoded.zeroize();
let mut secret = decoded?;
let digest = credential_digest(&secret);
zero(&mut secret);
let expected = self
.bootstrap_digest
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if digest != *expected {
return Err(CredentialError::InvalidBootstrap);
}
Ok(())
}
pub async fn bind_websocket(self: &Arc<Self>, port: u16) -> std::io::Result<CodexServerHandle> {
let listener =
TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port)).await?;
let address = listener.local_addr()?;
let shutdown = Arc::new(Notify::new());
let shutting_down = Arc::new(AtomicBool::new(false));
let stopped = Arc::new(AtomicBool::new(false));
let task = tokio::spawn(run_tcp_accept_loop(
listener,
self.clone(),
address.to_string(),
self.forwarded_ports.clone(),
shutdown.clone(),
shutting_down.clone(),
stopped.clone(),
));
Ok(CodexServerHandle {
kind: CodexEndpointKind::WebSocket(address),
shutdown,
shutting_down,
stopped,
task: Some(task),
})
}
#[cfg(unix)]
pub async fn bind_unix(
self: &Arc<Self>,
socket_path: impl Into<PathBuf>,
) -> std::io::Result<CodexServerHandle> {
use std::os::unix::fs::{DirBuilderExt, MetadataExt, PermissionsExt};
let socket_path = socket_path.into();
let parent = socket_path.parent().ok_or_else(|| {
std::io::Error::new(std::io::ErrorKind::InvalidInput, "socket has no parent")
})?;
match std::fs::symlink_metadata(parent) {
Ok(metadata) => {
if metadata.file_type().is_symlink() || !metadata.is_dir() {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"Unix endpoint directory must be a real directory",
));
}
if metadata.permissions().mode() & 0o777 != 0o700 {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"Unix endpoint directory must have mode 0700",
));
}
if metadata.uid() != unsafe { libc::geteuid() } {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"Unix endpoint directory must be owned by this user",
));
}
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
let mut builder = std::fs::DirBuilder::new();
builder.mode(0o700).create(parent)?;
}
Err(error) => return Err(error),
}
match std::fs::symlink_metadata(&socket_path) {
Ok(_) => {
return Err(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"refusing to replace an existing Unix endpoint",
))
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(error),
}
let listener = UnixListener::bind(&socket_path)?;
std::fs::set_permissions(&socket_path, std::fs::Permissions::from_mode(0o600))?;
let shutdown = Arc::new(Notify::new());
let shutting_down = Arc::new(AtomicBool::new(false));
let stopped = Arc::new(AtomicBool::new(false));
let task = tokio::spawn(run_unix_accept_loop(
listener,
self.clone(),
socket_path.clone(),
shutdown.clone(),
shutting_down.clone(),
stopped.clone(),
));
Ok(CodexServerHandle {
kind: CodexEndpointKind::Unix(socket_path),
shutdown,
shutting_down,
stopped,
task: Some(task),
})
}
fn authenticated_adapter(&self, grant: &CredentialGrant) -> Arc<CodexAppServerAdapter> {
debug_assert_eq!(grant.runtime_id, self.credentials.runtime_id);
debug_assert_eq!(grant.generation, self.credentials.generation);
let runtime = self
.coordinator
.client(grant.client_id.clone(), grant.authorization.clone());
CodexAppServerAdapter::with_mode_and_client_home(
runtime,
self.cwd.clone(),
self.client_home.clone().unwrap_or_else(|| self.cwd.clone()),
self.mode,
)
}
}
pub struct CodexServerHandle {
kind: CodexEndpointKind,
shutdown: Arc<Notify>,
shutting_down: Arc<AtomicBool>,
stopped: Arc<AtomicBool>,
task: Option<tokio::task::JoinHandle<()>>,
}
impl CodexServerHandle {
pub fn kind(&self) -> &CodexEndpointKind {
&self.kind
}
pub fn health(&self) -> CodexEndpointHealth {
if self.stopped.load(Ordering::SeqCst) {
CodexEndpointHealth::Stopped
} else if self.shutting_down.load(Ordering::SeqCst) {
CodexEndpointHealth::ShuttingDown
} else {
CodexEndpointHealth::Ready
}
}
pub async fn shutdown(mut self) {
self.shutting_down.store(true, Ordering::SeqCst);
self.shutdown.notify_waiters();
if let Some(task) = self.task.take() {
let _ = task.await;
}
}
}
impl Drop for CodexServerHandle {
fn drop(&mut self) {
self.shutting_down.store(true, Ordering::SeqCst);
self.shutdown.notify_waiters();
if let Some(task) = self.task.take() {
task.abort();
}
}
}
async fn run_tcp_accept_loop(
listener: TcpListener,
endpoint: Arc<CodexEndpoint>,
expected_host: String,
forwarded_ports: BTreeSet<u16>,
shutdown: Arc<Notify>,
shutting_down: Arc<AtomicBool>,
stopped: Arc<AtomicBool>,
) {
let connection_budget = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
loop {
tokio::select! {
_ = shutdown.notified() => break,
accepted = listener.accept() => match accepted {
Ok((stream, _)) => {
let Ok(permit) = connection_budget.clone().try_acquire_owned() else {
drop(stream);
continue;
};
let endpoint = endpoint.clone();
let expected_host = expected_host.clone();
let forwarded_ports = forwarded_ports.clone();
tokio::spawn(async move {
let _permit = permit;
let _ = serve_stream(stream, endpoint, HostRule::Tcp { expected_host, forwarded_ports }).await;
});
}
Err(_) => break,
}
}
}
shutting_down.store(true, Ordering::SeqCst);
stopped.store(true, Ordering::SeqCst);
}
#[cfg(unix)]
async fn run_unix_accept_loop(
listener: UnixListener,
endpoint: Arc<CodexEndpoint>,
socket_path: PathBuf,
shutdown: Arc<Notify>,
shutting_down: Arc<AtomicBool>,
stopped: Arc<AtomicBool>,
) {
let connection_budget = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
loop {
tokio::select! {
_ = shutdown.notified() => break,
accepted = listener.accept() => match accepted {
Ok((stream, _)) => {
let Ok(permit) = connection_budget.clone().try_acquire_owned() else {
drop(stream);
continue;
};
let endpoint = endpoint.clone();
tokio::spawn(async move {
let _permit = permit;
let _ = serve_stream(stream, endpoint, HostRule::Unix).await;
});
}
Err(_) => break,
}
}
}
let _ = std::fs::remove_file(socket_path);
shutting_down.store(true, Ordering::SeqCst);
stopped.store(true, Ordering::SeqCst);
}
enum HostRule {
Tcp {
expected_host: String,
forwarded_ports: BTreeSet<u16>,
},
#[cfg(unix)]
Unix,
}
trait LocalStream: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static {}
impl LocalStream for TcpStream {}
#[cfg(unix)]
impl LocalStream for UnixStream {}
struct HandshakeBoundedStream<S> {
inner: S,
bytes_read: usize,
complete: Arc<AtomicBool>,
}
impl<S: Unpin> Unpin for HandshakeBoundedStream<S> {}
impl<S: LocalStream> LocalStream for HandshakeBoundedStream<S> {}
impl<S: tokio::io::AsyncRead + Unpin> tokio::io::AsyncRead for HandshakeBoundedStream<S> {
fn poll_read(
mut self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
buffer: &mut tokio::io::ReadBuf<'_>,
) -> std::task::Poll<std::io::Result<()>> {
let before = buffer.filled().len();
match std::pin::Pin::new(&mut self.inner).poll_read(context, buffer) {
std::task::Poll::Ready(Ok(())) => {
if !self.complete.load(Ordering::Acquire) {
self.bytes_read = self
.bytes_read
.saturating_add(buffer.filled().len().saturating_sub(before));
if self.bytes_read > MAX_HANDSHAKE_BYTES {
return std::task::Poll::Ready(Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"WebSocket handshake exceeded byte limit",
)));
}
}
std::task::Poll::Ready(Ok(()))
}
other => other,
}
}
}
impl<S: tokio::io::AsyncWrite + Unpin> tokio::io::AsyncWrite for HandshakeBoundedStream<S> {
fn poll_write(
mut self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
buffer: &[u8],
) -> std::task::Poll<std::io::Result<usize>> {
std::pin::Pin::new(&mut self.inner).poll_write(context, buffer)
}
fn poll_flush(
mut self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::pin::Pin::new(&mut self.inner).poll_flush(context)
}
fn poll_shutdown(
mut self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::pin::Pin::new(&mut self.inner).poll_shutdown(context)
}
}
#[allow(clippy::result_large_err)]
async fn serve_stream<S: LocalStream>(
stream: S,
endpoint: Arc<CodexEndpoint>,
host_rule: HostRule,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let selected = Arc::new(Mutex::new(None::<CredentialGrant>));
let callback_selected = selected.clone();
let credentials = endpoint.credentials.clone();
let handshake_complete = Arc::new(AtomicBool::new(false));
let callback_complete = handshake_complete.clone();
let stream = HandshakeBoundedStream {
inner: stream,
bytes_read: 0,
complete: handshake_complete,
};
let websocket = tokio::time::timeout(
HANDSHAKE_TIMEOUT,
tokio_tungstenite::accept_hdr_async_with_config(
stream,
move |request: &Request, response: Response| -> Result<Response, ErrorResponse> {
match validate_upgrade(request, &host_rule, &credentials) {
Ok(grant) => {
*callback_selected
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(grant);
callback_complete.store(true, Ordering::Release);
Ok(response)
}
Err((status, message)) => Err(http_error(status, message)),
}
},
Some(websocket_config()),
),
)
.await
.map_err(|_| {
std::io::Error::new(
std::io::ErrorKind::TimedOut,
"WebSocket handshake timed out",
)
})??;
let grant = selected
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.ok_or("authenticated upgrade did not select a grant")?;
serve_websocket(websocket, endpoint, grant).await
}
async fn serve_websocket<S: LocalStream>(
mut websocket: tokio_tungstenite::WebSocketStream<S>,
endpoint: Arc<CodexEndpoint>,
mut grant: CredentialGrant,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let runtime = endpoint.authenticated_adapter(&grant);
let mut connection = runtime.connection();
loop {
tokio::select! {
changed = grant.revoked.changed() => {
if changed.is_err() || *grant.revoked.borrow() {
let _ = websocket.close(None).await;
break;
}
}
incoming = websocket.next() => {
let Some(message) = incoming else { break; };
match message? {
Message::Text(text) => {
let value: serde_json::Value = match serde_json::from_str(&text) {
Ok(value) => value,
Err(_) => {
websocket.send(Message::Text(serde_json::json!({"id":null,"error":{
"code":-32700,"message":"invalid JSON request","data":{"name":"parse_error"}
}}).to_string().into())).await?;
continue;
}
};
if contains_application_credential(&value) {
let id = value.get("id").cloned().unwrap_or(serde_json::Value::Null);
websocket.send(Message::Text(serde_json::json!({"id":id,"error":{
"code":-32030,"message":"application-message credentials are forbidden",
"data":{"name":"unauthenticated"}
}}).to_string().into())).await?;
continue;
}
let output = if value.get("method").is_some() {
connection.handle(value).await
} else {
connection.handle_server_response(&value).await?
};
for message in output {
websocket.send(Message::Text(message.to_string().into())).await?;
}
}
Message::Close(_) => break,
Message::Ping(payload) => websocket.send(Message::Pong(payload)).await?,
Message::Binary(_) => {
websocket.close(None).await?;
break;
}
Message::Pong(_) | Message::Frame(_) => {}
}
}
notifications = connection.next_notifications(), if connection_is_attached(&connection) => {
for message in notifications? {
websocket.send(Message::Text(message.to_string().into())).await?;
}
}
}
}
let _ = runtime.detach().await;
Ok(())
}
fn connection_is_attached(connection: &crate::CodexConnection) -> bool {
connection.is_attached()
}
fn validate_upgrade(
request: &Request,
host_rule: &HostRule,
credentials: &CodexCredentialRegistry,
) -> Result<CredentialGrant, (StatusCode, &'static str)> {
if request.uri().path() != "/" || request.uri().query().is_some() {
return Err((StatusCode::NOT_FOUND, "unsupported endpoint"));
}
if request.headers().contains_key("cookie")
|| request.headers().contains_key("sec-websocket-protocol")
|| request.headers().contains_key("transfer-encoding")
|| request.headers().contains_key("content-length")
{
return Err((StatusCode::BAD_REQUEST, "forbidden credential carrier"));
}
for name in [
"host",
"origin",
"authorization",
"content-length",
"x-supercode-permissions",
] {
if request.headers().get_all(name).iter().count() > 1 {
return Err((StatusCode::BAD_REQUEST, "duplicate security header"));
}
}
if request.uri().scheme().is_some() || request.uri().authority().is_some() {
return Err((StatusCode::BAD_REQUEST, "absolute-form target rejected"));
}
let host = request
.headers()
.get("host")
.and_then(|value| value.to_str().ok())
.ok_or((StatusCode::BAD_REQUEST, "missing Host"))?;
match host_rule {
HostRule::Tcp {
expected_host,
forwarded_ports,
} => {
let accepted_forward = host
.strip_prefix("127.0.0.1:")
.and_then(|port| port.parse::<u16>().ok())
.is_some_and(|port| forwarded_ports.contains(&port));
if host != expected_host && !accepted_forward {
return Err((StatusCode::BAD_REQUEST, "invalid loopback Host"));
}
}
#[cfg(unix)]
HostRule::Unix if host != "localhost" => {
return Err((StatusCode::BAD_REQUEST, "invalid Unix Host"));
}
#[cfg(unix)]
HostRule::Unix => {}
}
if request.headers().contains_key("origin") {
return Err((StatusCode::BAD_REQUEST, "Origin rejected"));
}
let authorization = request
.headers()
.get("authorization")
.and_then(|value| value.to_str().ok())
.and_then(|value| value.strip_prefix("Bearer "))
.ok_or((StatusCode::UNAUTHORIZED, "authentication failed"))?;
let mut grant = credentials
.authenticate(authorization)
.map_err(|error| match error {
AuthenticationError::Unauthenticated => {
(StatusCode::UNAUTHORIZED, "authentication failed")
}
AuthenticationError::Unauthorized => (StatusCode::FORBIDDEN, "client id rejected"),
AuthenticationError::Busy => (StatusCode::CONFLICT, "client channel busy"),
})?;
if let Some(requested) = request.headers().get("x-supercode-permissions") {
let requested = requested
.to_str()
.ok()
.and_then(|value| RuntimeAuthorization::parse_header(value).ok())
.ok_or((StatusCode::BAD_REQUEST, "invalid permission narrowing"))?;
grant.authorization = grant.authorization.restrict_to(&requested);
}
Ok(grant)
}
fn contains_application_credential(value: &serde_json::Value) -> bool {
match value {
serde_json::Value::Object(object) => object.iter().any(|(key, child)| {
matches!(
key.to_ascii_lowercase().as_str(),
"authorization" | "bearer" | "credential" | "token"
) || contains_application_credential(child)
}),
serde_json::Value::Array(array) => array.iter().any(contains_application_credential),
_ => false,
}
}
fn websocket_config() -> WebSocketConfig {
let mut config = WebSocketConfig::default();
config.write_buffer_size = 128 * 1024;
config.max_write_buffer_size = MAX_WRITE_BUFFER_BYTES;
config.max_message_size = Some(MAX_MESSAGE_BYTES);
config.max_frame_size = Some(MAX_FRAME_BYTES);
config.accept_unmasked_frames = false;
config
}
fn http_error(status: StatusCode, message: &'static str) -> ErrorResponse {
let (code, name) = match status {
StatusCode::UNAUTHORIZED => (-32030, "unauthenticated"),
StatusCode::FORBIDDEN => (-32031, "unauthorized"),
StatusCode::CONFLICT => (-32000, "busy"),
_ => (-32600, "invalid_request"),
};
tokio_tungstenite::tungstenite::http::Response::builder()
.status(status)
.header("content-type", "application/json")
.body(Some(
serde_json::json!({"error":{"code":code,"message":message,"data":{"name":name}}})
.to_string(),
))
.expect("static HTTP error response")
}
fn credential_digest(secret: &[u8; TOKEN_BYTES]) -> [u8; 32] {
*blake3::hash(secret).as_bytes()
}
fn random_generation() -> Result<[u8; 16], CredentialError> {
let mut generation = [0u8; 16];
getrandom::getrandom(&mut generation).map_err(|_| CredentialError::RandomSource)?;
Ok(generation)
}
fn new_bootstrap(generation: [u8; 16]) -> Result<CodexBootstrapCredential, CredentialError> {
let mut secret = [0u8; TOKEN_BYTES];
getrandom::getrandom(&mut secret).map_err(|_| CredentialError::RandomSource)?;
Ok(CodexBootstrapCredential {
digest: credential_digest(&secret),
secret,
generation,
revoked: false,
})
}
#[cfg(unix)]
fn write_bootstrap_secret(
runtime_dir: PathBuf,
secret: &[u8; TOKEN_BYTES],
generation: [u8; 16],
) -> Result<CodexBootstrapFile, CredentialError> {
use std::io::Write;
use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt};
let metadata = std::fs::symlink_metadata(&runtime_dir)?;
if metadata.file_type().is_symlink()
|| !metadata.is_dir()
|| metadata.permissions().mode() & 0o777 != 0o700
|| metadata.uid() != unsafe { libc::geteuid() }
{
return Err(CredentialError::Io(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"runtime credential directory must be owner-owned mode 0700 and not a symlink",
)));
}
for _ in 0..3 {
let mut name_bytes = [0u8; 16];
getrandom::getrandom(&mut name_bytes).map_err(|_| CredentialError::RandomSource)?;
let name = name_bytes
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
name_bytes.zeroize();
let path = runtime_dir.join(format!("bootstrap-{name}.credential"));
let mut options = std::fs::OpenOptions::new();
options
.write(true)
.create_new(true)
.mode(0o600)
.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
let mut file = match options.open(&path) {
Ok(file) => file,
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(error) => return Err(CredentialError::Io(error)),
};
let mut encoded = encode_hex(secret);
let write_result = file
.write_all(encoded.as_bytes())
.and_then(|_| file.sync_all());
encoded.zeroize();
if let Err(error) = write_result {
let _ = std::fs::remove_file(&path);
return Err(CredentialError::Io(error));
}
let written = match file.metadata() {
Ok(metadata) => metadata,
Err(error) => {
drop(file);
let _ = std::fs::remove_file(&path);
return Err(CredentialError::Io(error));
}
};
if !written.is_file()
|| written.permissions().mode() & 0o777 != 0o600
|| written.uid() != unsafe { libc::geteuid() }
{
drop(file);
let _ = std::fs::remove_file(&path);
return Err(CredentialError::Io(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"bootstrap credential file failed owner/mode verification",
)));
}
return Ok(CodexBootstrapFile { path, generation });
}
Err(CredentialError::Collision)
}
fn encode_hex(secret: &[u8; TOKEN_BYTES]) -> String {
const HEX: &[u8; 16] = b"0123456789abcdef";
let mut output = String::with_capacity(TOKEN_BYTES * 2);
for byte in secret {
output.push(HEX[(byte >> 4) as usize] as char);
output.push(HEX[(byte & 0x0f) as usize] as char);
}
output
}
fn decode_hex_secret(value: &str) -> Option<[u8; TOKEN_BYTES]> {
if value.len() != TOKEN_BYTES * 2 {
return None;
}
let mut output = [0u8; TOKEN_BYTES];
for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
output[index] = (hex_nibble(pair[0])? << 4) | hex_nibble(pair[1])?;
}
Some(output)
}
fn hex_nibble(value: u8) -> Option<u8> {
match value {
b'0'..=b'9' => Some(value - b'0'),
b'a'..=b'f' => Some(value - b'a' + 10),
_ => None,
}
}
fn zero(secret: &mut [u8; TOKEN_BYTES]) {
secret.zeroize();
}
#[cfg(test)]
mod tests {
use super::*;
use async_trait::async_trait;
use serde_json::{json, Value};
use std::collections::VecDeque;
use supercode::{
ChatMessage, FrontendAttachSnapshot, FrontendAttachment, FrontendConnectionState,
FrontendDisplayCapabilities, FrontendOperationInvocation, FrontendOperationResult,
FrontendResponse, FrontendRuntimeDescriptor, FrontendTurnState, RuntimePermission,
SdkError,
};
use tokio::sync::broadcast;
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
struct HandshakeRuntime;
#[async_trait]
impl SdkRuntime for HandshakeRuntime {
async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
Err(SdkError::UnsupportedOperation("test.describe".into()))
}
async fn attach(&self, _history_limit: usize) -> Result<FrontendAttachment, SdkError> {
Err(SdkError::UnsupportedOperation("test.attach".into()))
}
async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
Err(SdkError::UnsupportedOperation("test.input".into()))
}
async fn submit(&self, _prompt: String) -> Result<String, SdkError> {
Err(SdkError::UnsupportedOperation("test.submit".into()))
}
async fn interrupt(&self) -> Result<bool, SdkError> {
Err(SdkError::UnsupportedOperation("test.interrupt".into()))
}
async fn steer(&self, _prompt: String) -> Result<(), SdkError> {
Err(SdkError::UnsupportedOperation("test.steer".into()))
}
async fn respond(&self, _response: FrontendResponse) -> Result<(), SdkError> {
Err(SdkError::UnsupportedOperation("test.respond".into()))
}
async fn invoke(
&self,
_operation: FrontendOperationInvocation,
) -> Result<FrontendOperationResult, SdkError> {
Err(SdkError::UnsupportedOperation("test.invoke".into()))
}
}
struct FunctionalRuntime {
events: broadcast::Sender<supercode::FrontendEvent>,
}
impl FunctionalRuntime {
fn new() -> Arc<Self> {
let (events, _) = broadcast::channel(32);
Arc::new(Self { events })
}
fn descriptor() -> FrontendRuntimeDescriptor {
FrontendRuntimeDescriptor {
schema_version: 2,
session_id: "canonical-reconnect-session".into(),
source_harness: Some("claude-code".into()),
emulation_profile: Some("claude-code".into()),
active_modules: Vec::new(),
commands: Vec::new(),
operations: Vec::new(),
actions: supercode::FrontendActions {
submit: true,
interrupt: true,
steer: true,
respond: true,
detach: true,
close: true,
},
display: FrontendDisplayCapabilities {
event_kinds: vec!["text_delta".into()],
opaque_fallback: true,
},
model: "z-ai/glm-5.2".into(),
turn_state: FrontendTurnState::Idle,
connection_state: FrontendConnectionState::Connected,
extensions: BTreeMap::new(),
}
}
}
#[async_trait]
impl SdkRuntime for FunctionalRuntime {
async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
Ok(Self::descriptor())
}
async fn attach(&self, _history_limit: usize) -> Result<FrontendAttachment, SdkError> {
Ok(FrontendAttachment::from_snapshot(
FrontendAttachSnapshot {
descriptor: Self::descriptor(),
history: vec![ChatMessage::assistant("persisted history")],
history_cursor: 1,
replay: VecDeque::new(),
},
self.events.subscribe(),
))
}
async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
Ok(())
}
async fn submit(&self, _prompt: String) -> Result<String, SdkError> {
Ok(String::new())
}
async fn interrupt(&self) -> Result<bool, SdkError> {
Ok(true)
}
async fn steer(&self, _prompt: String) -> Result<(), SdkError> {
Ok(())
}
async fn respond(&self, _response: FrontendResponse) -> Result<(), SdkError> {
Ok(())
}
async fn close(&self) -> Result<(), SdkError> {
Ok(())
}
}
#[tokio::test]
async fn credential_is_redacted_distinct_and_revocable() {
let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
let mut first = registry
.issue("codex-a", RuntimeAuthorization::owner())
.unwrap();
let second = registry
.issue("codex-b", RuntimeAuthorization::observer())
.unwrap();
assert_eq!(format!("{first:?}"), "CodexClientCredential([REDACTED])");
assert_ne!(first.bearer(), second.bearer());
let grant = registry
.authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
.unwrap();
drop(grant);
registry.revoke(&mut first).await.unwrap();
assert!(registry
.authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
.is_err());
}
#[test]
fn credential_rng_failure_collision_retry_and_uniqueness_are_fail_closed() {
let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
assert!(matches!(
registry.issue_with_fill("failed", RuntimeAuthorization::observer(), |_| Err(
CredentialError::RandomSource
)),
Err(CredentialError::RandomSource)
));
assert!(registry
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty());
let first = registry
.issue_with_fill("first", RuntimeAuthorization::observer(), |secret| {
secret.fill(0x11);
Ok(())
})
.unwrap();
let mut attempt = 0usize;
let second = registry
.issue_with_fill("second", RuntimeAuthorization::observer(), |secret| {
attempt += 1;
secret.fill(if attempt == 1 { 0x11 } else { 0x22 });
Ok(())
})
.unwrap();
assert_eq!(attempt, 2, "collision must consume a fresh RNG fill");
assert_ne!(first.bearer(), second.bearer());
let mut unique = BTreeSet::new();
for index in 0..4096u64 {
let credential = registry
.issue(format!("client-{index}"), RuntimeAuthorization::observer())
.unwrap();
assert!(unique.insert(credential.digest));
}
}
#[test]
fn application_message_credentials_are_rejected_by_field_not_plain_text() {
assert!(contains_application_credential(
&serde_json::json!({"id":1,"method":"turn/start","params":{"token":"secret"}})
));
assert!(contains_application_credential(
&serde_json::json!({"authorization":"Bearer secret"})
));
assert!(!contains_application_credential(
&serde_json::json!({"input":[{"type":"text","text":"discuss token budgets"}]})
));
}
#[cfg(unix)]
#[tokio::test]
async fn child_only_credential_is_absent_from_argv_and_reused_command_state() {
use std::process::Stdio;
let registry = CodexCredentialRegistry::new("runtime".into(), [8; 16]);
let credential = registry
.issue("spawn-test", RuntimeAuthorization::observer())
.unwrap();
let expected = credential.bearer_value();
let variable = "SUPERCODE_CODEX_TEST_CHILD_ONLY";
let mut command = tokio::process::Command::new("sh");
command
.args(["-c", "printf %s \"$SUPERCODE_CODEX_TEST_CHILD_ONLY\""])
.stdout(Stdio::piped());
let first = credential
.spawn_tokio_child(&mut command, variable)
.unwrap()
.wait_with_output()
.await
.unwrap();
assert!(first.status.success());
assert!(
first.stdout == expected.as_bytes(),
"the child did not receive its exact scoped credential"
);
assert!(
!format!("{command:?}").contains(&expected),
"the reusable command debug state retained the credential"
);
let second = command.spawn().unwrap().wait_with_output().await.unwrap();
assert!(second.status.success());
assert!(
second.stdout.is_empty(),
"the reusable command leaked the credential into a later child"
);
}
#[cfg(unix)]
#[tokio::test]
async fn child_only_credential_is_absent_from_parent_process_table_and_diagnostics() {
use std::process::Stdio;
let registry = CodexCredentialRegistry::new("runtime".into(), [11; 16]);
let credential = registry
.issue("process-scan", RuntimeAuthorization::observer())
.unwrap();
let expected = credential.bearer_value();
let variable = "SUPERCODE_CODEX_TEST_PROCESS_SCAN";
assert!(std::env::var_os(variable).is_none());
let root = std::env::temp_dir().join(format!(
"supercode-codex-process-scan-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
std::fs::create_dir(&root).unwrap();
let child_only_output = root.join("child-only");
let mut command = tokio::process::Command::new("sh");
command
.args([
"-c",
"printf %s \"$SUPERCODE_CODEX_TEST_PROCESS_SCAN\" > \"$1\"; printf crash-diagnostic >&2; sleep 1",
"sh",
])
.arg(&child_only_output)
.stdout(Stdio::null())
.stderr(Stdio::piped());
let child = credential
.spawn_tokio_child(&mut command, variable)
.unwrap();
let pid = child.id().unwrap();
assert!(std::env::var_os(variable).is_none());
assert!(!format!("{command:?}").contains(&expected));
let process_table = tokio::process::Command::new("ps")
.args(["-o", "command=", "-p", &pid.to_string()])
.output()
.await
.unwrap();
assert!(process_table.status.success());
assert!(!String::from_utf8_lossy(&process_table.stdout).contains(&expected));
let output = child.wait_with_output().await.unwrap();
assert!(output.status.success());
assert_eq!(
std::fs::read(&child_only_output).unwrap(),
expected.as_bytes()
);
assert_eq!(String::from_utf8_lossy(&output.stderr), "crash-diagnostic");
assert!(!output
.stderr
.windows(expected.len())
.any(|bytes| bytes == expected.as_bytes()));
std::fs::remove_dir_all(root).unwrap();
}
#[cfg(windows)]
#[tokio::test]
async fn windows_child_credential_is_absent_from_process_table_and_diagnostics() {
use std::process::Stdio;
let registry = CodexCredentialRegistry::new("runtime".into(), [12; 16]);
let credential = registry
.issue("windows-process-scan", RuntimeAuthorization::observer())
.unwrap();
let expected = credential.bearer_value();
let variable = "SUPERCODE_CODEX_TEST_WINDOWS_PROCESS_SCAN";
assert!(std::env::var_os(variable).is_none());
let root = std::env::temp_dir().join(format!(
"supercode-codex-windows-process-scan-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
std::fs::create_dir(&root).unwrap();
let child_only_output = root.join("child-only");
let output_literal = child_only_output.display().to_string().replace('\'', "''");
let script = format!(
"[IO.File]::WriteAllText('{output_literal}', $env:{variable}); \
[Console]::Error.Write('crash-diagnostic'); Start-Sleep -Seconds 5"
);
let mut command = tokio::process::Command::new("powershell.exe");
command
.args(["-NoProfile", "-NonInteractive", "-Command"])
.arg(script)
.stdout(Stdio::null())
.stderr(Stdio::piped());
let child = credential
.spawn_tokio_child(&mut command, variable)
.unwrap();
let pid = child.id().unwrap();
assert!(std::env::var_os(variable).is_none());
assert!(!format!("{command:?}").contains(&expected));
let process_table = tokio::process::Command::new("powershell.exe")
.args([
"-NoProfile",
"-NonInteractive",
"-Command",
&format!("(Get-CimInstance Win32_Process -Filter 'ProcessId = {pid}').CommandLine"),
])
.output()
.await
.unwrap();
assert!(process_table.status.success());
let process_command = String::from_utf8_lossy(&process_table.stdout);
assert!(
process_command.contains(variable),
"CIM scan did not capture the live child command line: {process_command:?}"
);
assert!(!process_command.contains(&expected));
let output = child.wait_with_output().await.unwrap();
assert!(output.status.success());
assert_eq!(
std::fs::read(&child_only_output).unwrap(),
expected.as_bytes()
);
assert_eq!(
String::from_utf8_lossy(&output.stderr).trim(),
"crash-diagnostic"
);
assert!(!output
.stderr
.windows(expected.len())
.any(|bytes| bytes == expected.as_bytes()));
std::fs::remove_dir_all(root).unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn bootstrap_file_is_private_required_and_deleted() {
use std::os::unix::fs::PermissionsExt;
let root = std::env::temp_dir().join(format!(
"supercode-bootstrap-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
std::fs::create_dir(&root).unwrap();
std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700)).unwrap();
let (endpoint, mut bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let mut bootstrap_file = endpoint.write_bootstrap_file(&bootstrap, &root).unwrap();
let path = bootstrap_file.path.clone();
assert_eq!(
std::fs::metadata(&path).unwrap().permissions().mode() & 0o777,
0o600
);
assert_eq!(
format!("{bootstrap:?}"),
"CodexBootstrapCredential([REDACTED])"
);
let mut client = endpoint
.issue_interactive_from_file(&bootstrap_file, "stock")
.unwrap();
let (_foreign, foreign_bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "other", "/workspace").unwrap();
assert!(matches!(
endpoint.issue_interactive(&foreign_bootstrap, "forbidden"),
Err(CredentialError::InvalidBootstrap)
));
endpoint
.revoke_frontend_from_file(&bootstrap_file, &mut client)
.await
.unwrap();
let (replacement, replacement_file) = endpoint
.rotate_bootstrap_file(&mut bootstrap, &mut bootstrap_file, &root)
.unwrap();
assert!(!path.exists(), "rotation must delete the prior file");
assert!(matches!(
endpoint.issue_observer(&bootstrap, "old-bootstrap"),
Err(CredentialError::InvalidBootstrap)
));
assert_eq!(
std::fs::metadata(&replacement_file.path)
.unwrap()
.permissions()
.mode()
& 0o777,
0o600
);
let replacement_client = endpoint
.issue_interactive_from_file(&replacement_file, "replacement")
.unwrap();
drop(replacement_client);
drop(replacement_file);
drop(replacement);
drop(bootstrap_file);
assert!(!path.exists());
std::fs::remove_dir(root).unwrap();
}
#[cfg(unix)]
#[test]
fn bootstrap_file_rotation_rolls_back_when_old_file_cannot_be_removed() {
use std::os::unix::fs::PermissionsExt;
let root = std::env::temp_dir().join(format!(
"supercode-bootstrap-rollback-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
std::fs::create_dir(&root).unwrap();
std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700)).unwrap();
let (endpoint, mut bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let mut bootstrap_file = endpoint.write_bootstrap_file(&bootstrap, &root).unwrap();
let original_path = bootstrap_file.path.clone();
let error = match endpoint.rotate_bootstrap_file_with_remove(
&mut bootstrap,
&mut bootstrap_file,
root.clone(),
|_| {
Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"injected deletion failure",
))
},
) {
Ok(_) => panic!("injected removal failure must abort rotation"),
Err(error) => error,
};
assert!(matches!(error, CredentialError::Io(_)));
assert!(original_path.exists());
assert!(endpoint
.issue_observer(&bootstrap, "still-authorized")
.is_ok());
let files = std::fs::read_dir(&root)
.unwrap()
.filter_map(Result::ok)
.map(|entry| entry.path())
.collect::<Vec<_>>();
assert_eq!(files, vec![original_path.clone()]);
drop(bootstrap_file);
assert!(!original_path.exists());
std::fs::remove_dir(root).unwrap();
}
#[test]
fn in_memory_bootstrap_rotation_invalidates_old_value_before_return() {
let (endpoint, mut bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let replacement = endpoint.rotate_bootstrap(&mut bootstrap).unwrap();
assert!(matches!(
endpoint.issue_observer(&bootstrap, "old"),
Err(CredentialError::InvalidBootstrap)
));
assert!(endpoint.issue_observer(&replacement, "new").is_ok());
}
#[test]
fn upgrade_validation_is_fail_closed_and_permission_narrowing_only() {
let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
let credential = registry
.issue("codex", RuntimeAuthorization::owner())
.unwrap();
let request = |uri: &str, authorization: Option<&str>| {
let mut builder = Request::builder().uri(uri).header("host", "127.0.0.1:1234");
if let Some(authorization) = authorization {
builder = builder.header("authorization", authorization);
}
builder.body(()).unwrap()
};
let valid = || request("/", Some(&credential.bearer()));
let rule = || HostRule::Tcp {
expected_host: "127.0.0.1:1234".into(),
forwarded_ports: BTreeSet::from([4321]),
};
let grant = validate_upgrade(&valid(), &rule(), ®istry).unwrap();
assert!(grant.authorization.allows(RuntimePermission::Terminate));
drop(grant);
let unauthenticated = [
request("/", None),
request("/", Some("Basic nope")),
request("/", Some("Bearer malformed")),
request("/", Some("Bearer codex.00")),
];
let mut external = Vec::new();
for request in unauthenticated {
let (status, message) = match validate_upgrade(&request, &rule(), ®istry) {
Ok(_) => panic!("invalid credential must be rejected"),
Err(error) => error,
};
assert_eq!(status, StatusCode::UNAUTHORIZED);
let response = http_error(status, message);
assert!(response
.headers()
.get("access-control-allow-origin")
.is_none());
external.push(response.body().clone());
}
assert!(external.windows(2).all(|pair| pair[0] == pair[1]));
let actual = credential.bearer();
let wrong_client = actual.replacen("Bearer codex.", "Bearer impostor.", 1);
let (status, _) =
match validate_upgrade(&request("/", Some(&wrong_client)), &rule(), ®istry) {
Ok(_) => panic!("mismatched client id must be rejected"),
Err(error) => error,
};
assert_eq!(status, StatusCode::FORBIDDEN);
let mut narrowed = valid();
narrowed.headers_mut().insert(
"x-supercode-permissions",
"observe,terminate".parse().unwrap(),
);
let grant = validate_upgrade(&narrowed, &rule(), ®istry).unwrap();
assert!(grant.authorization.allows(RuntimePermission::Observe));
assert!(grant.authorization.allows(RuntimePermission::Terminate));
assert!(!grant.authorization.allows(RuntimePermission::Interact));
drop(grant);
let interactive = registry
.issue(
"interactive",
RuntimeAuthorization::new([
RuntimePermission::Observe,
RuntimePermission::Interact,
]),
)
.unwrap();
let mut cannot_expand = request("/", Some(&interactive.bearer()));
cannot_expand.headers_mut().insert(
"x-supercode-permissions",
"observe,terminate".parse().unwrap(),
);
let grant = validate_upgrade(&cannot_expand, &rule(), ®istry).unwrap();
assert!(grant.authorization.allows(RuntimePermission::Observe));
assert!(!grant.authorization.allows(RuntimePermission::Terminate));
drop(grant);
let mut forwarded = valid();
forwarded
.headers_mut()
.insert("host", "127.0.0.1:4321".parse().unwrap());
drop(validate_upgrade(&forwarded, &rule(), ®istry).unwrap());
let mut rejected = Vec::new();
rejected.push(request("/?token=nope", Some(&credential.bearer())));
let mut cookie = valid();
cookie
.headers_mut()
.insert("cookie", "x=y".parse().unwrap());
rejected.push(cookie);
let mut origin = valid();
origin
.headers_mut()
.insert("origin", "https://evil.invalid".parse().unwrap());
rejected.push(origin);
let mut foreign_host = valid();
foreign_host
.headers_mut()
.insert("host", "localhost:1234".parse().unwrap());
rejected.push(foreign_host);
for carrier in [
"content-length",
"transfer-encoding",
"sec-websocket-protocol",
] {
let mut forbidden = valid();
forbidden
.headers_mut()
.insert(carrier, "1".parse().unwrap());
rejected.push(forbidden);
}
for duplicate in [
"host",
"origin",
"authorization",
"content-length",
"x-supercode-permissions",
] {
let mut request = valid();
request
.headers_mut()
.append(duplicate, "duplicate".parse().unwrap());
rejected.push(request);
}
for request in rejected {
assert!(validate_upgrade(&request, &rule(), ®istry).is_err());
}
let absolute = Request::builder()
.uri("ws://127.0.0.1:1234/")
.header("host", "127.0.0.1:1234")
.header("authorization", credential.bearer())
.body(())
.unwrap();
assert!(validate_upgrade(&absolute, &rule(), ®istry).is_err());
}
#[tokio::test]
async fn half_open_handshakes_are_deadlined_and_connection_budget_recovers() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
let (endpoint, bootstrap) =
CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
let credential = endpoint
.issue_interactive(&bootstrap, "bounded-client")
.unwrap();
let handle = endpoint.bind_websocket(0).await.unwrap();
let address = match handle.kind() {
CodexEndpointKind::WebSocket(address) => *address,
#[cfg(unix)]
CodexEndpointKind::Unix(_) => unreachable!(),
};
let mut oversized = TcpStream::connect(address).await.unwrap();
let oversized_request = format!(
"GET / HTTP/1.1\r\nHost: {address}\r\nX-Oversized: {}\r\n\r\n",
"x".repeat(MAX_HANDSHAKE_BYTES)
);
oversized
.write_all(oversized_request.as_bytes())
.await
.unwrap();
let mut byte = [0u8; 1];
let oversized_closed = tokio::time::timeout(
std::time::Duration::from_millis(100),
oversized.read(&mut byte),
)
.await
.expect("an oversized handshake must be rejected before its deadline");
assert!(matches!(oversized_closed, Ok(0) | Err(_)));
let mut held = Vec::with_capacity(MAX_CONCURRENT_CONNECTIONS);
for _ in 0..MAX_CONCURRENT_CONNECTIONS {
held.push(TcpStream::connect(address).await.unwrap());
}
let mut overflow = TcpStream::connect(address).await.unwrap();
let overflow_closed = tokio::time::timeout(
std::time::Duration::from_millis(100),
overflow.read(&mut byte),
)
.await
.expect("an over-budget connection must be rejected promptly");
assert!(matches!(overflow_closed, Ok(0) | Err(_)));
tokio::time::sleep(HANDSHAKE_TIMEOUT + std::time::Duration::from_millis(100)).await;
for mut stream in held {
let closed = tokio::time::timeout(
std::time::Duration::from_millis(100),
stream.read(&mut byte),
)
.await
.expect("a half-open handshake must be closed at its deadline");
assert!(matches!(closed, Ok(0) | Err(_)));
}
let mut request = format!("ws://{address}/").into_client_request().unwrap();
request
.headers_mut()
.insert("authorization", credential.bearer().parse().unwrap());
let (mut websocket, _) = tokio_tungstenite::connect_async(request).await.unwrap();
websocket.close(None).await.unwrap();
handle.shutdown().await;
}
#[tokio::test]
async fn authenticated_loopback_websocket_is_ready_and_revocation_closes_it() {
let (endpoint, bootstrap) =
CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
let mut credential = endpoint
.issue_interactive(&bootstrap, "stock-codex")
.unwrap();
let handle = endpoint.bind_websocket(0).await.unwrap();
assert_eq!(handle.health(), CodexEndpointHealth::Ready);
let address = match handle.kind() {
CodexEndpointKind::WebSocket(address) => *address,
#[cfg(unix)]
CodexEndpointKind::Unix(_) => unreachable!(),
};
let mut request = format!("ws://{address}/").into_client_request().unwrap();
request
.headers_mut()
.insert("authorization", credential.bearer().parse().unwrap());
let (mut websocket, _) = tokio_tungstenite::connect_async(request).await.unwrap();
let mut competing = format!("ws://{address}/").into_client_request().unwrap();
competing
.headers_mut()
.insert("authorization", credential.bearer().parse().unwrap());
let competing_error = tokio_tungstenite::connect_async(competing)
.await
.expect_err("the same credential cannot own two live channels");
match competing_error {
tokio_tungstenite::tungstenite::Error::Http(response) => {
assert_eq!(response.status(), StatusCode::CONFLICT);
}
other => panic!("expected HTTP busy rejection, got {other}"),
}
websocket
.send(Message::Text("not-json".into()))
.await
.unwrap();
let parse_error = websocket
.next()
.await
.unwrap()
.unwrap()
.into_text()
.unwrap();
let parse_error: serde_json::Value = serde_json::from_str(&parse_error).unwrap();
assert_eq!(parse_error["error"]["code"], -32700);
assert_eq!(parse_error["error"]["data"]["name"], "parse_error");
websocket
.send(Message::Text(
serde_json::json!({"id":"initialize","method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
.to_string()
.into(),
))
.await
.unwrap();
let response = websocket
.next()
.await
.unwrap()
.unwrap()
.into_text()
.unwrap();
assert_eq!(
serde_json::from_str::<serde_json::Value>(&response).unwrap()["id"],
"initialize"
);
let notification = websocket
.next()
.await
.unwrap()
.unwrap()
.into_text()
.unwrap();
assert_eq!(
serde_json::from_str::<serde_json::Value>(¬ification).unwrap()["method"],
"remoteControl/status/changed"
);
endpoint
.revoke_frontend(&bootstrap, &mut credential)
.await
.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(1), async {
loop {
match websocket.next().await {
None | Some(Ok(Message::Close(_))) | Some(Err(_)) => break,
_ => {}
}
}
})
.await
.expect("revocation must close the active channel promptly");
handle.shutdown().await;
}
#[tokio::test]
async fn websocket_reconnect_preserves_deterministic_thread_identity() {
let (endpoint, bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let credential = endpoint
.issue_interactive(&bootstrap, "stock-codex")
.unwrap();
let handle = endpoint.bind_websocket(0).await.unwrap();
let address = match handle.kind() {
CodexEndpointKind::WebSocket(address) => *address,
#[cfg(unix)]
CodexEndpointKind::Unix(_) => unreachable!(),
};
let connect = || {
let bearer = credential.bearer();
async move {
let mut request = format!("ws://{address}/").into_client_request().unwrap();
request
.headers_mut()
.insert("authorization", bearer.parse().unwrap());
tokio_tungstenite::connect_async(request).await.unwrap().0
}
};
let mut first = connect().await;
first
.send(Message::Text(
json!({"id":1,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
.to_string()
.into(),
))
.await
.unwrap();
let _initialize = first.next().await.unwrap().unwrap();
let _status = first.next().await.unwrap().unwrap();
first
.send(Message::Text(
json!({"method":"initialized"}).to_string().into(),
))
.await
.unwrap();
first
.send(Message::Text(
json!({"id":2,"method":"thread/start","params":{}})
.to_string()
.into(),
))
.await
.unwrap();
let started: Value =
serde_json::from_str(&first.next().await.unwrap().unwrap().into_text().unwrap())
.unwrap();
let thread_id = started["result"]["thread"]["id"]
.as_str()
.unwrap()
.to_owned();
let _thread_started = first.next().await.unwrap().unwrap();
first.close(None).await.unwrap();
let mut inactive = {
let records = endpoint
.credentials
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
records
.get(&credential.digest)
.unwrap()
.active_channel
.subscribe()
};
tokio::time::timeout(std::time::Duration::from_secs(1), async {
while *inactive.borrow() {
inactive.changed().await.unwrap();
}
})
.await
.expect("closed socket must release its channel");
let mut second = connect().await;
second
.send(Message::Text(
json!({"id":3,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
.to_string()
.into(),
))
.await
.unwrap();
let _initialize = second.next().await.unwrap().unwrap();
let _status = second.next().await.unwrap().unwrap();
second
.send(Message::Text(
json!({"method":"initialized"}).to_string().into(),
))
.await
.unwrap();
second
.send(Message::Text(
json!({"id":4,"method":"thread/read","params":{"threadId":thread_id}})
.to_string()
.into(),
))
.await
.unwrap();
let resumed: Value =
serde_json::from_str(&second.next().await.unwrap().unwrap().into_text().unwrap())
.unwrap();
assert_eq!(resumed["result"]["thread"]["id"], thread_id);
second.close(None).await.unwrap();
handle.shutdown().await;
}
#[tokio::test]
async fn observer_and_competing_controller_fail_with_stable_codes() {
async fn initialize_and_attach(
adapter: Arc<CodexAppServerAdapter>,
) -> (crate::codex_app_server_v0_144::CodexConnection, String) {
let mut connection = adapter.connection();
connection
.handle(json!({"id":1,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}}))
.await;
connection.handle(json!({"method":"initialized"})).await;
let started = connection
.handle(json!({"id":2,"method":"thread/start","params":{}}))
.await;
let thread_id = started[0]["result"]["thread"]["id"]
.as_str()
.unwrap()
.to_owned();
(connection, thread_id)
}
let (endpoint, bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let observer = endpoint.issue_observer(&bootstrap, "observer").unwrap();
let first = endpoint.issue_interactive(&bootstrap, "first").unwrap();
let second = endpoint.issue_interactive(&bootstrap, "second").unwrap();
let observer_grant = endpoint
.credentials
.authenticate(observer.bearer().strip_prefix("Bearer ").unwrap())
.unwrap();
let first_grant = endpoint
.credentials
.authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
.unwrap();
let second_grant = endpoint
.credentials
.authenticate(second.bearer().strip_prefix("Bearer ").unwrap())
.unwrap();
let (mut observer_connection, observer_thread) =
initialize_and_attach(endpoint.authenticated_adapter(&observer_grant)).await;
let denied = observer_connection
.handle(json!({"id":3,"method":"turn/start","params":{"threadId":observer_thread,"input":[{"type":"text","text":"forbidden"}]}}))
.await;
assert_eq!(denied[0]["error"]["code"], -32031);
assert_eq!(denied[0]["error"]["data"]["name"], "unauthorized");
let (mut first_connection, first_thread) =
initialize_and_attach(endpoint.authenticated_adapter(&first_grant)).await;
let accepted = first_connection
.handle(json!({"id":4,"method":"turn/start","params":{"threadId":first_thread,"input":[{"type":"text","text":"claim"}]}}))
.await;
assert!(accepted[0].get("result").is_some());
let (mut second_connection, second_thread) =
initialize_and_attach(endpoint.authenticated_adapter(&second_grant)).await;
let conflict = second_connection
.handle(json!({"id":5,"method":"turn/start","params":{"threadId":second_thread,"input":[{"type":"text","text":"compete"}]}}))
.await;
assert_eq!(conflict[0]["error"]["code"], -32032);
assert_eq!(conflict[0]["error"]["data"]["name"], "controller_required");
}
#[cfg(unix)]
#[tokio::test]
async fn authenticated_unix_websocket_uses_protected_files() {
use std::os::unix::fs::PermissionsExt;
let (endpoint, bootstrap) =
CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
let credential = endpoint
.issue_interactive(&bootstrap, "stock-codex-unix")
.unwrap();
let root = std::env::temp_dir().join(format!(
"supercode-codex-unix-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
let socket = root.join("adapter.sock");
let handle = endpoint.bind_unix(&socket).await.unwrap();
assert_eq!(
std::fs::metadata(&root).unwrap().permissions().mode() & 0o777,
0o700
);
assert_eq!(
std::fs::metadata(&socket).unwrap().permissions().mode() & 0o777,
0o600
);
let stream = UnixStream::connect(&socket).await.unwrap();
let mut request = "ws://localhost/".into_client_request().unwrap();
request
.headers_mut()
.insert("authorization", credential.bearer().parse().unwrap());
let (mut websocket, _) = tokio_tungstenite::client_async(request, stream)
.await
.unwrap();
websocket
.send(Message::Text(
serde_json::json!({"id":"initialize","method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
.to_string()
.into(),
))
.await
.unwrap();
assert!(websocket.next().await.unwrap().unwrap().is_text());
handle.shutdown().await;
assert!(!socket.exists());
std::fs::remove_dir(&root).unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn unix_binding_refuses_precreated_destinations_and_symlink_parents() {
use std::os::unix::fs::{symlink, PermissionsExt};
let (endpoint, _bootstrap) =
CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
let base = std::env::temp_dir().join(format!(
"supercode-codex-precreate-{}-{}",
std::process::id(),
epoch_millis_for_test()
));
std::fs::create_dir(&base).unwrap();
std::fs::set_permissions(&base, std::fs::Permissions::from_mode(0o700)).unwrap();
let occupied = base.join("occupied.sock");
std::fs::write(&occupied, b"do not replace").unwrap();
let error = match endpoint.bind_unix(&occupied).await {
Ok(_) => panic!("precreated destination must be rejected"),
Err(error) => error,
};
assert_eq!(error.kind(), std::io::ErrorKind::AlreadyExists);
assert_eq!(std::fs::read(&occupied).unwrap(), b"do not replace");
let real_parent = base.join("real");
std::fs::create_dir(&real_parent).unwrap();
std::fs::set_permissions(&real_parent, std::fs::Permissions::from_mode(0o700)).unwrap();
let linked_parent = base.join("linked");
symlink(&real_parent, &linked_parent).unwrap();
let error = match endpoint
.bind_unix(linked_parent.join("endpoint.sock"))
.await
{
Ok(_) => panic!("symlink parent must be rejected"),
Err(error) => error,
};
assert_eq!(error.kind(), std::io::ErrorKind::PermissionDenied);
assert!(!real_parent.join("endpoint.sock").exists());
std::fs::remove_file(linked_parent).unwrap();
std::fs::remove_dir(real_parent).unwrap();
std::fs::remove_file(occupied).unwrap();
std::fs::remove_dir(base).unwrap();
}
fn epoch_millis_for_test() -> u128 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis()
}
}