diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 339fc68867..b3096dd9b0 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -4006,6 +4006,7 @@ name = "codex-terminal-browser" version = "0.0.0" dependencies = [ "anyhow", + "codex-network-proxy", "codex-protocol", "codex-sandboxing", "codex-utils-absolute-path", @@ -4014,7 +4015,6 @@ dependencies = [ "futures", "libc", "pretty_assertions", - "reqwest 0.12.28", "serde", "serde_json", "sha2 0.10.9", diff --git a/codex-rs/terminal-browser/Cargo.toml b/codex-rs/terminal-browser/Cargo.toml index 88a8c1a1d2..ace066471b 100644 --- a/codex-rs/terminal-browser/Cargo.toml +++ b/codex-rs/terminal-browser/Cargo.toml @@ -14,6 +14,7 @@ workspace = true [dependencies] anyhow = { workspace = true } +codex-network-proxy = { workspace = true } codex-protocol = { workspace = true } codex-sandboxing = { workspace = true } codex-utils-absolute-path = { workspace = true } @@ -21,12 +22,11 @@ codex-utils-path-uri = { workspace = true } codex-utils-pty = { workspace = true } futures = { workspace = true } libc = { workspace = true } -reqwest = { workspace = true, features = ["json"] } serde = { workspace = true, features = ["derive"] } serde_json = { workspace = true } sha2 = { workspace = true } tempfile = { workspace = true } -tokio = { workspace = true, features = ["macros", "net", "process", "rt", "sync", "time"] } +tokio = { workspace = true, features = ["io-util", "macros", "net", "process", "rt", "sync", "time"] } tokio-tungstenite = { workspace = true } tracing = { workspace = true } url = { workspace = true } diff --git a/codex-rs/terminal-browser/src/cdp.rs b/codex-rs/terminal-browser/src/cdp.rs index 2e14fc5402..a27ccfec6c 100644 --- a/codex-rs/terminal-browser/src/cdp.rs +++ b/codex-rs/terminal-browser/src/cdp.rs @@ -9,28 +9,46 @@ use std::time::Duration; use anyhow::Context; use anyhow::Result; use anyhow::bail; -use futures::SinkExt; -use futures::StreamExt; +use serde::Deserialize; use serde_json::Value; use serde_json::json; +use tokio::io::AsyncRead; +use tokio::io::AsyncReadExt; +use tokio::io::AsyncWrite; +use tokio::io::AsyncWriteExt; +#[cfg(test)] use tokio::net::TcpStream; use tokio::sync::broadcast; use tokio::sync::mpsc; use tokio::sync::oneshot; use tokio::time::timeout; + +#[cfg(test)] +use futures::SinkExt; +#[cfg(test)] +use futures::StreamExt; +#[cfg(test)] use tokio_tungstenite::MaybeTlsStream; +#[cfg(test)] use tokio_tungstenite::WebSocketStream; +#[cfg(test)] use tokio_tungstenite::connect_async; +#[cfg(test)] use tokio_tungstenite::tungstenite::Message; +#[cfg(test)] type Socket = WebSocketStream>; type PendingCalls = Arc>>>>; +type PageSessionId = Arc>>; -const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); -const CALL_TIMEOUT: Duration = Duration::from_secs(15); +#[cfg(test)] +const CONNECT_TIMEOUT: Duration = Duration::from_secs(/*secs*/ 10); +const CALL_TIMEOUT: Duration = Duration::from_secs(/*secs*/ 15); const OUTBOUND_CAPACITY: usize = 64; const EVENT_CAPACITY: usize = 256; const MAX_EVENT_BYTES: usize = 16 * 1024; +const MAX_CDP_FRAME_BYTES: usize = 16 * 1024 * 1024; +const PIPE_READ_BUFFER_BYTES: usize = 8 * 1024; #[derive(Clone, Debug, PartialEq)] pub(crate) enum CdpEvent { @@ -38,20 +56,28 @@ pub(crate) enum CdpEvent { Disconnected(String), } +pub(crate) struct ConnectedPage { + pub(crate) client: CdpClient, + pub(crate) title: String, +} + #[derive(Clone)] pub(crate) struct CdpClient { - outbound: mpsc::Sender, + outbound: mpsc::Sender, pending: PendingCalls, events: broadcast::Sender, next_id: Arc, + page_session_id: PageSessionId, call_timeout: Duration, } impl CdpClient { + #[cfg(test)] pub(crate) async fn connect(websocket_url: &str) -> Result { Self::connect_with_call_timeout(websocket_url, CALL_TIMEOUT).await } + #[cfg(test)] async fn connect_with_call_timeout( websocket_url: &str, call_timeout: Duration, @@ -59,42 +85,153 @@ impl CdpClient { let (socket, _) = timeout(CONNECT_TIMEOUT, connect_async(websocket_url)) .await .context("timed out connecting to Carbonyl DevTools")??; - let (outbound, outbound_rx) = mpsc::channel(OUTBOUND_CAPACITY); - let (events, _) = broadcast::channel(EVENT_CAPACITY); - let client = Self { - outbound, - pending: Arc::new(Mutex::new(HashMap::new())), - events, - next_id: Arc::new(AtomicU64::new(/*v*/ 1)), - call_timeout, - }; - tokio::spawn(run_pump( + let (client, outbound_rx) = Self::new(call_timeout); + tokio::spawn(run_websocket_pump( socket, outbound_rx, client.pending.clone(), client.events.clone(), + client.page_session_id.clone(), )); - client.call("Page.enable", json!({})).await?; - client.call("Runtime.enable", json!({})).await?; - client.call("DOM.enable", json!({})).await?; - client.call("Accessibility.enable", json!({})).await?; - client - .call("Page.setLifecycleEventsEnabled", json!({ "enabled": true })) - .await?; + client.initialize_page().await?; Ok(client) } + #[cfg(unix)] + pub(crate) async fn connect_pipe( + reader: std::os::unix::net::UnixStream, + writer: std::os::unix::net::UnixStream, + ) -> Result { + Self::connect_pipe_io( + tokio::net::UnixStream::from_std(reader) + .context("adopt Carbonyl DevTools output pipe")?, + tokio::net::UnixStream::from_std(writer) + .context("adopt Carbonyl DevTools input pipe")?, + CALL_TIMEOUT, + ) + .await + } + + async fn connect_pipe_io( + reader: R, + writer: W, + call_timeout: Duration, + ) -> Result + where + R: AsyncRead + Send + Unpin + 'static, + W: AsyncWrite + Send + Unpin + 'static, + { + let (client, outbound_rx) = Self::new(call_timeout); + tokio::spawn(run_pipe_pump( + reader, + writer, + outbound_rx, + client.pending.clone(), + client.events.clone(), + client.page_session_id.clone(), + )); + + let target_list: TargetList = + serde_json::from_value(client.call_root("Target.getTargets", json!({})).await?) + .context("decode Carbonyl page targets")?; + let target = target_list + .target_infos + .into_iter() + .find(|target| target.kind == "page") + .context("Carbonyl DevTools did not expose a page target")?; + let attached: AttachedTarget = serde_json::from_value( + client + .call_root( + "Target.attachToTarget", + json!({ + "targetId": target.target_id, + "flatten": true, + }), + ) + .await?, + ) + .context("decode Carbonyl page session")?; + *client + .page_session_id + .lock() + .unwrap_or_else(PoisonError::into_inner) = Some(attached.session_id); + client.initialize_page().await?; + Ok(ConnectedPage { + client, + title: target.title, + }) + } + + fn new(call_timeout: Duration) -> (Self, mpsc::Receiver) { + let (outbound, outbound_rx) = mpsc::channel(OUTBOUND_CAPACITY); + let (events, _) = broadcast::channel(EVENT_CAPACITY); + ( + Self { + outbound, + pending: Arc::new(Mutex::new(/*t*/ HashMap::new())), + events, + next_id: Arc::new(AtomicU64::new(/*v*/ 1)), + page_session_id: Arc::new(Mutex::new(/*t*/ None)), + call_timeout, + }, + outbound_rx, + ) + } + + async fn initialize_page(&self) -> Result<()> { + self.call("Page.enable", json!({})).await?; + self.call("Runtime.enable", json!({})).await?; + self.call("DOM.enable", json!({})).await?; + self.call("Accessibility.enable", json!({})).await?; + self.call("Page.setLifecycleEventsEnabled", json!({ "enabled": true })) + .await?; + Ok(()) + } + pub(crate) fn subscribe_events(&self) -> broadcast::Receiver { self.events.subscribe() } pub(crate) async fn call(&self, method: &str, params: Value) -> Result { + let page_session_id = self + .page_session_id + .lock() + .unwrap_or_else(PoisonError::into_inner) + .clone(); + self.call_with_session(method, params, page_session_id.as_deref()) + .await + } + + pub(crate) async fn call_browser(&self, method: &str, params: Value) -> Result { + self.call_root(method, params).await + } + + async fn call_root(&self, method: &str, params: Value) -> Result { + self.call_with_session(method, params, /*session_id*/ None) + .await + } + + async fn call_with_session( + &self, + method: &str, + params: Value, + session_id: Option<&str>, + ) -> Result { let id = self.next_id.fetch_add(/*val*/ 1, Ordering::Relaxed); - let request = json!({ + let mut request = json!({ "id": id, "method": method, "params": params, }); + if let Some(session_id) = session_id { + request["sessionId"] = Value::String(session_id.to_string()); + } + let encoded = serde_json::to_string(&request).context("encode CDP request")?; + anyhow::ensure!( + encoded.len() <= MAX_CDP_FRAME_BYTES, + "CDP {method} request exceeded the {MAX_CDP_FRAME_BYTES}-byte frame limit" + ); + let (response_tx, response_rx) = oneshot::channel(); self.pending .lock() @@ -106,12 +243,7 @@ impl CdpClient { }; let result = timeout(self.call_timeout, async { - if self - .outbound - .send(Message::Text(request.to_string().into())) - .await - .is_err() - { + if self.outbound.send(encoded).await.is_err() { bail!("Carbonyl closed the DevTools connection"); } response_rx @@ -154,6 +286,28 @@ impl CdpClient { } } +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +struct TargetList { + target_infos: Vec, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +struct TargetInfo { + target_id: String, + #[serde(rename = "type")] + kind: String, + #[serde(default)] + title: String, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +struct AttachedTarget { + session_id: String, +} + struct PendingCallGuard { pending: PendingCalls, id: u64, @@ -168,22 +322,112 @@ impl Drop for PendingCallGuard { } } -async fn run_pump( - mut socket: Socket, - mut outbound: mpsc::Receiver, +struct PipeFrameDecoder { + frame: Vec, + max_frame_bytes: usize, +} + +impl PipeFrameDecoder { + fn new(max_frame_bytes: usize) -> Self { + Self { + frame: Vec::with_capacity(PIPE_READ_BUFFER_BYTES.min(max_frame_bytes)), + max_frame_bytes, + } + } + + fn push(&mut self, chunk: &[u8]) -> Result> { + let mut messages = Vec::new(); + for &byte in chunk { + if byte == 0 { + let message = serde_json::from_slice(&self.frame) + .context("failed to decode DevTools pipe message")?; + messages.push(message); + self.frame.clear(); + continue; + } + anyhow::ensure!( + self.frame.len() < self.max_frame_bytes, + "DevTools pipe message exceeded the {}-byte frame limit", + self.max_frame_bytes + ); + self.frame.push(byte); + } + Ok(messages) + } +} + +async fn run_pipe_pump( + mut reader: R, + writer: W, + outbound: mpsc::Receiver, pending: PendingCalls, events: broadcast::Sender, + page_session_id: PageSessionId, +) where + R: AsyncRead + Send + Unpin + 'static, + W: AsyncWrite + Send + Unpin + 'static, +{ + let mut writer_task = tokio::spawn(run_pipe_writer(writer, outbound)); + let mut decoder = PipeFrameDecoder::new(MAX_CDP_FRAME_BYTES); + let mut read_buffer = [0_u8; PIPE_READ_BUFFER_BYTES]; + let disconnect_reason = loop { + tokio::select! { + incoming = reader.read(&mut read_buffer) => match incoming { + Ok(0) => break "Carbonyl closed the DevTools connection".to_string(), + Ok(count) => match decoder.push(&read_buffer[..count]) { + Ok(messages) => { + for message in messages { + dispatch_message(message, &pending, &events, &page_session_id); + } + } + Err(error) => break error.to_string(), + }, + Err(error) => break format!("DevTools pipe connection failed: {error}"), + }, + writer_result = &mut writer_task => break match writer_result { + Ok(reason) => reason, + Err(error) => format!("DevTools pipe writer stopped: {error}"), + }, + } + }; + + writer_task.abort(); + finish_disconnect(&pending, &events, disconnect_reason); +} + +async fn run_pipe_writer(mut writer: W, mut outbound: mpsc::Receiver) -> String +where + W: AsyncWrite + Send + Unpin + 'static, +{ + while let Some(message) = outbound.recv().await { + if let Err(error) = writer.write_all(message.as_bytes()).await { + return format!("failed to send DevTools pipe message: {error}"); + } + if let Err(error) = writer.write_all(&[0]).await { + return format!("failed to delimit DevTools pipe message: {error}"); + } + } + "DevTools client closed".to_string() +} + +#[cfg(test)] +async fn run_websocket_pump( + mut socket: Socket, + mut outbound: mpsc::Receiver, + pending: PendingCalls, + events: broadcast::Sender, + page_session_id: PageSessionId, ) { let disconnect_reason = loop { tokio::select! { outgoing = outbound.recv() => match outgoing { Some(message) => { - if let Err(error) = socket.send(message).await { + if let Err(error) = socket.send(Message::Text(message.into())).await { break format!("failed to send DevTools message: {error}"); } } None => { - let _ = socket.close(None).await; + let _ = socket.close(/*msg*/ None).await; break "DevTools client closed".to_string(); } }, @@ -193,7 +437,7 @@ async fn run_pump( Ok(response) => response, Err(error) => break format!("failed to decode DevTools message: {error}"), }; - dispatch_message(response, &pending, &events); + dispatch_message(response, &pending, &events, &page_session_id); } Some(Ok(Message::Ping(bytes))) => { if let Err(error) = socket.send(Message::Pong(bytes)).await { @@ -210,11 +454,15 @@ async fn run_pump( } }; - fail_pending(&pending, &disconnect_reason); - let _ = events.send(CdpEvent::Disconnected(disconnect_reason)); + finish_disconnect(&pending, &events, disconnect_reason); } -fn dispatch_message(response: Value, pending: &PendingCalls, events: &broadcast::Sender) { +fn dispatch_message( + response: Value, + pending: &PendingCalls, + events: &broadcast::Sender, + page_session_id: &PageSessionId, +) { let response_tx = response.get("id").and_then(Value::as_u64).and_then(|id| { pending .lock() @@ -223,12 +471,23 @@ fn dispatch_message(response: Value, pending: &PendingCalls, events: &broadcast: }); if let Some(response_tx) = response_tx { let _ = response_tx.send(Ok(response)); - } else if should_broadcast_event(&response) { + return; + } + let page_session_id = page_session_id + .lock() + .unwrap_or_else(PoisonError::into_inner) + .clone(); + if should_broadcast_event(&response, page_session_id.as_deref()) { let _ = events.send(CdpEvent::Message(response)); } } -fn should_broadcast_event(message: &Value) -> bool { +fn should_broadcast_event(message: &Value, page_session_id: Option<&str>) -> bool { + if let Some(page_session_id) = page_session_id + && message.get("sessionId").and_then(Value::as_str) != Some(page_session_id) + { + return false; + } let Some(method) = message.get("method").and_then(Value::as_str) else { return false; }; @@ -242,6 +501,11 @@ fn should_broadcast_event(message: &Value) -> bool { ) && serde_json::to_vec(message).is_ok_and(|encoded| encoded.len() <= MAX_EVENT_BYTES) } +fn finish_disconnect(pending: &PendingCalls, events: &broadcast::Sender, reason: String) { + fail_pending(pending, &reason); + let _ = events.send(CdpEvent::Disconnected(reason)); +} + fn fail_pending(pending: &PendingCalls, reason: &str) { let pending = std::mem::take(&mut *pending.lock().unwrap_or_else(PoisonError::into_inner)); for response_tx in pending.into_values() { diff --git a/codex-rs/terminal-browser/src/cdp_tests.rs b/codex-rs/terminal-browser/src/cdp_tests.rs index 19c16ca70d..ea85155245 100644 --- a/codex-rs/terminal-browser/src/cdp_tests.rs +++ b/codex-rs/terminal-browser/src/cdp_tests.rs @@ -5,6 +5,10 @@ use futures::StreamExt; use pretty_assertions::assert_eq; use serde_json::Value; use serde_json::json; +use tokio::io::AsyncRead; +use tokio::io::AsyncReadExt; +use tokio::io::AsyncWrite; +use tokio::io::AsyncWriteExt; use tokio::net::TcpListener; use tokio::task::JoinHandle; use tokio_tungstenite::accept_async; @@ -12,6 +16,9 @@ use tokio_tungstenite::tungstenite::Message; use super::CdpClient; use super::CdpEvent; +use super::PipeFrameDecoder; + +const TEST_PIPE_CAPACITY: usize = 64 * 1024; async fn test_server( handler: impl FnOnce(tokio_tungstenite::WebSocketStream) -> JoinHandle<()> @@ -275,3 +282,173 @@ async fn concurrent_calls_match_reverse_order_responses() { )); server.await.expect("server task"); } + +async fn receive_pipe_request(reader: &mut (impl AsyncRead + Unpin)) -> Value { + let mut frame = Vec::new(); + loop { + let byte = reader.read_u8().await.expect("read pipe request"); + if byte == 0 { + return serde_json::from_slice(&frame).expect("decode pipe request"); + } + frame.push(byte); + } +} + +fn framed_pipe_message(message: &Value) -> Vec { + let mut frame = serde_json::to_vec(message).expect("encode pipe message"); + frame.push(/*value*/ 0); + frame +} + +async fn respond_pipe(writer: &mut (impl AsyncWrite + Unpin), request: &Value, result: Value) { + let mut response = json!({ "id": request["id"], "result": result }); + if let Some(session_id) = request.get("sessionId") { + response["sessionId"] = session_id.clone(); + } + writer + .write_all(&framed_pipe_message(&response)) + .await + .expect("write pipe response"); +} + +async fn complete_pipe_initialization( + reader: &mut (impl AsyncRead + Unpin), + writer: &mut (impl AsyncWrite + Unpin), +) { + let targets = receive_pipe_request(reader).await; + assert_eq!(targets["method"], "Target.getTargets"); + assert_eq!(targets.get("sessionId"), None); + let response = json!({ + "id": targets["id"], + "result": { + "targetInfos": [{ + "targetId": "page-target", + "type": "page", + "title": "Pipe page" + }] + } + }); + let response = framed_pipe_message(&response); + let split = response.len() / 2; + writer + .write_all(&response[..split]) + .await + .expect("write first response fragment"); + tokio::task::yield_now().await; + writer + .write_all(&response[split..]) + .await + .expect("write second response fragment"); + + let attach = receive_pipe_request(reader).await; + assert_eq!(attach["method"], "Target.attachToTarget"); + assert_eq!(attach["params"]["targetId"], "page-target"); + assert_eq!(attach["params"]["flatten"], true); + assert_eq!(attach.get("sessionId"), None); + respond_pipe(writer, &attach, json!({ "sessionId": "page-session" })).await; + + for expected_method in [ + "Page.enable", + "Runtime.enable", + "DOM.enable", + "Accessibility.enable", + "Page.setLifecycleEventsEnabled", + ] { + let request = receive_pipe_request(reader).await; + assert_eq!(request["method"], expected_method); + assert_eq!(request["sessionId"], "page-session"); + respond_pipe(writer, &request, json!({})).await; + } +} + +#[tokio::test] +async fn pipe_transport_attaches_and_routes_page_messages_by_session() { + let (client_pipe, server_pipe) = tokio::io::duplex(TEST_PIPE_CAPACITY); + let (client_read, client_write) = tokio::io::split(client_pipe); + let (mut server_read, mut server_write) = tokio::io::split(server_pipe); + let server = tokio::spawn(async move { + complete_pipe_initialization(&mut server_read, &mut server_write).await; + let browser_request = receive_pipe_request(&mut server_read).await; + assert_eq!(browser_request["method"], "Test.browserCommand"); + assert_eq!(browser_request.get("sessionId"), None); + respond_pipe(&mut server_write, &browser_request, json!({})).await; + + let request = receive_pipe_request(&mut server_read).await; + assert_eq!(request["method"], "Test.command"); + assert_eq!(request["sessionId"], "page-session"); + + let messages = [ + json!({ + "method": "Page.loadEventFired", + "params": { "timestamp": 1 }, + "sessionId": "another-session" + }), + json!({ + "method": "Page.loadEventFired", + "params": { "timestamp": 2 }, + "sessionId": "page-session" + }), + json!({ + "id": request["id"], + "result": { "ok": true }, + "sessionId": "page-session" + }), + ]; + let mut frames = Vec::new(); + for message in messages { + frames.extend(framed_pipe_message(&message)); + } + server_write + .write_all(&frames) + .await + .expect("write combined pipe frames"); + }); + let connection = + CdpClient::connect_pipe_io(client_read, client_write, Duration::from_secs(/*secs*/ 5)) + .await + .expect("connect pipe client"); + assert_eq!(connection.title, "Pipe page"); + let mut events = connection.client.subscribe_events(); + connection + .client + .call_browser("Test.browserCommand", json!({})) + .await + .expect("browser-target pipe command"); + + let result = connection + .client + .call("Test.command", json!({})) + .await + .expect("pipe command response"); + + assert_eq!(result, json!({ "ok": true })); + assert_eq!( + events.recv().await.expect("page-session lifecycle event"), + CdpEvent::Message(json!({ + "method": "Page.loadEventFired", + "params": { "timestamp": 2 }, + "sessionId": "page-session" + })) + ); + server.await.expect("server task"); +} + +#[test] +fn pipe_frame_decoder_handles_fragments_and_enforces_its_limit() { + let mut decoder = PipeFrameDecoder::new(/*max_frame_bytes*/ 8); + + assert_eq!( + decoder.push(b"{\"id\"").expect("first fragment"), + Vec::::new() + ); + assert_eq!( + decoder.push(b":1}\0{\"id\":2}\0").expect("complete frames"), + vec![json!({ "id": 1 }), json!({ "id": 2 })] + ); + + let mut decoder = PipeFrameDecoder::new(/*max_frame_bytes*/ 4); + let error = decoder + .push(b"12345") + .expect_err("oversized frame must fail"); + assert!(error.to_string().contains("4-byte frame limit")); +} diff --git a/codex-rs/terminal-browser/src/devtools.rs b/codex-rs/terminal-browser/src/devtools.rs index ec960d5a70..caeb1ef072 100644 --- a/codex-rs/terminal-browser/src/devtools.rs +++ b/codex-rs/terminal-browser/src/devtools.rs @@ -1,127 +1,11 @@ -use std::path::Path; -use std::time::Duration; -use std::time::Instant; - use anyhow::Context; use anyhow::Result; -use reqwest::Client; -use serde::Deserialize; -use url::Host; -use url::Url; use crate::cdp::CdpClient; -#[derive(Deserialize)] -#[serde(rename_all = "camelCase")] -struct DevtoolsTarget { - #[serde(rename = "type")] - kind: String, - #[serde(default)] - title: String, - web_socket_debugger_url: Option, -} - -pub(crate) struct PageTarget { - pub(crate) title: String, - pub(crate) websocket_url: String, -} - -pub(crate) fn prepare_devtools_active_port(profile: &Path) -> Result<()> { - let active_port_path = profile.join("DevToolsActivePort"); - let metadata = match std::fs::symlink_metadata(&active_port_path) { - Ok(metadata) => metadata, - Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()), - Err(error) => { - return Err(error).context("inspect Carbonyl DevToolsActivePort before startup"); - } - }; - anyhow::ensure!( - !metadata.file_type().is_symlink(), - "refusing symbolic-link Carbonyl DevToolsActivePort" - ); - anyhow::ensure!( - metadata.is_file(), - "refusing non-file Carbonyl DevToolsActivePort" - ); - std::fs::remove_file(&active_port_path) - .context("remove stale Carbonyl DevToolsActivePort before startup") -} - -pub(crate) async fn discover_page_target(profile: &Path) -> Result { - let client = Client::builder() - .no_proxy() - .timeout(Duration::from_secs(/*secs*/ 2)) - .build()?; - let deadline = Instant::now() + Duration::from_secs(/*secs*/ 12); - let active_port_path = profile.join("DevToolsActivePort"); - let port = loop { - match std::fs::read_to_string(&active_port_path) { - Ok(contents) => { - let port = contents - .lines() - .next() - .context("DevToolsActivePort did not contain a port")? - .parse::() - .context("DevToolsActivePort contained an invalid port")?; - anyhow::ensure!(port != 0, "DevToolsActivePort contained port zero"); - break port; - } - Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} - Err(error) => { - return Err(error).context("read Carbonyl DevToolsActivePort"); - } - } - anyhow::ensure!( - Instant::now() < deadline, - "timed out waiting for Carbonyl DevToolsActivePort" - ); - tokio::time::sleep(Duration::from_millis(/*millis*/ 100)).await; - }; - let endpoint = format!("http://127.0.0.1:{port}/json/list"); - loop { - if let Ok(response) = client.get(&endpoint).send().await - && let Ok(targets) = response.json::>().await - { - for target in targets { - if target.kind != "page" { - continue; - } - let Some(websocket_url) = target.web_socket_debugger_url else { - continue; - }; - return Ok(PageTarget { - title: target.title, - websocket_url: validated_websocket_url(&websocket_url, port)?, - }); - } - } - anyhow::ensure!( - Instant::now() < deadline, - "timed out waiting for Carbonyl DevTools on {endpoint}" - ); - tokio::time::sleep(Duration::from_millis(/*millis*/ 100)).await; - } -} - -pub(crate) fn validated_websocket_url(websocket_url: &str, expected_port: u16) -> Result { - let parsed = Url::parse(websocket_url).context("parse Carbonyl DevTools WebSocket URL")?; - anyhow::ensure!(parsed.scheme() == "ws", "Carbonyl DevTools URL must use ws"); - let loopback = match parsed.host() { - Some(Host::Ipv4(address)) => address.is_loopback(), - Some(Host::Ipv6(address)) => address.is_loopback(), - Some(Host::Domain(_)) | None => false, - }; - anyhow::ensure!(loopback, "Carbonyl DevTools URL must use a loopback host"); - anyhow::ensure!( - parsed.port() == Some(expected_port), - "Carbonyl DevTools URL used an unexpected port" - ); - Ok(parsed.to_string()) -} - pub(crate) async fn deny_downloads(client: &CdpClient) -> Result<()> { if client - .call( + .call_browser( "Browser.setDownloadBehavior", serde_json::json!({ "behavior": "deny" }), ) diff --git a/codex-rs/terminal-browser/src/process.rs b/codex-rs/terminal-browser/src/process.rs index ddc224e739..ccae2bdc5c 100644 --- a/codex-rs/terminal-browser/src/process.rs +++ b/codex-rs/terminal-browser/src/process.rs @@ -7,15 +7,18 @@ use anyhow::Context; use anyhow::Result; use anyhow::anyhow; use codex_utils_pty::ProcessHandle; -use codex_utils_pty::spawn_pty_process; +#[cfg(unix)] +use codex_utils_pty::pty::ChildFdMapping; +#[cfg(unix)] +use codex_utils_pty::pty::spawn_process_with_fd_mappings; use tokio::sync::mpsc; use tokio::sync::watch; use tokio::task::JoinHandle; use crate::cdp::CdpClient; +#[cfg(unix)] +use crate::cdp::ConnectedPage; use crate::devtools::deny_downloads; -use crate::devtools::discover_page_target; -use crate::devtools::prepare_devtools_active_port; use crate::handles::BrowserHandles; use crate::network::BrowserNetworkPolicy; use crate::profile::BrowserProfileLock; @@ -192,6 +195,7 @@ impl Inner { Ok(path) } + #[cfg(unix)] #[expect( clippy::await_holding_invalid_type, reason = "the session guard is explicitly dropped before the conditional close await" @@ -206,9 +210,7 @@ impl Inner { let persistent_profile = self.selected_profile_resources()?; let runtime = BrowserRuntime::create(persistent_profile.as_ref().map(|(profile, _lock)| profile))?; - prepare_devtools_active_port(&runtime.profile)?; let args = carbonyl_args( - /*debugging_port*/ 0, runtime.profile.as_path().to_string_lossy().as_ref(), &config.network_policy, config.render_mode, @@ -225,14 +227,32 @@ impl Inner { &config.network_policy, &self.launch_context, )?; + let (browser_read, codex_write) = std::os::unix::net::UnixStream::pair() + .context("create Carbonyl DevTools input pipe")?; + let (codex_read, browser_write) = std::os::unix::net::UnixStream::pair() + .context("create Carbonyl DevTools output pipe")?; + codex_write + .set_nonblocking(/*nonblocking*/ true) + .context("configure Carbonyl DevTools input pipe")?; + codex_read + .set_nonblocking(/*nonblocking*/ true) + .context("configure Carbonyl DevTools output pipe")?; + let child_fd_mappings = [ + ChildFdMapping::new(&browser_read, /*child_fd*/ 3), + ChildFdMapping::new(&browser_write, /*child_fd*/ 4), + ]; + // `dup2` clears `FD_CLOEXEC` on these fixed targets. Bubblewrap's monitor close-fd path + // explicitly passes any other non-CLOEXEC fds on to the child, so the Linux sandbox + // wrapper needs no fd-specific flag or ambient-descriptor widening for the CDP pipe pair. let size = *self.resize_tx.borrow(); - let spawned = match spawn_pty_process( + let spawned = match spawn_process_with_fd_mappings( &launch.program, &launch.args, launch.cwd.as_path(), &launch.env, &launch.arg0, size.into(), + &child_fd_mappings, ) .await { @@ -241,6 +261,8 @@ impl Inner { return Err(error).context("launch Carbonyl"); } }; + drop(browser_read); + drop(browser_write); let process = Arc::new(spawned.session); let writer = process.writer_sender(); @@ -270,12 +292,11 @@ impl Inner { startup.set_output_task(output_task); let mut exit_rx = spawned.exit_rx; - let target = tokio::select! { - target = discover_page_target(&runtime.profile) => target, + let connection = tokio::select! { + connection = CdpClient::connect_pipe(codex_read, codex_write) => connection, exit = &mut exit_rx => Err(anyhow!("Carbonyl exited during startup: {exit:?}")), }; - let target = target?; - let cdp = CdpClient::connect(&target.websocket_url).await?; + let ConnectedPage { client: cdp, title } = connection?; deny_downloads(&cdp).await?; startup.set_navigation_policy_task(spawn_navigation_policy_task(self.clone(), cdp.clone())); let exit_inner = self.clone(); @@ -314,8 +335,8 @@ impl Inner { self.update_view(|view| { if !self.terminated.load(Ordering::SeqCst) { view.status = BrowserStatus::Running; - if !target.title.is_empty() { - view.title = Some(target.title); + if !title.is_empty() { + view.title = Some(title); } } }); @@ -326,6 +347,11 @@ impl Inner { Ok(()) } + #[cfg(not(unix))] + async fn start_session(self: &Arc, _config: SessionConfig) -> Result<()> { + anyhow::bail!("Carbonyl terminal browsing is only supported on macOS and Linux") + } + fn selected_profile_resources( &self, ) -> Result< @@ -516,14 +542,12 @@ async fn wait_for_exit(process: Arc) { } pub(crate) fn carbonyl_args( - debugging_port: u16, profile: &str, network_policy: &BrowserNetworkPolicy, render_mode: RenderMode, ) -> Vec { let mut args = vec![ - "--remote-debugging-address=127.0.0.1".to_string(), - format!("--remote-debugging-port={debugging_port}"), + "--remote-debugging-pipe".to_string(), format!("--user-data-dir={profile}"), "--disable-extensions".to_string(), "--disable-background-networking".to_string(), diff --git a/codex-rs/terminal-browser/src/process_tests.rs b/codex-rs/terminal-browser/src/process_tests.rs index 784fb50c10..7de7482ab3 100644 --- a/codex-rs/terminal-browser/src/process_tests.rs +++ b/codex-rs/terminal-browser/src/process_tests.rs @@ -5,13 +5,14 @@ use std::sync::atomic::Ordering; use codex_utils_pty::ProcessDriver; use codex_utils_pty::spawn_from_driver; -use pretty_assertions::assert_eq; use tokio::sync::broadcast; use tokio::sync::mpsc; use tokio::sync::oneshot; use super::StartupGuard; -use crate::devtools::prepare_devtools_active_port; +use super::carbonyl_args; +use crate::network::BrowserNetworkPolicy; +use crate::session::RenderMode; #[tokio::test] async fn canceled_startup_terminates_the_process_and_aborts_output() { @@ -47,52 +48,22 @@ async fn canceled_startup_terminates_the_process_and_aborts_output() { } #[test] -fn stale_devtools_active_port_is_removed_before_spawn() { - let profile = tempfile::tempdir().expect("profile directory"); - let active_port = profile.path().join("DevToolsActivePort"); - std::fs::write(&active_port, "9222\n").expect("stale active port"); +fn carbonyl_uses_only_the_inherited_devtools_pipe() { + let args = carbonyl_args( + "/tmp/profile", + &BrowserNetworkPolicy::Direct, + RenderMode::NativeText, + ); - prepare_devtools_active_port(profile.path()).expect("prepare active port"); - - assert!(!active_port.exists()); - prepare_devtools_active_port(profile.path()).expect("missing active port is valid"); -} - -#[test] -fn non_file_devtools_active_port_is_rejected() { - let profile = tempfile::tempdir().expect("profile directory"); - let active_port = profile.path().join("DevToolsActivePort"); - std::fs::create_dir(&active_port).expect("active port directory"); - - let error = prepare_devtools_active_port(profile.path()) - .expect_err("active port directory must be rejected"); - - assert_eq!( - error.to_string(), - "refusing non-file Carbonyl DevToolsActivePort" - ); -} - -#[cfg(unix)] -#[test] -fn symbolic_link_devtools_active_port_is_rejected() { - use std::os::unix::fs::symlink; - - let profile = tempfile::tempdir().expect("profile directory"); - let target = profile.path().join("target"); - let active_port = profile.path().join("DevToolsActivePort"); - std::fs::write(&target, "9222\n").expect("target active port"); - symlink(&target, &active_port).expect("active port symlink"); - - let error = prepare_devtools_active_port(profile.path()) - .expect_err("active port symlink must be rejected"); - - assert_eq!( - error.to_string(), - "refusing symbolic-link Carbonyl DevToolsActivePort" - ); - assert_eq!( - std::fs::read_to_string(target).expect("target remains"), - "9222\n" + assert!(args.contains(&"--remote-debugging-pipe".to_string())); + assert!( + !args + .iter() + .any(|arg| arg.starts_with("--remote-debugging-address=")) + ); + assert!( + !args + .iter() + .any(|arg| arg.starts_with("--remote-debugging-port=")) ); } diff --git a/codex-rs/terminal-browser/src/terminal_browser_tests.rs b/codex-rs/terminal-browser/src/terminal_browser_tests.rs index 4a871e49bc..3adf770857 100644 --- a/codex-rs/terminal-browser/src/terminal_browser_tests.rs +++ b/codex-rs/terminal-browser/src/terminal_browser_tests.rs @@ -6,7 +6,6 @@ use crate::BrowserNetworkPolicy; use crate::TerminalBrowser; use crate::TerminalSize; use crate::actions::bounded_snapshot_json; -use crate::devtools::validated_websocket_url; use crate::handles::BrowserHandles; use crate::human_control::HumanControlStateTransition; use crate::process::carbonyl_args; @@ -331,6 +330,28 @@ async fn configured_real_carbonyl_opens_local_page_and_snapshots() { }; assert!(snapshot.contains("Smoke button")); + let render_deadline = + std::time::Instant::now() + std::time::Duration::from_secs(/*secs*/ 5); + loop { + let view = browser.view(); + let rendered = (0..view.screen.rows) + .map(|row| { + (0..view.screen.cols) + .filter_map(|col| view.screen.cell(row, col)) + .map(|cell| cell.text.as_str()) + .collect::() + }) + .collect::>() + .join("\n"); + if rendered.contains("Smoke button") { + break; + } + assert!( + std::time::Instant::now() < render_deadline, + "Carbonyl did not render the smoke page in its PTY:\n{rendered}" + ); + tokio::time::sleep(std::time::Duration::from_millis(/*millis*/ 20)).await; + } browser.close().await; server.join().expect("smoke server thread"); } @@ -406,7 +427,6 @@ fn terminal_query_responder_handles_queries_split_across_chunks() { #[test] fn carbonyl_uses_the_managed_proxy_only_when_policy_provides_one() { let direct = carbonyl_args( - /*debugging_port*/ 9_222, "/tmp/profile", &BrowserNetworkPolicy::Direct, RenderMode::NativeText, @@ -415,7 +435,6 @@ fn carbonyl_uses_the_managed_proxy_only_when_policy_provides_one() { let http_addr = "127.0.0.1:43128".parse().expect("valid proxy address"); let proxied = carbonyl_args( - /*debugging_port*/ 9_222, "/tmp/profile", &BrowserNetworkPolicy::ManagedProxy { http_addr }, RenderMode::NativeText, @@ -424,39 +443,6 @@ fn carbonyl_uses_the_managed_proxy_only_when_policy_provides_one() { assert!(proxied.contains(&"--proxy-bypass-list=<-loopback>".to_string())); } -#[test] -fn carbonyl_devtools_websocket_must_match_the_private_loopback_listener() { - assert_eq!( - validated_websocket_url( - "ws://127.0.0.1:9222/devtools/page/1", - /*expected_port*/ 9_222, - ) - .unwrap(), - "ws://127.0.0.1:9222/devtools/page/1" - ); - assert!( - validated_websocket_url( - "ws://192.0.2.10:9222/devtools/page/1", - /*expected_port*/ 9_222, - ) - .is_err() - ); - assert!( - validated_websocket_url( - "ws://127.0.0.1:9223/devtools/page/1", - /*expected_port*/ 9_222, - ) - .is_err() - ); - assert!( - validated_websocket_url( - "wss://127.0.0.1:9222/devtools/page/1", - /*expected_port*/ 9_222, - ) - .is_err() - ); -} - #[test] fn browser_node_handles_reject_unknown_and_stale_ids() { let mut handles = BrowserHandles::default();