feat(terminal-browser): use cdp pipe transport

This commit is contained in:
Felipe Coury
2026-07-06 16:41:41 -03:00
parent 23484dd8d5
commit 2046175b2c
8 changed files with 565 additions and 259 deletions

2
codex-rs/Cargo.lock generated
View File

@@ -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",

View File

@@ -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 }

View File

@@ -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<MaybeTlsStream<TcpStream>>;
type PendingCalls = Arc<Mutex<HashMap<u64, oneshot::Sender<Result<Value, String>>>>>;
type PageSessionId = Arc<Mutex<Option<String>>>;
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<Message>,
outbound: mpsc::Sender<String>,
pending: PendingCalls,
events: broadcast::Sender<CdpEvent>,
next_id: Arc<AtomicU64>,
page_session_id: PageSessionId,
call_timeout: Duration,
}
impl CdpClient {
#[cfg(test)]
pub(crate) async fn connect(websocket_url: &str) -> Result<Self> {
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<ConnectedPage> {
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<R, W>(
reader: R,
writer: W,
call_timeout: Duration,
) -> Result<ConnectedPage>
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<String>) {
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<CdpEvent> {
self.events.subscribe()
}
pub(crate) async fn call(&self, method: &str, params: Value) -> Result<Value> {
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<Value> {
self.call_root(method, params).await
}
async fn call_root(&self, method: &str, params: Value) -> Result<Value> {
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<Value> {
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<TargetInfo>,
}
#[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<Message>,
struct PipeFrameDecoder {
frame: Vec<u8>,
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<Vec<Value>> {
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<R, W>(
mut reader: R,
writer: W,
outbound: mpsc::Receiver<String>,
pending: PendingCalls,
events: broadcast::Sender<CdpEvent>,
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<W>(mut writer: W, mut outbound: mpsc::Receiver<String>) -> 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<String>,
pending: PendingCalls,
events: broadcast::Sender<CdpEvent>,
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<CdpEvent>) {
fn dispatch_message(
response: Value,
pending: &PendingCalls,
events: &broadcast::Sender<CdpEvent>,
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<CdpEvent>, 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() {

View File

@@ -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<tokio::net::TcpStream>) -> 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<u8> {
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::<Value>::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"));
}

View File

@@ -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<String>,
}
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<PageTarget> {
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::<u16>()
.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::<Vec<DevtoolsTarget>>().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<String> {
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" }),
)

View File

@@ -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<Self>, _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<ProcessHandle>) {
}
pub(crate) fn carbonyl_args(
debugging_port: u16,
profile: &str,
network_policy: &BrowserNetworkPolicy,
render_mode: RenderMode,
) -> Vec<String> {
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(),

View File

@@ -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="))
);
}

View File

@@ -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::<String>()
})
.collect::<Vec<_>>()
.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();