1use std::collections::{BTreeMap, BTreeSet};
7use std::fmt;
8use std::net::{IpAddr, Ipv4Addr, SocketAddr};
9use std::path::PathBuf;
10use std::sync::atomic::{AtomicBool, Ordering};
11use std::sync::{Arc, Mutex};
12
13use futures::{SinkExt, StreamExt};
14use supercode::{
15 CoordinatedRuntime, RuntimeAuthorization, RuntimeClientId, SdkRuntime,
16 DEFAULT_RUNTIME_LEASE_TTL_MS,
17};
18use tokio::net::{TcpListener, TcpStream};
19#[cfg(unix)]
20use tokio::net::{UnixListener, UnixStream};
21use tokio::sync::{watch, Notify};
22use tokio_tungstenite::tungstenite::handshake::server::{ErrorResponse, Request, Response};
23use tokio_tungstenite::tungstenite::http::StatusCode;
24use tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
25use tokio_tungstenite::tungstenite::Message;
26use zeroize::Zeroize;
27use zeroize::Zeroizing;
28
29use crate::codex_app_server_v0_144::{CodexAppServerAdapter, CodexCompatibilityMode};
30
31const TOKEN_BYTES: usize = 32;
32const MAX_MESSAGE_BYTES: usize = 16 * 1024 * 1024;
33const MAX_FRAME_BYTES: usize = 1024 * 1024;
34const MAX_WRITE_BUFFER_BYTES: usize = 2 * MAX_MESSAGE_BYTES;
35const MAX_HANDSHAKE_BYTES: usize = 64 * 1024;
36const MAX_CONCURRENT_CONNECTIONS: usize = 32;
37#[cfg(not(test))]
38const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
39#[cfg(test)]
40const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
41
42#[derive(Debug, thiserror::Error)]
43pub enum CredentialError {
44 #[error("operating system random source failed")]
45 RandomSource,
46 #[error("credential generation collided repeatedly")]
47 Collision,
48 #[error("credential is already revoked")]
49 Revoked,
50 #[error("bootstrap credential is invalid for this runtime generation")]
51 InvalidBootstrap,
52 #[error("invalid runtime client identity: {0}")]
53 InvalidClient(String),
54 #[error("private credential file operation failed: {0}")]
55 Io(#[from] std::io::Error),
56 #[error("endpoint configuration is already shared")]
57 ConfigurationLocked,
58}
59
60struct CredentialGrant {
61 client_id: RuntimeClientId,
62 authorization: RuntimeAuthorization,
63 revoked: watch::Receiver<bool>,
64 runtime_id: String,
65 generation: [u8; 16],
66 active_channel: watch::Sender<bool>,
67}
68
69impl Drop for CredentialGrant {
70 fn drop(&mut self) {
71 self.active_channel.send_replace(false);
72 }
73}
74
75struct CredentialRecord {
76 client_id: RuntimeClientId,
77 authorization: RuntimeAuthorization,
78 revoked: watch::Sender<bool>,
79 active_channel: watch::Sender<bool>,
80}
81
82struct CodexCredentialRegistry {
84 records: Mutex<BTreeMap<[u8; 32], CredentialRecord>>,
85 runtime_id: String,
86 generation: [u8; 16],
87}
88
89pub struct CodexBootstrapCredential {
92 secret: [u8; TOKEN_BYTES],
93 digest: [u8; 32],
94 generation: [u8; 16],
95 revoked: bool,
96}
97
98impl fmt::Debug for CodexBootstrapCredential {
99 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
100 formatter.write_str("CodexBootstrapCredential([REDACTED])")
101 }
102}
103
104impl Drop for CodexBootstrapCredential {
105 fn drop(&mut self) {
106 zero(&mut self.secret);
107 }
108}
109
110pub struct CodexBootstrapFile {
114 path: PathBuf,
115 generation: [u8; 16],
116}
117
118impl Drop for CodexBootstrapFile {
119 fn drop(&mut self) {
120 let _ = std::fs::remove_file(&self.path);
121 }
122}
123
124impl CodexCredentialRegistry {
125 fn new(runtime_id: String, generation: [u8; 16]) -> Arc<Self> {
126 Arc::new(Self {
127 records: Mutex::new(BTreeMap::new()),
128 runtime_id,
129 generation,
130 })
131 }
132
133 fn issue(
135 self: &Arc<Self>,
136 client_id: impl Into<String>,
137 authorization: RuntimeAuthorization,
138 ) -> Result<CodexClientCredential, CredentialError> {
139 self.issue_with_fill(client_id, authorization, |secret| {
140 getrandom::getrandom(secret).map_err(|_| CredentialError::RandomSource)
141 })
142 }
143
144 fn issue_with_fill(
145 self: &Arc<Self>,
146 client_id: impl Into<String>,
147 authorization: RuntimeAuthorization,
148 mut fill: impl FnMut(&mut [u8; TOKEN_BYTES]) -> Result<(), CredentialError>,
149 ) -> Result<CodexClientCredential, CredentialError> {
150 let client_id = RuntimeClientId::parse(client_id.into())
151 .map_err(|error| CredentialError::InvalidClient(error.to_string()))?;
152 for _ in 0..3 {
153 let mut secret = [0u8; TOKEN_BYTES];
154 fill(&mut secret)?;
155 let digest = credential_digest(&secret);
156 let mut records = self
157 .records
158 .lock()
159 .unwrap_or_else(std::sync::PoisonError::into_inner);
160 if records.contains_key(&digest) {
161 zero(&mut secret);
162 continue;
163 }
164 let (revoked, _) = watch::channel(false);
165 let (active_channel, _) = watch::channel(false);
166 records.insert(
167 digest,
168 CredentialRecord {
169 client_id: client_id.clone(),
170 authorization,
171 revoked,
172 active_channel,
173 },
174 );
175 return Ok(CodexClientCredential {
176 secret,
177 digest,
178 client_id,
179 revoked: false,
180 });
181 }
182 Err(CredentialError::Collision)
183 }
184
185 async fn revoke(&self, credential: &mut CodexClientCredential) -> Result<(), CredentialError> {
186 if credential.revoked {
187 return Err(CredentialError::Revoked);
188 }
189 let record = self
190 .records
191 .lock()
192 .unwrap_or_else(std::sync::PoisonError::into_inner)
193 .remove(&credential.digest)
194 .ok_or(CredentialError::Revoked)?;
195 let mut active_channel = record.active_channel.subscribe();
196 let _ = record.revoked.send(true);
197 while *active_channel.borrow() {
198 if active_channel.changed().await.is_err() {
199 break;
200 }
201 }
202 credential.revoked = true;
203 zero(&mut credential.secret);
204 Ok(())
205 }
206
207 async fn rotate(
208 self: &Arc<Self>,
209 credential: &mut CodexClientCredential,
210 client_id: impl Into<String>,
211 authorization: RuntimeAuthorization,
212 ) -> Result<CodexClientCredential, CredentialError> {
213 let mut replacement = self.issue(client_id, authorization)?;
214 if let Err(error) = self.revoke(credential).await {
215 let _ = self.revoke(&mut replacement).await;
216 return Err(error);
217 }
218 Ok(replacement)
219 }
220
221 async fn revoke_all(&self) {
222 let records = {
223 let mut records = self
224 .records
225 .lock()
226 .unwrap_or_else(std::sync::PoisonError::into_inner);
227 std::mem::take(&mut *records)
228 };
229 let mut channels = Vec::with_capacity(records.len());
230 for (_, record) in records {
231 let mut active = record.active_channel.subscribe();
232 let _ = record.revoked.send(true);
233 channels.push(async move {
234 while *active.borrow() {
235 if active.changed().await.is_err() {
236 break;
237 }
238 }
239 });
240 }
241 futures::future::join_all(channels).await;
242 }
243
244 fn authenticate(&self, token: &str) -> Result<CredentialGrant, AuthenticationError> {
245 let (presented_client_id, secret_text) = token
246 .split_once('.')
247 .ok_or(AuthenticationError::Unauthenticated)?;
248 let mut secret =
249 decode_hex_secret(secret_text).ok_or(AuthenticationError::Unauthenticated)?;
250 let digest = credential_digest(&secret);
251 zero(&mut secret);
252 let mut records = self
253 .records
254 .lock()
255 .unwrap_or_else(std::sync::PoisonError::into_inner);
256 let record = records
257 .get_mut(&digest)
258 .ok_or(AuthenticationError::Unauthenticated)?;
259 if record.client_id.as_str() != presented_client_id {
260 return Err(AuthenticationError::Unauthorized);
261 }
262 if *record.active_channel.borrow() {
263 return Err(AuthenticationError::Busy);
264 }
265 record.active_channel.send_replace(true);
266 Ok(CredentialGrant {
267 client_id: record.client_id.clone(),
268 authorization: record.authorization.clone(),
269 revoked: record.revoked.subscribe(),
270 runtime_id: self.runtime_id.clone(),
271 generation: self.generation,
272 active_channel: record.active_channel.clone(),
273 })
274 }
275}
276
277#[derive(Debug, Clone, Copy, PartialEq, Eq)]
278enum AuthenticationError {
279 Unauthenticated,
280 Unauthorized,
281 Busy,
282}
283
284pub struct CodexClientCredential {
287 secret: [u8; TOKEN_BYTES],
288 digest: [u8; 32],
289 client_id: RuntimeClientId,
290 revoked: bool,
291}
292
293impl CodexClientCredential {
294 pub fn spawn_tokio_child(
298 &self,
299 command: &mut tokio::process::Command,
300 variable: &str,
301 ) -> std::io::Result<tokio::process::Child> {
302 let bearer = Zeroizing::new(self.bearer_value());
303 command.env(variable, bearer.as_str());
304 let child = command.spawn();
305 command.env_remove(variable);
306 child
307 }
308
309 #[cfg(test)]
310 fn bearer(&self) -> String {
311 format!("Bearer {}", self.bearer_value())
312 }
313
314 fn bearer_value(&self) -> String {
315 format!("{}.{}", self.client_id.as_str(), encode_hex(&self.secret))
316 }
317}
318
319impl fmt::Debug for CodexClientCredential {
320 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
321 formatter.write_str("CodexClientCredential([REDACTED])")
322 }
323}
324
325impl Drop for CodexClientCredential {
326 fn drop(&mut self) {
327 zero(&mut self.secret);
328 }
329}
330
331#[derive(Debug, Clone, PartialEq, Eq)]
332pub enum CodexEndpointKind {
333 WebSocket(SocketAddr),
334 #[cfg(unix)]
335 Unix(PathBuf),
336}
337
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339pub enum CodexEndpointHealth {
340 Ready,
341 ShuttingDown,
342 Stopped,
343}
344
345pub struct CodexEndpoint {
347 coordinator: Arc<CoordinatedRuntime>,
348 credentials: Arc<CodexCredentialRegistry>,
349 bootstrap_digest: Mutex<[u8; 32]>,
350 cwd: PathBuf,
351 client_home: Option<PathBuf>,
352 mode: CodexCompatibilityMode,
353 forwarded_ports: BTreeSet<u16>,
354}
355
356impl CodexEndpoint {
357 pub fn new(
358 runtime: Arc<dyn SdkRuntime>,
359 runtime_id: impl Into<String>,
360 cwd: impl Into<PathBuf>,
361 ) -> Result<(Arc<Self>, CodexBootstrapCredential), CredentialError> {
362 let generation = random_generation()?;
363 let bootstrap = new_bootstrap(generation)?;
364 let endpoint = Arc::new(Self {
365 coordinator: CoordinatedRuntime::with_lease_ttl(runtime, DEFAULT_RUNTIME_LEASE_TTL_MS),
366 credentials: CodexCredentialRegistry::new(runtime_id.into(), generation),
367 bootstrap_digest: Mutex::new(bootstrap.digest),
368 cwd: cwd.into(),
369 client_home: None,
370 mode: CodexCompatibilityMode::TracedOnly,
371 forwarded_ports: BTreeSet::new(),
372 });
373 Ok((endpoint, bootstrap))
374 }
375
376 pub fn with_mode(
377 runtime: Arc<dyn SdkRuntime>,
378 runtime_id: impl Into<String>,
379 cwd: impl Into<PathBuf>,
380 mode: CodexCompatibilityMode,
381 ) -> Result<(Arc<Self>, CodexBootstrapCredential), CredentialError> {
382 let generation = random_generation()?;
383 let bootstrap = new_bootstrap(generation)?;
384 let endpoint = Arc::new(Self {
385 coordinator: CoordinatedRuntime::with_lease_ttl(runtime, DEFAULT_RUNTIME_LEASE_TTL_MS),
386 credentials: CodexCredentialRegistry::new(runtime_id.into(), generation),
387 bootstrap_digest: Mutex::new(bootstrap.digest),
388 cwd: cwd.into(),
389 client_home: None,
390 mode,
391 forwarded_ports: BTreeSet::new(),
392 });
393 Ok((endpoint, bootstrap))
394 }
395
396 pub fn with_forwarded_ports(
399 mut self: Arc<Self>,
400 ports: impl IntoIterator<Item = u16>,
401 ) -> Result<Arc<Self>, CredentialError> {
402 Arc::get_mut(&mut self)
403 .ok_or(CredentialError::ConfigurationLocked)?
404 .forwarded_ports
405 .extend(ports.into_iter().filter(|port| *port != 0));
406 Ok(self)
407 }
408
409 pub fn with_client_home(
412 mut self: Arc<Self>,
413 path: impl Into<PathBuf>,
414 ) -> Result<Arc<Self>, CredentialError> {
415 Arc::get_mut(&mut self)
416 .ok_or(CredentialError::ConfigurationLocked)?
417 .client_home = Some(path.into());
418 Ok(self)
419 }
420
421 pub fn issue_interactive(
422 self: &Arc<Self>,
423 bootstrap: &CodexBootstrapCredential,
424 client_id: impl Into<String>,
425 ) -> Result<CodexClientCredential, CredentialError> {
426 self.authenticate_bootstrap(bootstrap)?;
427 self.credentials.issue(
428 client_id,
429 RuntimeAuthorization::new([
430 supercode::RuntimePermission::Observe,
431 supercode::RuntimePermission::Interact,
432 supercode::RuntimePermission::Approve,
433 ]),
434 )
435 }
436
437 pub fn issue_observer(
438 self: &Arc<Self>,
439 bootstrap: &CodexBootstrapCredential,
440 client_id: impl Into<String>,
441 ) -> Result<CodexClientCredential, CredentialError> {
442 self.authenticate_bootstrap(bootstrap)?;
443 self.credentials
444 .issue(client_id, RuntimeAuthorization::observer())
445 }
446
447 pub async fn rotate_interactive(
448 self: &Arc<Self>,
449 bootstrap: &CodexBootstrapCredential,
450 credential: &mut CodexClientCredential,
451 client_id: impl Into<String>,
452 ) -> Result<CodexClientCredential, CredentialError> {
453 self.authenticate_bootstrap(bootstrap)?;
454 self.credentials
455 .rotate(
456 credential,
457 client_id,
458 RuntimeAuthorization::new([
459 supercode::RuntimePermission::Observe,
460 supercode::RuntimePermission::Interact,
461 supercode::RuntimePermission::Approve,
462 ]),
463 )
464 .await
465 }
466
467 pub fn rotate_bootstrap(
471 &self,
472 bootstrap: &mut CodexBootstrapCredential,
473 ) -> Result<CodexBootstrapCredential, CredentialError> {
474 let replacement = new_bootstrap(self.credentials.generation)?;
475 let mut expected = self
476 .bootstrap_digest
477 .lock()
478 .unwrap_or_else(std::sync::PoisonError::into_inner);
479 if bootstrap.revoked
480 || bootstrap.generation != self.credentials.generation
481 || *expected != bootstrap.digest
482 || credential_digest(&bootstrap.secret) != *expected
483 {
484 return Err(CredentialError::InvalidBootstrap);
485 }
486 *expected = replacement.digest;
487 bootstrap.revoked = true;
488 zero(&mut bootstrap.secret);
489 Ok(replacement)
490 }
491
492 #[cfg(unix)]
496 pub fn rotate_bootstrap_file(
497 &self,
498 bootstrap: &mut CodexBootstrapCredential,
499 bootstrap_file: &mut CodexBootstrapFile,
500 runtime_dir: impl Into<PathBuf>,
501 ) -> Result<(CodexBootstrapCredential, CodexBootstrapFile), CredentialError> {
502 self.rotate_bootstrap_file_with_remove(
503 bootstrap,
504 bootstrap_file,
505 runtime_dir.into(),
506 |path| std::fs::remove_file(path),
507 )
508 }
509
510 #[cfg(unix)]
511 fn rotate_bootstrap_file_with_remove(
512 &self,
513 bootstrap: &mut CodexBootstrapCredential,
514 bootstrap_file: &mut CodexBootstrapFile,
515 runtime_dir: PathBuf,
516 remove_old: impl FnOnce(&std::path::Path) -> std::io::Result<()>,
517 ) -> Result<(CodexBootstrapCredential, CodexBootstrapFile), CredentialError> {
518 self.authenticate_bootstrap(bootstrap)?;
519 self.authenticate_bootstrap_file(bootstrap_file)?;
520 let replacement = new_bootstrap(self.credentials.generation)?;
521 let replacement_file =
522 write_bootstrap_secret(runtime_dir, &replacement.secret, replacement.generation)?;
523 let mut expected = self
524 .bootstrap_digest
525 .lock()
526 .unwrap_or_else(std::sync::PoisonError::into_inner);
527 if *expected != bootstrap.digest || credential_digest(&bootstrap.secret) != *expected {
528 drop(replacement_file);
529 return Err(CredentialError::InvalidBootstrap);
530 }
531 remove_old(&bootstrap_file.path)?;
535 *expected = replacement.digest;
536 bootstrap.revoked = true;
537 zero(&mut bootstrap.secret);
538 bootstrap_file.generation = [0; 16];
539 Ok((replacement, replacement_file))
540 }
541
542 pub async fn revoke_frontend(
545 &self,
546 bootstrap: &CodexBootstrapCredential,
547 credential: &mut CodexClientCredential,
548 ) -> Result<(), CredentialError> {
549 self.authenticate_bootstrap(bootstrap)?;
550 self.credentials.revoke(credential).await
551 }
552
553 #[cfg(unix)]
556 pub fn write_bootstrap_file(
557 &self,
558 bootstrap: &CodexBootstrapCredential,
559 runtime_dir: impl Into<PathBuf>,
560 ) -> Result<CodexBootstrapFile, CredentialError> {
561 self.authenticate_bootstrap(bootstrap)?;
562 write_bootstrap_secret(runtime_dir.into(), &bootstrap.secret, bootstrap.generation)
563 }
564
565 #[cfg(unix)]
568 pub fn issue_interactive_from_file(
569 self: &Arc<Self>,
570 bootstrap_file: &CodexBootstrapFile,
571 client_id: impl Into<String>,
572 ) -> Result<CodexClientCredential, CredentialError> {
573 self.authenticate_bootstrap_file(bootstrap_file)?;
574 self.credentials.issue(
575 client_id,
576 RuntimeAuthorization::new([
577 supercode::RuntimePermission::Observe,
578 supercode::RuntimePermission::Interact,
579 supercode::RuntimePermission::Approve,
580 ]),
581 )
582 }
583
584 #[cfg(unix)]
585 pub async fn revoke_frontend_from_file(
586 &self,
587 bootstrap_file: &CodexBootstrapFile,
588 credential: &mut CodexClientCredential,
589 ) -> Result<(), CredentialError> {
590 self.authenticate_bootstrap_file(bootstrap_file)?;
591 self.credentials.revoke(credential).await
592 }
593
594 pub async fn terminate(
596 self: &Arc<Self>,
597 bootstrap: &CodexBootstrapCredential,
598 ) -> Result<(), CredentialError> {
599 self.authenticate_bootstrap(bootstrap)?;
600 self.credentials.revoke_all().await;
601 let client_id = RuntimeClientId::parse("supercode-bootstrap-operator")
602 .map_err(|error| CredentialError::InvalidClient(error.to_string()))?;
603 self.coordinator
604 .client(client_id, RuntimeAuthorization::owner())
605 .close()
606 .await
607 .map_err(|_| CredentialError::InvalidBootstrap)
608 }
609
610 fn authenticate_bootstrap(
611 &self,
612 bootstrap: &CodexBootstrapCredential,
613 ) -> Result<(), CredentialError> {
614 if bootstrap.revoked || bootstrap.generation != self.credentials.generation {
615 return Err(CredentialError::InvalidBootstrap);
616 }
617 let expected = self
618 .bootstrap_digest
619 .lock()
620 .unwrap_or_else(std::sync::PoisonError::into_inner);
621 if *expected != bootstrap.digest || credential_digest(&bootstrap.secret) != *expected {
622 return Err(CredentialError::InvalidBootstrap);
623 }
624 Ok(())
625 }
626
627 #[cfg(unix)]
628 fn authenticate_bootstrap_file(
629 &self,
630 bootstrap_file: &CodexBootstrapFile,
631 ) -> Result<(), CredentialError> {
632 if bootstrap_file.generation != self.credentials.generation {
633 return Err(CredentialError::InvalidBootstrap);
634 }
635 let mut encoded = std::fs::read(&bootstrap_file.path)?;
636 let decoded = std::str::from_utf8(&encoded)
637 .ok()
638 .and_then(decode_hex_secret)
639 .ok_or(CredentialError::InvalidBootstrap);
640 encoded.zeroize();
641 let mut secret = decoded?;
642 let digest = credential_digest(&secret);
643 zero(&mut secret);
644 let expected = self
645 .bootstrap_digest
646 .lock()
647 .unwrap_or_else(std::sync::PoisonError::into_inner);
648 if digest != *expected {
649 return Err(CredentialError::InvalidBootstrap);
650 }
651 Ok(())
652 }
653
654 pub async fn bind_websocket(self: &Arc<Self>, port: u16) -> std::io::Result<CodexServerHandle> {
656 let listener =
657 TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port)).await?;
658 let address = listener.local_addr()?;
659 let shutdown = Arc::new(Notify::new());
660 let shutting_down = Arc::new(AtomicBool::new(false));
661 let stopped = Arc::new(AtomicBool::new(false));
662 let task = tokio::spawn(run_tcp_accept_loop(
663 listener,
664 self.clone(),
665 address.to_string(),
666 self.forwarded_ports.clone(),
667 shutdown.clone(),
668 shutting_down.clone(),
669 stopped.clone(),
670 ));
671 Ok(CodexServerHandle {
672 kind: CodexEndpointKind::WebSocket(address),
673 shutdown,
674 shutting_down,
675 stopped,
676 task: Some(task),
677 })
678 }
679
680 #[cfg(unix)]
681 pub async fn bind_unix(
682 self: &Arc<Self>,
683 socket_path: impl Into<PathBuf>,
684 ) -> std::io::Result<CodexServerHandle> {
685 use std::os::unix::fs::{DirBuilderExt, MetadataExt, PermissionsExt};
686
687 let socket_path = socket_path.into();
688 let parent = socket_path.parent().ok_or_else(|| {
689 std::io::Error::new(std::io::ErrorKind::InvalidInput, "socket has no parent")
690 })?;
691 match std::fs::symlink_metadata(parent) {
692 Ok(metadata) => {
693 if metadata.file_type().is_symlink() || !metadata.is_dir() {
694 return Err(std::io::Error::new(
695 std::io::ErrorKind::PermissionDenied,
696 "Unix endpoint directory must be a real directory",
697 ));
698 }
699 if metadata.permissions().mode() & 0o777 != 0o700 {
700 return Err(std::io::Error::new(
701 std::io::ErrorKind::PermissionDenied,
702 "Unix endpoint directory must have mode 0700",
703 ));
704 }
705 if metadata.uid() != unsafe { libc::geteuid() } {
708 return Err(std::io::Error::new(
709 std::io::ErrorKind::PermissionDenied,
710 "Unix endpoint directory must be owned by this user",
711 ));
712 }
713 }
714 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
715 let mut builder = std::fs::DirBuilder::new();
716 builder.mode(0o700).create(parent)?;
717 }
718 Err(error) => return Err(error),
719 }
720 match std::fs::symlink_metadata(&socket_path) {
721 Ok(_) => {
722 return Err(std::io::Error::new(
723 std::io::ErrorKind::AlreadyExists,
724 "refusing to replace an existing Unix endpoint",
725 ))
726 }
727 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
728 Err(error) => return Err(error),
729 }
730 let listener = UnixListener::bind(&socket_path)?;
731 std::fs::set_permissions(&socket_path, std::fs::Permissions::from_mode(0o600))?;
732 let shutdown = Arc::new(Notify::new());
733 let shutting_down = Arc::new(AtomicBool::new(false));
734 let stopped = Arc::new(AtomicBool::new(false));
735 let task = tokio::spawn(run_unix_accept_loop(
736 listener,
737 self.clone(),
738 socket_path.clone(),
739 shutdown.clone(),
740 shutting_down.clone(),
741 stopped.clone(),
742 ));
743 Ok(CodexServerHandle {
744 kind: CodexEndpointKind::Unix(socket_path),
745 shutdown,
746 shutting_down,
747 stopped,
748 task: Some(task),
749 })
750 }
751
752 fn authenticated_adapter(&self, grant: &CredentialGrant) -> Arc<CodexAppServerAdapter> {
753 debug_assert_eq!(grant.runtime_id, self.credentials.runtime_id);
754 debug_assert_eq!(grant.generation, self.credentials.generation);
755 let runtime = self
756 .coordinator
757 .client(grant.client_id.clone(), grant.authorization.clone());
758 CodexAppServerAdapter::with_mode_and_client_home(
759 runtime,
760 self.cwd.clone(),
761 self.client_home.clone().unwrap_or_else(|| self.cwd.clone()),
762 self.mode,
763 )
764 }
765}
766
767pub struct CodexServerHandle {
768 kind: CodexEndpointKind,
769 shutdown: Arc<Notify>,
770 shutting_down: Arc<AtomicBool>,
771 stopped: Arc<AtomicBool>,
772 task: Option<tokio::task::JoinHandle<()>>,
773}
774
775impl CodexServerHandle {
776 pub fn kind(&self) -> &CodexEndpointKind {
777 &self.kind
778 }
779
780 pub fn health(&self) -> CodexEndpointHealth {
781 if self.stopped.load(Ordering::SeqCst) {
782 CodexEndpointHealth::Stopped
783 } else if self.shutting_down.load(Ordering::SeqCst) {
784 CodexEndpointHealth::ShuttingDown
785 } else {
786 CodexEndpointHealth::Ready
787 }
788 }
789
790 pub async fn shutdown(mut self) {
791 self.shutting_down.store(true, Ordering::SeqCst);
792 self.shutdown.notify_waiters();
793 if let Some(task) = self.task.take() {
794 let _ = task.await;
795 }
796 }
797}
798
799impl Drop for CodexServerHandle {
800 fn drop(&mut self) {
801 self.shutting_down.store(true, Ordering::SeqCst);
802 self.shutdown.notify_waiters();
803 if let Some(task) = self.task.take() {
804 task.abort();
805 }
806 }
807}
808
809async fn run_tcp_accept_loop(
810 listener: TcpListener,
811 endpoint: Arc<CodexEndpoint>,
812 expected_host: String,
813 forwarded_ports: BTreeSet<u16>,
814 shutdown: Arc<Notify>,
815 shutting_down: Arc<AtomicBool>,
816 stopped: Arc<AtomicBool>,
817) {
818 let connection_budget = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
819 loop {
820 tokio::select! {
821 _ = shutdown.notified() => break,
822 accepted = listener.accept() => match accepted {
823 Ok((stream, _)) => {
824 let Ok(permit) = connection_budget.clone().try_acquire_owned() else {
825 drop(stream);
826 continue;
827 };
828 let endpoint = endpoint.clone();
829 let expected_host = expected_host.clone();
830 let forwarded_ports = forwarded_ports.clone();
831 tokio::spawn(async move {
832 let _permit = permit;
833 let _ = serve_stream(stream, endpoint, HostRule::Tcp { expected_host, forwarded_ports }).await;
834 });
835 }
836 Err(_) => break,
837 }
838 }
839 }
840 shutting_down.store(true, Ordering::SeqCst);
841 stopped.store(true, Ordering::SeqCst);
842}
843
844#[cfg(unix)]
845async fn run_unix_accept_loop(
846 listener: UnixListener,
847 endpoint: Arc<CodexEndpoint>,
848 socket_path: PathBuf,
849 shutdown: Arc<Notify>,
850 shutting_down: Arc<AtomicBool>,
851 stopped: Arc<AtomicBool>,
852) {
853 let connection_budget = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
854 loop {
855 tokio::select! {
856 _ = shutdown.notified() => break,
857 accepted = listener.accept() => match accepted {
858 Ok((stream, _)) => {
859 let Ok(permit) = connection_budget.clone().try_acquire_owned() else {
860 drop(stream);
861 continue;
862 };
863 let endpoint = endpoint.clone();
864 tokio::spawn(async move {
865 let _permit = permit;
866 let _ = serve_stream(stream, endpoint, HostRule::Unix).await;
867 });
868 }
869 Err(_) => break,
870 }
871 }
872 }
873 let _ = std::fs::remove_file(socket_path);
874 shutting_down.store(true, Ordering::SeqCst);
875 stopped.store(true, Ordering::SeqCst);
876}
877
878enum HostRule {
879 Tcp {
880 expected_host: String,
881 forwarded_ports: BTreeSet<u16>,
882 },
883 #[cfg(unix)]
884 Unix,
885}
886
887trait LocalStream: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static {}
888impl LocalStream for TcpStream {}
889#[cfg(unix)]
890impl LocalStream for UnixStream {}
891
892struct HandshakeBoundedStream<S> {
893 inner: S,
894 bytes_read: usize,
895 complete: Arc<AtomicBool>,
896}
897
898impl<S: Unpin> Unpin for HandshakeBoundedStream<S> {}
899impl<S: LocalStream> LocalStream for HandshakeBoundedStream<S> {}
900
901impl<S: tokio::io::AsyncRead + Unpin> tokio::io::AsyncRead for HandshakeBoundedStream<S> {
902 fn poll_read(
903 mut self: std::pin::Pin<&mut Self>,
904 context: &mut std::task::Context<'_>,
905 buffer: &mut tokio::io::ReadBuf<'_>,
906 ) -> std::task::Poll<std::io::Result<()>> {
907 let before = buffer.filled().len();
908 match std::pin::Pin::new(&mut self.inner).poll_read(context, buffer) {
909 std::task::Poll::Ready(Ok(())) => {
910 if !self.complete.load(Ordering::Acquire) {
911 self.bytes_read = self
912 .bytes_read
913 .saturating_add(buffer.filled().len().saturating_sub(before));
914 if self.bytes_read > MAX_HANDSHAKE_BYTES {
915 return std::task::Poll::Ready(Err(std::io::Error::new(
916 std::io::ErrorKind::InvalidData,
917 "WebSocket handshake exceeded byte limit",
918 )));
919 }
920 }
921 std::task::Poll::Ready(Ok(()))
922 }
923 other => other,
924 }
925 }
926}
927
928impl<S: tokio::io::AsyncWrite + Unpin> tokio::io::AsyncWrite for HandshakeBoundedStream<S> {
929 fn poll_write(
930 mut self: std::pin::Pin<&mut Self>,
931 context: &mut std::task::Context<'_>,
932 buffer: &[u8],
933 ) -> std::task::Poll<std::io::Result<usize>> {
934 std::pin::Pin::new(&mut self.inner).poll_write(context, buffer)
935 }
936
937 fn poll_flush(
938 mut self: std::pin::Pin<&mut Self>,
939 context: &mut std::task::Context<'_>,
940 ) -> std::task::Poll<std::io::Result<()>> {
941 std::pin::Pin::new(&mut self.inner).poll_flush(context)
942 }
943
944 fn poll_shutdown(
945 mut self: std::pin::Pin<&mut Self>,
946 context: &mut std::task::Context<'_>,
947 ) -> std::task::Poll<std::io::Result<()>> {
948 std::pin::Pin::new(&mut self.inner).poll_shutdown(context)
949 }
950}
951
952#[allow(clippy::result_large_err)]
953async fn serve_stream<S: LocalStream>(
954 stream: S,
955 endpoint: Arc<CodexEndpoint>,
956 host_rule: HostRule,
957) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
958 let selected = Arc::new(Mutex::new(None::<CredentialGrant>));
959 let callback_selected = selected.clone();
960 let credentials = endpoint.credentials.clone();
961 let handshake_complete = Arc::new(AtomicBool::new(false));
962 let callback_complete = handshake_complete.clone();
963 let stream = HandshakeBoundedStream {
964 inner: stream,
965 bytes_read: 0,
966 complete: handshake_complete,
967 };
968 let websocket = tokio::time::timeout(
969 HANDSHAKE_TIMEOUT,
970 tokio_tungstenite::accept_hdr_async_with_config(
971 stream,
972 move |request: &Request, response: Response| -> Result<Response, ErrorResponse> {
973 match validate_upgrade(request, &host_rule, &credentials) {
974 Ok(grant) => {
975 *callback_selected
976 .lock()
977 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(grant);
978 callback_complete.store(true, Ordering::Release);
979 Ok(response)
980 }
981 Err((status, message)) => Err(http_error(status, message)),
982 }
983 },
984 Some(websocket_config()),
985 ),
986 )
987 .await
988 .map_err(|_| {
989 std::io::Error::new(
990 std::io::ErrorKind::TimedOut,
991 "WebSocket handshake timed out",
992 )
993 })??;
994 let grant = selected
995 .lock()
996 .unwrap_or_else(std::sync::PoisonError::into_inner)
997 .take()
998 .ok_or("authenticated upgrade did not select a grant")?;
999 serve_websocket(websocket, endpoint, grant).await
1000}
1001
1002async fn serve_websocket<S: LocalStream>(
1003 mut websocket: tokio_tungstenite::WebSocketStream<S>,
1004 endpoint: Arc<CodexEndpoint>,
1005 mut grant: CredentialGrant,
1006) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
1007 let runtime = endpoint.authenticated_adapter(&grant);
1008 let mut connection = runtime.connection();
1009 loop {
1010 tokio::select! {
1011 changed = grant.revoked.changed() => {
1012 if changed.is_err() || *grant.revoked.borrow() {
1013 let _ = websocket.close(None).await;
1014 break;
1015 }
1016 }
1017 incoming = websocket.next() => {
1018 let Some(message) = incoming else { break; };
1019 match message? {
1020 Message::Text(text) => {
1021 let value: serde_json::Value = match serde_json::from_str(&text) {
1022 Ok(value) => value,
1023 Err(_) => {
1024 websocket.send(Message::Text(serde_json::json!({"id":null,"error":{
1025 "code":-32700,"message":"invalid JSON request","data":{"name":"parse_error"}
1026 }}).to_string().into())).await?;
1027 continue;
1028 }
1029 };
1030 if contains_application_credential(&value) {
1031 let id = value.get("id").cloned().unwrap_or(serde_json::Value::Null);
1032 websocket.send(Message::Text(serde_json::json!({"id":id,"error":{
1033 "code":-32030,"message":"application-message credentials are forbidden",
1034 "data":{"name":"unauthenticated"}
1035 }}).to_string().into())).await?;
1036 continue;
1037 }
1038 let output = if value.get("method").is_some() {
1039 connection.handle(value).await
1040 } else {
1041 connection.handle_server_response(&value).await?
1042 };
1043 for message in output {
1044 websocket.send(Message::Text(message.to_string().into())).await?;
1045 }
1046 }
1047 Message::Close(_) => break,
1048 Message::Ping(payload) => websocket.send(Message::Pong(payload)).await?,
1049 Message::Binary(_) => {
1050 websocket.close(None).await?;
1051 break;
1052 }
1053 Message::Pong(_) | Message::Frame(_) => {}
1054 }
1055 }
1056 notifications = connection.next_notifications(), if connection_is_attached(&connection) => {
1057 for message in notifications? {
1058 websocket.send(Message::Text(message.to_string().into())).await?;
1059 }
1060 }
1061 }
1062 }
1063 let _ = runtime.detach().await;
1064 Ok(())
1065}
1066
1067fn connection_is_attached(connection: &crate::CodexConnection) -> bool {
1068 connection.is_attached()
1069}
1070
1071fn validate_upgrade(
1072 request: &Request,
1073 host_rule: &HostRule,
1074 credentials: &CodexCredentialRegistry,
1075) -> Result<CredentialGrant, (StatusCode, &'static str)> {
1076 if request.uri().path() != "/" || request.uri().query().is_some() {
1077 return Err((StatusCode::NOT_FOUND, "unsupported endpoint"));
1078 }
1079 if request.headers().contains_key("cookie")
1080 || request.headers().contains_key("sec-websocket-protocol")
1081 || request.headers().contains_key("transfer-encoding")
1082 || request.headers().contains_key("content-length")
1083 {
1084 return Err((StatusCode::BAD_REQUEST, "forbidden credential carrier"));
1085 }
1086 for name in [
1087 "host",
1088 "origin",
1089 "authorization",
1090 "content-length",
1091 "x-supercode-permissions",
1092 ] {
1093 if request.headers().get_all(name).iter().count() > 1 {
1094 return Err((StatusCode::BAD_REQUEST, "duplicate security header"));
1095 }
1096 }
1097 if request.uri().scheme().is_some() || request.uri().authority().is_some() {
1098 return Err((StatusCode::BAD_REQUEST, "absolute-form target rejected"));
1099 }
1100 let host = request
1101 .headers()
1102 .get("host")
1103 .and_then(|value| value.to_str().ok())
1104 .ok_or((StatusCode::BAD_REQUEST, "missing Host"))?;
1105 match host_rule {
1106 HostRule::Tcp {
1107 expected_host,
1108 forwarded_ports,
1109 } => {
1110 let accepted_forward = host
1111 .strip_prefix("127.0.0.1:")
1112 .and_then(|port| port.parse::<u16>().ok())
1113 .is_some_and(|port| forwarded_ports.contains(&port));
1114 if host != expected_host && !accepted_forward {
1115 return Err((StatusCode::BAD_REQUEST, "invalid loopback Host"));
1116 }
1117 }
1118 #[cfg(unix)]
1119 HostRule::Unix if host != "localhost" => {
1120 return Err((StatusCode::BAD_REQUEST, "invalid Unix Host"));
1121 }
1122 #[cfg(unix)]
1123 HostRule::Unix => {}
1124 }
1125 if request.headers().contains_key("origin") {
1126 return Err((StatusCode::BAD_REQUEST, "Origin rejected"));
1127 }
1128 let authorization = request
1129 .headers()
1130 .get("authorization")
1131 .and_then(|value| value.to_str().ok())
1132 .and_then(|value| value.strip_prefix("Bearer "))
1133 .ok_or((StatusCode::UNAUTHORIZED, "authentication failed"))?;
1134 let mut grant = credentials
1135 .authenticate(authorization)
1136 .map_err(|error| match error {
1137 AuthenticationError::Unauthenticated => {
1138 (StatusCode::UNAUTHORIZED, "authentication failed")
1139 }
1140 AuthenticationError::Unauthorized => (StatusCode::FORBIDDEN, "client id rejected"),
1141 AuthenticationError::Busy => (StatusCode::CONFLICT, "client channel busy"),
1142 })?;
1143 if let Some(requested) = request.headers().get("x-supercode-permissions") {
1144 let requested = requested
1145 .to_str()
1146 .ok()
1147 .and_then(|value| RuntimeAuthorization::parse_header(value).ok())
1148 .ok_or((StatusCode::BAD_REQUEST, "invalid permission narrowing"))?;
1149 grant.authorization = grant.authorization.restrict_to(&requested);
1150 }
1151 Ok(grant)
1152}
1153
1154fn contains_application_credential(value: &serde_json::Value) -> bool {
1155 match value {
1156 serde_json::Value::Object(object) => object.iter().any(|(key, child)| {
1157 matches!(
1158 key.to_ascii_lowercase().as_str(),
1159 "authorization" | "bearer" | "credential" | "token"
1160 ) || contains_application_credential(child)
1161 }),
1162 serde_json::Value::Array(array) => array.iter().any(contains_application_credential),
1163 _ => false,
1164 }
1165}
1166
1167fn websocket_config() -> WebSocketConfig {
1168 let mut config = WebSocketConfig::default();
1169 config.write_buffer_size = 128 * 1024;
1170 config.max_write_buffer_size = MAX_WRITE_BUFFER_BYTES;
1171 config.max_message_size = Some(MAX_MESSAGE_BYTES);
1172 config.max_frame_size = Some(MAX_FRAME_BYTES);
1173 config.accept_unmasked_frames = false;
1174 config
1175}
1176
1177fn http_error(status: StatusCode, message: &'static str) -> ErrorResponse {
1178 let (code, name) = match status {
1179 StatusCode::UNAUTHORIZED => (-32030, "unauthenticated"),
1180 StatusCode::FORBIDDEN => (-32031, "unauthorized"),
1181 StatusCode::CONFLICT => (-32000, "busy"),
1182 _ => (-32600, "invalid_request"),
1183 };
1184 tokio_tungstenite::tungstenite::http::Response::builder()
1185 .status(status)
1186 .header("content-type", "application/json")
1187 .body(Some(
1188 serde_json::json!({"error":{"code":code,"message":message,"data":{"name":name}}})
1189 .to_string(),
1190 ))
1191 .expect("static HTTP error response")
1192}
1193
1194fn credential_digest(secret: &[u8; TOKEN_BYTES]) -> [u8; 32] {
1195 *blake3::hash(secret).as_bytes()
1196}
1197
1198fn random_generation() -> Result<[u8; 16], CredentialError> {
1199 let mut generation = [0u8; 16];
1200 getrandom::getrandom(&mut generation).map_err(|_| CredentialError::RandomSource)?;
1201 Ok(generation)
1202}
1203
1204fn new_bootstrap(generation: [u8; 16]) -> Result<CodexBootstrapCredential, CredentialError> {
1205 let mut secret = [0u8; TOKEN_BYTES];
1206 getrandom::getrandom(&mut secret).map_err(|_| CredentialError::RandomSource)?;
1207 Ok(CodexBootstrapCredential {
1208 digest: credential_digest(&secret),
1209 secret,
1210 generation,
1211 revoked: false,
1212 })
1213}
1214
1215#[cfg(unix)]
1216fn write_bootstrap_secret(
1217 runtime_dir: PathBuf,
1218 secret: &[u8; TOKEN_BYTES],
1219 generation: [u8; 16],
1220) -> Result<CodexBootstrapFile, CredentialError> {
1221 use std::io::Write;
1222 use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt};
1223
1224 let metadata = std::fs::symlink_metadata(&runtime_dir)?;
1225 if metadata.file_type().is_symlink()
1226 || !metadata.is_dir()
1227 || metadata.permissions().mode() & 0o777 != 0o700
1228 || metadata.uid() != unsafe { libc::geteuid() }
1231 {
1232 return Err(CredentialError::Io(std::io::Error::new(
1233 std::io::ErrorKind::PermissionDenied,
1234 "runtime credential directory must be owner-owned mode 0700 and not a symlink",
1235 )));
1236 }
1237
1238 for _ in 0..3 {
1239 let mut name_bytes = [0u8; 16];
1240 getrandom::getrandom(&mut name_bytes).map_err(|_| CredentialError::RandomSource)?;
1241 let name = name_bytes
1242 .iter()
1243 .map(|byte| format!("{byte:02x}"))
1244 .collect::<String>();
1245 name_bytes.zeroize();
1246 let path = runtime_dir.join(format!("bootstrap-{name}.credential"));
1247 let mut options = std::fs::OpenOptions::new();
1248 options
1249 .write(true)
1250 .create_new(true)
1251 .mode(0o600)
1252 .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
1253 let mut file = match options.open(&path) {
1254 Ok(file) => file,
1255 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
1256 Err(error) => return Err(CredentialError::Io(error)),
1257 };
1258 let mut encoded = encode_hex(secret);
1259 let write_result = file
1260 .write_all(encoded.as_bytes())
1261 .and_then(|_| file.sync_all());
1262 encoded.zeroize();
1263 if let Err(error) = write_result {
1264 let _ = std::fs::remove_file(&path);
1265 return Err(CredentialError::Io(error));
1266 }
1267 let written = match file.metadata() {
1268 Ok(metadata) => metadata,
1269 Err(error) => {
1270 drop(file);
1271 let _ = std::fs::remove_file(&path);
1272 return Err(CredentialError::Io(error));
1273 }
1274 };
1275 if !written.is_file()
1276 || written.permissions().mode() & 0o777 != 0o600
1277 || written.uid() != unsafe { libc::geteuid() }
1280 {
1281 drop(file);
1282 let _ = std::fs::remove_file(&path);
1283 return Err(CredentialError::Io(std::io::Error::new(
1284 std::io::ErrorKind::PermissionDenied,
1285 "bootstrap credential file failed owner/mode verification",
1286 )));
1287 }
1288 return Ok(CodexBootstrapFile { path, generation });
1289 }
1290 Err(CredentialError::Collision)
1291}
1292
1293fn encode_hex(secret: &[u8; TOKEN_BYTES]) -> String {
1294 const HEX: &[u8; 16] = b"0123456789abcdef";
1295 let mut output = String::with_capacity(TOKEN_BYTES * 2);
1296 for byte in secret {
1297 output.push(HEX[(byte >> 4) as usize] as char);
1298 output.push(HEX[(byte & 0x0f) as usize] as char);
1299 }
1300 output
1301}
1302
1303fn decode_hex_secret(value: &str) -> Option<[u8; TOKEN_BYTES]> {
1304 if value.len() != TOKEN_BYTES * 2 {
1305 return None;
1306 }
1307 let mut output = [0u8; TOKEN_BYTES];
1308 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1309 output[index] = (hex_nibble(pair[0])? << 4) | hex_nibble(pair[1])?;
1310 }
1311 Some(output)
1312}
1313
1314fn hex_nibble(value: u8) -> Option<u8> {
1315 match value {
1316 b'0'..=b'9' => Some(value - b'0'),
1317 b'a'..=b'f' => Some(value - b'a' + 10),
1318 _ => None,
1319 }
1320}
1321
1322fn zero(secret: &mut [u8; TOKEN_BYTES]) {
1323 secret.zeroize();
1324}
1325
1326#[cfg(test)]
1327mod tests {
1328 use super::*;
1329 use async_trait::async_trait;
1330 use serde_json::{json, Value};
1331 use std::collections::VecDeque;
1332 use supercode::{
1333 ChatMessage, FrontendAttachSnapshot, FrontendAttachment, FrontendConnectionState,
1334 FrontendDisplayCapabilities, FrontendOperationInvocation, FrontendOperationResult,
1335 FrontendResponse, FrontendRuntimeDescriptor, FrontendTurnState, RuntimePermission,
1336 SdkError,
1337 };
1338 use tokio::sync::broadcast;
1339 use tokio_tungstenite::tungstenite::client::IntoClientRequest;
1340
1341 struct HandshakeRuntime;
1342
1343 #[async_trait]
1344 impl SdkRuntime for HandshakeRuntime {
1345 async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
1346 Err(SdkError::UnsupportedOperation("test.describe".into()))
1347 }
1348
1349 async fn attach(&self, _history_limit: usize) -> Result<FrontendAttachment, SdkError> {
1350 Err(SdkError::UnsupportedOperation("test.attach".into()))
1351 }
1352
1353 async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
1354 Err(SdkError::UnsupportedOperation("test.input".into()))
1355 }
1356
1357 async fn submit(&self, _prompt: String) -> Result<String, SdkError> {
1358 Err(SdkError::UnsupportedOperation("test.submit".into()))
1359 }
1360
1361 async fn interrupt(&self) -> Result<bool, SdkError> {
1362 Err(SdkError::UnsupportedOperation("test.interrupt".into()))
1363 }
1364
1365 async fn steer(&self, _prompt: String) -> Result<(), SdkError> {
1366 Err(SdkError::UnsupportedOperation("test.steer".into()))
1367 }
1368
1369 async fn respond(&self, _response: FrontendResponse) -> Result<(), SdkError> {
1370 Err(SdkError::UnsupportedOperation("test.respond".into()))
1371 }
1372
1373 async fn invoke(
1374 &self,
1375 _operation: FrontendOperationInvocation,
1376 ) -> Result<FrontendOperationResult, SdkError> {
1377 Err(SdkError::UnsupportedOperation("test.invoke".into()))
1378 }
1379 }
1380
1381 struct FunctionalRuntime {
1382 events: broadcast::Sender<supercode::FrontendEvent>,
1383 }
1384
1385 impl FunctionalRuntime {
1386 fn new() -> Arc<Self> {
1387 let (events, _) = broadcast::channel(32);
1388 Arc::new(Self { events })
1389 }
1390
1391 fn descriptor() -> FrontendRuntimeDescriptor {
1392 FrontendRuntimeDescriptor {
1393 schema_version: 2,
1394 session_id: "canonical-reconnect-session".into(),
1395 source_harness: Some("claude-code".into()),
1396 emulation_profile: Some("claude-code".into()),
1397 active_modules: Vec::new(),
1398 commands: Vec::new(),
1399 operations: Vec::new(),
1400 actions: supercode::FrontendActions {
1401 submit: true,
1402 interrupt: true,
1403 steer: true,
1404 respond: true,
1405 detach: true,
1406 close: true,
1407 },
1408 display: FrontendDisplayCapabilities {
1409 event_kinds: vec!["text_delta".into()],
1410 opaque_fallback: true,
1411 },
1412 model: "z-ai/glm-5.2".into(),
1413 turn_state: FrontendTurnState::Idle,
1414 connection_state: FrontendConnectionState::Connected,
1415 extensions: BTreeMap::new(),
1416 }
1417 }
1418 }
1419
1420 #[async_trait]
1421 impl SdkRuntime for FunctionalRuntime {
1422 async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
1423 Ok(Self::descriptor())
1424 }
1425
1426 async fn attach(&self, _history_limit: usize) -> Result<FrontendAttachment, SdkError> {
1427 Ok(FrontendAttachment::from_snapshot(
1428 FrontendAttachSnapshot {
1429 descriptor: Self::descriptor(),
1430 history: vec![ChatMessage::assistant("persisted history")],
1431 history_cursor: 1,
1432 replay: VecDeque::new(),
1433 },
1434 self.events.subscribe(),
1435 ))
1436 }
1437
1438 async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
1439 Ok(())
1440 }
1441
1442 async fn submit(&self, _prompt: String) -> Result<String, SdkError> {
1443 Ok(String::new())
1444 }
1445
1446 async fn interrupt(&self) -> Result<bool, SdkError> {
1447 Ok(true)
1448 }
1449
1450 async fn steer(&self, _prompt: String) -> Result<(), SdkError> {
1451 Ok(())
1452 }
1453
1454 async fn respond(&self, _response: FrontendResponse) -> Result<(), SdkError> {
1455 Ok(())
1456 }
1457
1458 async fn close(&self) -> Result<(), SdkError> {
1459 Ok(())
1460 }
1461 }
1462
1463 #[tokio::test]
1464 async fn credential_is_redacted_distinct_and_revocable() {
1465 let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
1466 let mut first = registry
1467 .issue("codex-a", RuntimeAuthorization::owner())
1468 .unwrap();
1469 let second = registry
1470 .issue("codex-b", RuntimeAuthorization::observer())
1471 .unwrap();
1472 assert_eq!(format!("{first:?}"), "CodexClientCredential([REDACTED])");
1473 assert_ne!(first.bearer(), second.bearer());
1474 let grant = registry
1475 .authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
1476 .unwrap();
1477 drop(grant);
1478 registry.revoke(&mut first).await.unwrap();
1479 assert!(registry
1480 .authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
1481 .is_err());
1482 }
1483
1484 #[test]
1485 fn credential_rng_failure_collision_retry_and_uniqueness_are_fail_closed() {
1486 let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
1487 assert!(matches!(
1488 registry.issue_with_fill("failed", RuntimeAuthorization::observer(), |_| Err(
1489 CredentialError::RandomSource
1490 )),
1491 Err(CredentialError::RandomSource)
1492 ));
1493 assert!(registry
1494 .records
1495 .lock()
1496 .unwrap_or_else(std::sync::PoisonError::into_inner)
1497 .is_empty());
1498
1499 let first = registry
1500 .issue_with_fill("first", RuntimeAuthorization::observer(), |secret| {
1501 secret.fill(0x11);
1502 Ok(())
1503 })
1504 .unwrap();
1505 let mut attempt = 0usize;
1506 let second = registry
1507 .issue_with_fill("second", RuntimeAuthorization::observer(), |secret| {
1508 attempt += 1;
1509 secret.fill(if attempt == 1 { 0x11 } else { 0x22 });
1510 Ok(())
1511 })
1512 .unwrap();
1513 assert_eq!(attempt, 2, "collision must consume a fresh RNG fill");
1514 assert_ne!(first.bearer(), second.bearer());
1515
1516 let mut unique = BTreeSet::new();
1517 for index in 0..4096u64 {
1518 let credential = registry
1519 .issue(format!("client-{index}"), RuntimeAuthorization::observer())
1520 .unwrap();
1521 assert!(unique.insert(credential.digest));
1522 }
1523 }
1524
1525 #[test]
1526 fn application_message_credentials_are_rejected_by_field_not_plain_text() {
1527 assert!(contains_application_credential(
1528 &serde_json::json!({"id":1,"method":"turn/start","params":{"token":"secret"}})
1529 ));
1530 assert!(contains_application_credential(
1531 &serde_json::json!({"authorization":"Bearer secret"})
1532 ));
1533 assert!(!contains_application_credential(
1534 &serde_json::json!({"input":[{"type":"text","text":"discuss token budgets"}]})
1535 ));
1536 }
1537
1538 #[cfg(unix)]
1539 #[tokio::test]
1540 async fn child_only_credential_is_absent_from_argv_and_reused_command_state() {
1541 use std::process::Stdio;
1542
1543 let registry = CodexCredentialRegistry::new("runtime".into(), [8; 16]);
1544 let credential = registry
1545 .issue("spawn-test", RuntimeAuthorization::observer())
1546 .unwrap();
1547 let expected = credential.bearer_value();
1548 let variable = "SUPERCODE_CODEX_TEST_CHILD_ONLY";
1549 let mut command = tokio::process::Command::new("sh");
1550 command
1551 .args(["-c", "printf %s \"$SUPERCODE_CODEX_TEST_CHILD_ONLY\""])
1552 .stdout(Stdio::piped());
1553
1554 let first = credential
1555 .spawn_tokio_child(&mut command, variable)
1556 .unwrap()
1557 .wait_with_output()
1558 .await
1559 .unwrap();
1560 assert!(first.status.success());
1561 assert!(
1562 first.stdout == expected.as_bytes(),
1563 "the child did not receive its exact scoped credential"
1564 );
1565 assert!(
1566 !format!("{command:?}").contains(&expected),
1567 "the reusable command debug state retained the credential"
1568 );
1569
1570 let second = command.spawn().unwrap().wait_with_output().await.unwrap();
1571 assert!(second.status.success());
1572 assert!(
1573 second.stdout.is_empty(),
1574 "the reusable command leaked the credential into a later child"
1575 );
1576 }
1577
1578 #[cfg(unix)]
1579 #[tokio::test]
1580 async fn child_only_credential_is_absent_from_parent_process_table_and_diagnostics() {
1581 use std::process::Stdio;
1582
1583 let registry = CodexCredentialRegistry::new("runtime".into(), [11; 16]);
1584 let credential = registry
1585 .issue("process-scan", RuntimeAuthorization::observer())
1586 .unwrap();
1587 let expected = credential.bearer_value();
1588 let variable = "SUPERCODE_CODEX_TEST_PROCESS_SCAN";
1589 assert!(std::env::var_os(variable).is_none());
1590 let root = std::env::temp_dir().join(format!(
1591 "supercode-codex-process-scan-{}-{}",
1592 std::process::id(),
1593 epoch_millis_for_test()
1594 ));
1595 std::fs::create_dir(&root).unwrap();
1596 let child_only_output = root.join("child-only");
1597 let mut command = tokio::process::Command::new("sh");
1598 command
1599 .args([
1600 "-c",
1601 "printf %s \"$SUPERCODE_CODEX_TEST_PROCESS_SCAN\" > \"$1\"; printf crash-diagnostic >&2; sleep 1",
1602 "sh",
1603 ])
1604 .arg(&child_only_output)
1605 .stdout(Stdio::null())
1606 .stderr(Stdio::piped());
1607 let child = credential
1608 .spawn_tokio_child(&mut command, variable)
1609 .unwrap();
1610 let pid = child.id().unwrap();
1611
1612 assert!(std::env::var_os(variable).is_none());
1613 assert!(!format!("{command:?}").contains(&expected));
1614 let process_table = tokio::process::Command::new("ps")
1615 .args(["-o", "command=", "-p", &pid.to_string()])
1616 .output()
1617 .await
1618 .unwrap();
1619 assert!(process_table.status.success());
1620 assert!(!String::from_utf8_lossy(&process_table.stdout).contains(&expected));
1621
1622 let output = child.wait_with_output().await.unwrap();
1623 assert!(output.status.success());
1624 assert_eq!(
1625 std::fs::read(&child_only_output).unwrap(),
1626 expected.as_bytes()
1627 );
1628 assert_eq!(String::from_utf8_lossy(&output.stderr), "crash-diagnostic");
1629 assert!(!output
1630 .stderr
1631 .windows(expected.len())
1632 .any(|bytes| bytes == expected.as_bytes()));
1633 std::fs::remove_dir_all(root).unwrap();
1634 }
1635
1636 #[cfg(windows)]
1637 #[tokio::test]
1638 async fn windows_child_credential_is_absent_from_process_table_and_diagnostics() {
1639 use std::process::Stdio;
1640
1641 let registry = CodexCredentialRegistry::new("runtime".into(), [12; 16]);
1642 let credential = registry
1643 .issue("windows-process-scan", RuntimeAuthorization::observer())
1644 .unwrap();
1645 let expected = credential.bearer_value();
1646 let variable = "SUPERCODE_CODEX_TEST_WINDOWS_PROCESS_SCAN";
1647 assert!(std::env::var_os(variable).is_none());
1648 let root = std::env::temp_dir().join(format!(
1649 "supercode-codex-windows-process-scan-{}-{}",
1650 std::process::id(),
1651 epoch_millis_for_test()
1652 ));
1653 std::fs::create_dir(&root).unwrap();
1654 let child_only_output = root.join("child-only");
1655 let output_literal = child_only_output.display().to_string().replace('\'', "''");
1656 let script = format!(
1657 "[IO.File]::WriteAllText('{output_literal}', $env:{variable}); \
1658 [Console]::Error.Write('crash-diagnostic'); Start-Sleep -Seconds 5"
1659 );
1660 let mut command = tokio::process::Command::new("powershell.exe");
1661 command
1662 .args(["-NoProfile", "-NonInteractive", "-Command"])
1663 .arg(script)
1664 .stdout(Stdio::null())
1665 .stderr(Stdio::piped());
1666 let child = credential
1667 .spawn_tokio_child(&mut command, variable)
1668 .unwrap();
1669 let pid = child.id().unwrap();
1670
1671 assert!(std::env::var_os(variable).is_none());
1672 assert!(!format!("{command:?}").contains(&expected));
1673 let process_table = tokio::process::Command::new("powershell.exe")
1674 .args([
1675 "-NoProfile",
1676 "-NonInteractive",
1677 "-Command",
1678 &format!("(Get-CimInstance Win32_Process -Filter 'ProcessId = {pid}').CommandLine"),
1679 ])
1680 .output()
1681 .await
1682 .unwrap();
1683 assert!(process_table.status.success());
1684 let process_command = String::from_utf8_lossy(&process_table.stdout);
1685 assert!(
1686 process_command.contains(variable),
1687 "CIM scan did not capture the live child command line: {process_command:?}"
1688 );
1689 assert!(!process_command.contains(&expected));
1690
1691 let output = child.wait_with_output().await.unwrap();
1692 assert!(output.status.success());
1693 assert_eq!(
1694 std::fs::read(&child_only_output).unwrap(),
1695 expected.as_bytes()
1696 );
1697 assert_eq!(
1698 String::from_utf8_lossy(&output.stderr).trim(),
1699 "crash-diagnostic"
1700 );
1701 assert!(!output
1702 .stderr
1703 .windows(expected.len())
1704 .any(|bytes| bytes == expected.as_bytes()));
1705 std::fs::remove_dir_all(root).unwrap();
1706 }
1707
1708 #[cfg(unix)]
1709 #[tokio::test]
1710 async fn bootstrap_file_is_private_required_and_deleted() {
1711 use std::os::unix::fs::PermissionsExt;
1712
1713 let root = std::env::temp_dir().join(format!(
1714 "supercode-bootstrap-{}-{}",
1715 std::process::id(),
1716 epoch_millis_for_test()
1717 ));
1718 std::fs::create_dir(&root).unwrap();
1719 std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700)).unwrap();
1720 let (endpoint, mut bootstrap) =
1721 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
1722 let mut bootstrap_file = endpoint.write_bootstrap_file(&bootstrap, &root).unwrap();
1723 let path = bootstrap_file.path.clone();
1724 assert_eq!(
1725 std::fs::metadata(&path).unwrap().permissions().mode() & 0o777,
1726 0o600
1727 );
1728 assert_eq!(
1729 format!("{bootstrap:?}"),
1730 "CodexBootstrapCredential([REDACTED])"
1731 );
1732
1733 let mut client = endpoint
1734 .issue_interactive_from_file(&bootstrap_file, "stock")
1735 .unwrap();
1736 let (_foreign, foreign_bootstrap) =
1737 CodexEndpoint::new(FunctionalRuntime::new(), "other", "/workspace").unwrap();
1738 assert!(matches!(
1739 endpoint.issue_interactive(&foreign_bootstrap, "forbidden"),
1740 Err(CredentialError::InvalidBootstrap)
1741 ));
1742
1743 endpoint
1744 .revoke_frontend_from_file(&bootstrap_file, &mut client)
1745 .await
1746 .unwrap();
1747 let (replacement, replacement_file) = endpoint
1748 .rotate_bootstrap_file(&mut bootstrap, &mut bootstrap_file, &root)
1749 .unwrap();
1750 assert!(!path.exists(), "rotation must delete the prior file");
1751 assert!(matches!(
1752 endpoint.issue_observer(&bootstrap, "old-bootstrap"),
1753 Err(CredentialError::InvalidBootstrap)
1754 ));
1755 assert_eq!(
1756 std::fs::metadata(&replacement_file.path)
1757 .unwrap()
1758 .permissions()
1759 .mode()
1760 & 0o777,
1761 0o600
1762 );
1763 let replacement_client = endpoint
1764 .issue_interactive_from_file(&replacement_file, "replacement")
1765 .unwrap();
1766 drop(replacement_client);
1767 drop(replacement_file);
1768 drop(replacement);
1769 drop(bootstrap_file);
1770 assert!(!path.exists());
1771 std::fs::remove_dir(root).unwrap();
1772 }
1773
1774 #[cfg(unix)]
1775 #[test]
1776 fn bootstrap_file_rotation_rolls_back_when_old_file_cannot_be_removed() {
1777 use std::os::unix::fs::PermissionsExt;
1778
1779 let root = std::env::temp_dir().join(format!(
1780 "supercode-bootstrap-rollback-{}-{}",
1781 std::process::id(),
1782 epoch_millis_for_test()
1783 ));
1784 std::fs::create_dir(&root).unwrap();
1785 std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700)).unwrap();
1786 let (endpoint, mut bootstrap) =
1787 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
1788 let mut bootstrap_file = endpoint.write_bootstrap_file(&bootstrap, &root).unwrap();
1789 let original_path = bootstrap_file.path.clone();
1790
1791 let error = match endpoint.rotate_bootstrap_file_with_remove(
1792 &mut bootstrap,
1793 &mut bootstrap_file,
1794 root.clone(),
1795 |_| {
1796 Err(std::io::Error::new(
1797 std::io::ErrorKind::PermissionDenied,
1798 "injected deletion failure",
1799 ))
1800 },
1801 ) {
1802 Ok(_) => panic!("injected removal failure must abort rotation"),
1803 Err(error) => error,
1804 };
1805 assert!(matches!(error, CredentialError::Io(_)));
1806 assert!(original_path.exists());
1807 assert!(endpoint
1808 .issue_observer(&bootstrap, "still-authorized")
1809 .is_ok());
1810 let files = std::fs::read_dir(&root)
1811 .unwrap()
1812 .filter_map(Result::ok)
1813 .map(|entry| entry.path())
1814 .collect::<Vec<_>>();
1815 assert_eq!(files, vec![original_path.clone()]);
1816
1817 drop(bootstrap_file);
1818 assert!(!original_path.exists());
1819 std::fs::remove_dir(root).unwrap();
1820 }
1821
1822 #[test]
1823 fn in_memory_bootstrap_rotation_invalidates_old_value_before_return() {
1824 let (endpoint, mut bootstrap) =
1825 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
1826 let replacement = endpoint.rotate_bootstrap(&mut bootstrap).unwrap();
1827 assert!(matches!(
1828 endpoint.issue_observer(&bootstrap, "old"),
1829 Err(CredentialError::InvalidBootstrap)
1830 ));
1831 assert!(endpoint.issue_observer(&replacement, "new").is_ok());
1832 }
1833
1834 #[test]
1835 fn upgrade_validation_is_fail_closed_and_permission_narrowing_only() {
1836 let registry = CodexCredentialRegistry::new("runtime".into(), [7; 16]);
1837 let credential = registry
1838 .issue("codex", RuntimeAuthorization::owner())
1839 .unwrap();
1840 let request = |uri: &str, authorization: Option<&str>| {
1841 let mut builder = Request::builder().uri(uri).header("host", "127.0.0.1:1234");
1842 if let Some(authorization) = authorization {
1843 builder = builder.header("authorization", authorization);
1844 }
1845 builder.body(()).unwrap()
1846 };
1847 let valid = || request("/", Some(&credential.bearer()));
1848 let rule = || HostRule::Tcp {
1849 expected_host: "127.0.0.1:1234".into(),
1850 forwarded_ports: BTreeSet::from([4321]),
1851 };
1852 let grant = validate_upgrade(&valid(), &rule(), ®istry).unwrap();
1853 assert!(grant.authorization.allows(RuntimePermission::Terminate));
1854 drop(grant);
1855
1856 let unauthenticated = [
1857 request("/", None),
1858 request("/", Some("Basic nope")),
1859 request("/", Some("Bearer malformed")),
1860 request("/", Some("Bearer codex.00")),
1861 ];
1862 let mut external = Vec::new();
1863 for request in unauthenticated {
1864 let (status, message) = match validate_upgrade(&request, &rule(), ®istry) {
1865 Ok(_) => panic!("invalid credential must be rejected"),
1866 Err(error) => error,
1867 };
1868 assert_eq!(status, StatusCode::UNAUTHORIZED);
1869 let response = http_error(status, message);
1870 assert!(response
1871 .headers()
1872 .get("access-control-allow-origin")
1873 .is_none());
1874 external.push(response.body().clone());
1875 }
1876 assert!(external.windows(2).all(|pair| pair[0] == pair[1]));
1877
1878 let actual = credential.bearer();
1879 let wrong_client = actual.replacen("Bearer codex.", "Bearer impostor.", 1);
1880 let (status, _) =
1881 match validate_upgrade(&request("/", Some(&wrong_client)), &rule(), ®istry) {
1882 Ok(_) => panic!("mismatched client id must be rejected"),
1883 Err(error) => error,
1884 };
1885 assert_eq!(status, StatusCode::FORBIDDEN);
1886
1887 let mut narrowed = valid();
1888 narrowed.headers_mut().insert(
1889 "x-supercode-permissions",
1890 "observe,terminate".parse().unwrap(),
1891 );
1892 let grant = validate_upgrade(&narrowed, &rule(), ®istry).unwrap();
1893 assert!(grant.authorization.allows(RuntimePermission::Observe));
1894 assert!(grant.authorization.allows(RuntimePermission::Terminate));
1895 assert!(!grant.authorization.allows(RuntimePermission::Interact));
1896 drop(grant);
1897
1898 let interactive = registry
1899 .issue(
1900 "interactive",
1901 RuntimeAuthorization::new([
1902 RuntimePermission::Observe,
1903 RuntimePermission::Interact,
1904 ]),
1905 )
1906 .unwrap();
1907 let mut cannot_expand = request("/", Some(&interactive.bearer()));
1908 cannot_expand.headers_mut().insert(
1909 "x-supercode-permissions",
1910 "observe,terminate".parse().unwrap(),
1911 );
1912 let grant = validate_upgrade(&cannot_expand, &rule(), ®istry).unwrap();
1913 assert!(grant.authorization.allows(RuntimePermission::Observe));
1914 assert!(!grant.authorization.allows(RuntimePermission::Terminate));
1915 drop(grant);
1916
1917 let mut forwarded = valid();
1918 forwarded
1919 .headers_mut()
1920 .insert("host", "127.0.0.1:4321".parse().unwrap());
1921 drop(validate_upgrade(&forwarded, &rule(), ®istry).unwrap());
1922
1923 let mut rejected = Vec::new();
1924 rejected.push(request("/?token=nope", Some(&credential.bearer())));
1925 let mut cookie = valid();
1926 cookie
1927 .headers_mut()
1928 .insert("cookie", "x=y".parse().unwrap());
1929 rejected.push(cookie);
1930 let mut origin = valid();
1931 origin
1932 .headers_mut()
1933 .insert("origin", "https://evil.invalid".parse().unwrap());
1934 rejected.push(origin);
1935 let mut foreign_host = valid();
1936 foreign_host
1937 .headers_mut()
1938 .insert("host", "localhost:1234".parse().unwrap());
1939 rejected.push(foreign_host);
1940 for carrier in [
1941 "content-length",
1942 "transfer-encoding",
1943 "sec-websocket-protocol",
1944 ] {
1945 let mut forbidden = valid();
1946 forbidden
1947 .headers_mut()
1948 .insert(carrier, "1".parse().unwrap());
1949 rejected.push(forbidden);
1950 }
1951 for duplicate in [
1952 "host",
1953 "origin",
1954 "authorization",
1955 "content-length",
1956 "x-supercode-permissions",
1957 ] {
1958 let mut request = valid();
1959 request
1960 .headers_mut()
1961 .append(duplicate, "duplicate".parse().unwrap());
1962 rejected.push(request);
1963 }
1964 for request in rejected {
1965 assert!(validate_upgrade(&request, &rule(), ®istry).is_err());
1966 }
1967
1968 let absolute = Request::builder()
1969 .uri("ws://127.0.0.1:1234/")
1970 .header("host", "127.0.0.1:1234")
1971 .header("authorization", credential.bearer())
1972 .body(())
1973 .unwrap();
1974 assert!(validate_upgrade(&absolute, &rule(), ®istry).is_err());
1975 }
1976
1977 #[tokio::test]
1978 async fn half_open_handshakes_are_deadlined_and_connection_budget_recovers() {
1979 use tokio::io::{AsyncReadExt, AsyncWriteExt};
1980 use tokio_tungstenite::tungstenite::client::IntoClientRequest;
1981
1982 let (endpoint, bootstrap) =
1983 CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
1984 let credential = endpoint
1985 .issue_interactive(&bootstrap, "bounded-client")
1986 .unwrap();
1987 let handle = endpoint.bind_websocket(0).await.unwrap();
1988 let address = match handle.kind() {
1989 CodexEndpointKind::WebSocket(address) => *address,
1990 #[cfg(unix)]
1991 CodexEndpointKind::Unix(_) => unreachable!(),
1992 };
1993
1994 let mut oversized = TcpStream::connect(address).await.unwrap();
1995 let oversized_request = format!(
1996 "GET / HTTP/1.1\r\nHost: {address}\r\nX-Oversized: {}\r\n\r\n",
1997 "x".repeat(MAX_HANDSHAKE_BYTES)
1998 );
1999 oversized
2000 .write_all(oversized_request.as_bytes())
2001 .await
2002 .unwrap();
2003 let mut byte = [0u8; 1];
2004 let oversized_closed = tokio::time::timeout(
2005 std::time::Duration::from_millis(100),
2006 oversized.read(&mut byte),
2007 )
2008 .await
2009 .expect("an oversized handshake must be rejected before its deadline");
2010 assert!(matches!(oversized_closed, Ok(0) | Err(_)));
2011
2012 let mut held = Vec::with_capacity(MAX_CONCURRENT_CONNECTIONS);
2013 for _ in 0..MAX_CONCURRENT_CONNECTIONS {
2014 held.push(TcpStream::connect(address).await.unwrap());
2015 }
2016 let mut overflow = TcpStream::connect(address).await.unwrap();
2017 let overflow_closed = tokio::time::timeout(
2018 std::time::Duration::from_millis(100),
2019 overflow.read(&mut byte),
2020 )
2021 .await
2022 .expect("an over-budget connection must be rejected promptly");
2023 assert!(matches!(overflow_closed, Ok(0) | Err(_)));
2024
2025 tokio::time::sleep(HANDSHAKE_TIMEOUT + std::time::Duration::from_millis(100)).await;
2026 for mut stream in held {
2027 let closed = tokio::time::timeout(
2028 std::time::Duration::from_millis(100),
2029 stream.read(&mut byte),
2030 )
2031 .await
2032 .expect("a half-open handshake must be closed at its deadline");
2033 assert!(matches!(closed, Ok(0) | Err(_)));
2034 }
2035
2036 let mut request = format!("ws://{address}/").into_client_request().unwrap();
2037 request
2038 .headers_mut()
2039 .insert("authorization", credential.bearer().parse().unwrap());
2040 let (mut websocket, _) = tokio_tungstenite::connect_async(request).await.unwrap();
2041 websocket.close(None).await.unwrap();
2042 handle.shutdown().await;
2043 }
2044
2045 #[tokio::test]
2046 async fn authenticated_loopback_websocket_is_ready_and_revocation_closes_it() {
2047 let (endpoint, bootstrap) =
2048 CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
2049 let mut credential = endpoint
2050 .issue_interactive(&bootstrap, "stock-codex")
2051 .unwrap();
2052 let handle = endpoint.bind_websocket(0).await.unwrap();
2053 assert_eq!(handle.health(), CodexEndpointHealth::Ready);
2054 let address = match handle.kind() {
2055 CodexEndpointKind::WebSocket(address) => *address,
2056 #[cfg(unix)]
2057 CodexEndpointKind::Unix(_) => unreachable!(),
2058 };
2059 let mut request = format!("ws://{address}/").into_client_request().unwrap();
2060 request
2061 .headers_mut()
2062 .insert("authorization", credential.bearer().parse().unwrap());
2063 let (mut websocket, _) = tokio_tungstenite::connect_async(request).await.unwrap();
2064 let mut competing = format!("ws://{address}/").into_client_request().unwrap();
2065 competing
2066 .headers_mut()
2067 .insert("authorization", credential.bearer().parse().unwrap());
2068 let competing_error = tokio_tungstenite::connect_async(competing)
2069 .await
2070 .expect_err("the same credential cannot own two live channels");
2071 match competing_error {
2072 tokio_tungstenite::tungstenite::Error::Http(response) => {
2073 assert_eq!(response.status(), StatusCode::CONFLICT);
2074 }
2075 other => panic!("expected HTTP busy rejection, got {other}"),
2076 }
2077 websocket
2078 .send(Message::Text("not-json".into()))
2079 .await
2080 .unwrap();
2081 let parse_error = websocket
2082 .next()
2083 .await
2084 .unwrap()
2085 .unwrap()
2086 .into_text()
2087 .unwrap();
2088 let parse_error: serde_json::Value = serde_json::from_str(&parse_error).unwrap();
2089 assert_eq!(parse_error["error"]["code"], -32700);
2090 assert_eq!(parse_error["error"]["data"]["name"], "parse_error");
2091 websocket
2092 .send(Message::Text(
2093 serde_json::json!({"id":"initialize","method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
2094 .to_string()
2095 .into(),
2096 ))
2097 .await
2098 .unwrap();
2099 let response = websocket
2100 .next()
2101 .await
2102 .unwrap()
2103 .unwrap()
2104 .into_text()
2105 .unwrap();
2106 assert_eq!(
2107 serde_json::from_str::<serde_json::Value>(&response).unwrap()["id"],
2108 "initialize"
2109 );
2110 let notification = websocket
2111 .next()
2112 .await
2113 .unwrap()
2114 .unwrap()
2115 .into_text()
2116 .unwrap();
2117 assert_eq!(
2118 serde_json::from_str::<serde_json::Value>(¬ification).unwrap()["method"],
2119 "remoteControl/status/changed"
2120 );
2121
2122 endpoint
2123 .revoke_frontend(&bootstrap, &mut credential)
2124 .await
2125 .unwrap();
2126 tokio::time::timeout(std::time::Duration::from_secs(1), async {
2127 loop {
2128 match websocket.next().await {
2129 None | Some(Ok(Message::Close(_))) | Some(Err(_)) => break,
2130 _ => {}
2131 }
2132 }
2133 })
2134 .await
2135 .expect("revocation must close the active channel promptly");
2136 handle.shutdown().await;
2137 }
2138
2139 #[tokio::test]
2140 async fn websocket_reconnect_preserves_deterministic_thread_identity() {
2141 let (endpoint, bootstrap) =
2142 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
2143 let credential = endpoint
2144 .issue_interactive(&bootstrap, "stock-codex")
2145 .unwrap();
2146 let handle = endpoint.bind_websocket(0).await.unwrap();
2147 let address = match handle.kind() {
2148 CodexEndpointKind::WebSocket(address) => *address,
2149 #[cfg(unix)]
2150 CodexEndpointKind::Unix(_) => unreachable!(),
2151 };
2152
2153 let connect = || {
2154 let bearer = credential.bearer();
2155 async move {
2156 let mut request = format!("ws://{address}/").into_client_request().unwrap();
2157 request
2158 .headers_mut()
2159 .insert("authorization", bearer.parse().unwrap());
2160 tokio_tungstenite::connect_async(request).await.unwrap().0
2161 }
2162 };
2163 let mut first = connect().await;
2164 first
2165 .send(Message::Text(
2166 json!({"id":1,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
2167 .to_string()
2168 .into(),
2169 ))
2170 .await
2171 .unwrap();
2172 let _initialize = first.next().await.unwrap().unwrap();
2173 let _status = first.next().await.unwrap().unwrap();
2174 first
2175 .send(Message::Text(
2176 json!({"method":"initialized"}).to_string().into(),
2177 ))
2178 .await
2179 .unwrap();
2180 first
2181 .send(Message::Text(
2182 json!({"id":2,"method":"thread/start","params":{}})
2183 .to_string()
2184 .into(),
2185 ))
2186 .await
2187 .unwrap();
2188 let started: Value =
2189 serde_json::from_str(&first.next().await.unwrap().unwrap().into_text().unwrap())
2190 .unwrap();
2191 let thread_id = started["result"]["thread"]["id"]
2192 .as_str()
2193 .unwrap()
2194 .to_owned();
2195 let _thread_started = first.next().await.unwrap().unwrap();
2196 first.close(None).await.unwrap();
2197
2198 let mut inactive = {
2199 let records = endpoint
2200 .credentials
2201 .records
2202 .lock()
2203 .unwrap_or_else(std::sync::PoisonError::into_inner);
2204 records
2205 .get(&credential.digest)
2206 .unwrap()
2207 .active_channel
2208 .subscribe()
2209 };
2210 tokio::time::timeout(std::time::Duration::from_secs(1), async {
2211 while *inactive.borrow() {
2212 inactive.changed().await.unwrap();
2213 }
2214 })
2215 .await
2216 .expect("closed socket must release its channel");
2217
2218 let mut second = connect().await;
2219 second
2220 .send(Message::Text(
2221 json!({"id":3,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
2222 .to_string()
2223 .into(),
2224 ))
2225 .await
2226 .unwrap();
2227 let _initialize = second.next().await.unwrap().unwrap();
2228 let _status = second.next().await.unwrap().unwrap();
2229 second
2230 .send(Message::Text(
2231 json!({"method":"initialized"}).to_string().into(),
2232 ))
2233 .await
2234 .unwrap();
2235 second
2236 .send(Message::Text(
2237 json!({"id":4,"method":"thread/read","params":{"threadId":thread_id}})
2238 .to_string()
2239 .into(),
2240 ))
2241 .await
2242 .unwrap();
2243 let resumed: Value =
2244 serde_json::from_str(&second.next().await.unwrap().unwrap().into_text().unwrap())
2245 .unwrap();
2246 assert_eq!(resumed["result"]["thread"]["id"], thread_id);
2247 second.close(None).await.unwrap();
2248 handle.shutdown().await;
2249 }
2250
2251 #[tokio::test]
2252 async fn observer_and_competing_controller_fail_with_stable_codes() {
2253 async fn initialize_and_attach(
2254 adapter: Arc<CodexAppServerAdapter>,
2255 ) -> (crate::codex_app_server_v0_144::CodexConnection, String) {
2256 let mut connection = adapter.connection();
2257 connection
2258 .handle(json!({"id":1,"method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}}))
2259 .await;
2260 connection.handle(json!({"method":"initialized"})).await;
2261 let started = connection
2262 .handle(json!({"id":2,"method":"thread/start","params":{}}))
2263 .await;
2264 let thread_id = started[0]["result"]["thread"]["id"]
2265 .as_str()
2266 .unwrap()
2267 .to_owned();
2268 (connection, thread_id)
2269 }
2270
2271 let (endpoint, bootstrap) =
2272 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
2273 let observer = endpoint.issue_observer(&bootstrap, "observer").unwrap();
2274 let first = endpoint.issue_interactive(&bootstrap, "first").unwrap();
2275 let second = endpoint.issue_interactive(&bootstrap, "second").unwrap();
2276 let observer_grant = endpoint
2277 .credentials
2278 .authenticate(observer.bearer().strip_prefix("Bearer ").unwrap())
2279 .unwrap();
2280 let first_grant = endpoint
2281 .credentials
2282 .authenticate(first.bearer().strip_prefix("Bearer ").unwrap())
2283 .unwrap();
2284 let second_grant = endpoint
2285 .credentials
2286 .authenticate(second.bearer().strip_prefix("Bearer ").unwrap())
2287 .unwrap();
2288 let (mut observer_connection, observer_thread) =
2289 initialize_and_attach(endpoint.authenticated_adapter(&observer_grant)).await;
2290 let denied = observer_connection
2291 .handle(json!({"id":3,"method":"turn/start","params":{"threadId":observer_thread,"input":[{"type":"text","text":"forbidden"}]}}))
2292 .await;
2293 assert_eq!(denied[0]["error"]["code"], -32031);
2294 assert_eq!(denied[0]["error"]["data"]["name"], "unauthorized");
2295
2296 let (mut first_connection, first_thread) =
2297 initialize_and_attach(endpoint.authenticated_adapter(&first_grant)).await;
2298 let accepted = first_connection
2299 .handle(json!({"id":4,"method":"turn/start","params":{"threadId":first_thread,"input":[{"type":"text","text":"claim"}]}}))
2300 .await;
2301 assert!(accepted[0].get("result").is_some());
2302
2303 let (mut second_connection, second_thread) =
2304 initialize_and_attach(endpoint.authenticated_adapter(&second_grant)).await;
2305 let conflict = second_connection
2306 .handle(json!({"id":5,"method":"turn/start","params":{"threadId":second_thread,"input":[{"type":"text","text":"compete"}]}}))
2307 .await;
2308 assert_eq!(conflict[0]["error"]["code"], -32032);
2309 assert_eq!(conflict[0]["error"]["data"]["name"], "controller_required");
2310 }
2311
2312 #[cfg(unix)]
2313 #[tokio::test]
2314 async fn authenticated_unix_websocket_uses_protected_files() {
2315 use std::os::unix::fs::PermissionsExt;
2316
2317 let (endpoint, bootstrap) =
2318 CodexEndpoint::new(Arc::new(HandshakeRuntime), "runtime", "/workspace").unwrap();
2319 let credential = endpoint
2320 .issue_interactive(&bootstrap, "stock-codex-unix")
2321 .unwrap();
2322 let root = std::env::temp_dir().join(format!(
2323 "supercode-codex-unix-{}-{}",
2324 std::process::id(),
2325 epoch_millis_for_test()
2326 ));
2327 let socket = root.join("adapter.sock");
2328 let handle = endpoint.bind_unix(&socket).await.unwrap();
2329 assert_eq!(
2330 std::fs::metadata(&root).unwrap().permissions().mode() & 0o777,
2331 0o700
2332 );
2333 assert_eq!(
2334 std::fs::metadata(&socket).unwrap().permissions().mode() & 0o777,
2335 0o600
2336 );
2337 let stream = UnixStream::connect(&socket).await.unwrap();
2338 let mut request = "ws://localhost/".into_client_request().unwrap();
2339 request
2340 .headers_mut()
2341 .insert("authorization", credential.bearer().parse().unwrap());
2342 let (mut websocket, _) = tokio_tungstenite::client_async(request, stream)
2343 .await
2344 .unwrap();
2345 websocket
2346 .send(Message::Text(
2347 serde_json::json!({"id":"initialize","method":"initialize","params":{"clientInfo":{"version":"0.144.4"}}})
2348 .to_string()
2349 .into(),
2350 ))
2351 .await
2352 .unwrap();
2353 assert!(websocket.next().await.unwrap().unwrap().is_text());
2354 handle.shutdown().await;
2355 assert!(!socket.exists());
2356 std::fs::remove_dir(&root).unwrap();
2357 }
2358
2359 #[cfg(unix)]
2360 #[tokio::test]
2361 async fn unix_binding_refuses_precreated_destinations_and_symlink_parents() {
2362 use std::os::unix::fs::{symlink, PermissionsExt};
2363
2364 let (endpoint, _bootstrap) =
2365 CodexEndpoint::new(FunctionalRuntime::new(), "runtime", "/workspace").unwrap();
2366 let base = std::env::temp_dir().join(format!(
2367 "supercode-codex-precreate-{}-{}",
2368 std::process::id(),
2369 epoch_millis_for_test()
2370 ));
2371 std::fs::create_dir(&base).unwrap();
2372 std::fs::set_permissions(&base, std::fs::Permissions::from_mode(0o700)).unwrap();
2373 let occupied = base.join("occupied.sock");
2374 std::fs::write(&occupied, b"do not replace").unwrap();
2375 let error = match endpoint.bind_unix(&occupied).await {
2376 Ok(_) => panic!("precreated destination must be rejected"),
2377 Err(error) => error,
2378 };
2379 assert_eq!(error.kind(), std::io::ErrorKind::AlreadyExists);
2380 assert_eq!(std::fs::read(&occupied).unwrap(), b"do not replace");
2381
2382 let real_parent = base.join("real");
2383 std::fs::create_dir(&real_parent).unwrap();
2384 std::fs::set_permissions(&real_parent, std::fs::Permissions::from_mode(0o700)).unwrap();
2385 let linked_parent = base.join("linked");
2386 symlink(&real_parent, &linked_parent).unwrap();
2387 let error = match endpoint
2388 .bind_unix(linked_parent.join("endpoint.sock"))
2389 .await
2390 {
2391 Ok(_) => panic!("symlink parent must be rejected"),
2392 Err(error) => error,
2393 };
2394 assert_eq!(error.kind(), std::io::ErrorKind::PermissionDenied);
2395 assert!(!real_parent.join("endpoint.sock").exists());
2396
2397 std::fs::remove_file(linked_parent).unwrap();
2398 std::fs::remove_dir(real_parent).unwrap();
2399 std::fs::remove_file(occupied).unwrap();
2400 std::fs::remove_dir(base).unwrap();
2401 }
2402
2403 fn epoch_millis_for_test() -> u128 {
2404 std::time::SystemTime::now()
2405 .duration_since(std::time::UNIX_EPOCH)
2406 .unwrap()
2407 .as_millis()
2408 }
2409}