use std::io::Error;
use std::net::Ipv4Addr;
use std::sync::Arc;
use http::HeaderValue;
use http::header::CACHE_CONTROL;
use http::header::CONTENT_TYPE;
use serde_json::Value;
use serde_json::json;
use tokio::io::AsyncReadExt;
use tokio::io::AsyncWriteExt;
use tokio::net::TcpListener;
use tokio::net::TcpStream;
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use topcoat::Result;
use topcoat::context::Cx;
use topcoat::context::app_context;
use topcoat::router::Body;
use topcoat::router::IntoResponse;
use topcoat::router::Response;
use topcoat::router::Router;
use topcoat::router::RouterBuilderDiscoverExt;
use topcoat::router::headers;
use topcoat::router::not_found;
use topcoat::router::page;
use topcoat::router::path_param;
use topcoat::router::route;
use topcoat::router::to_bytes;
use topcoat::router::uri;
use crate::WebOptions;
use crate::bridge::AppServerBridge;
use crate::components::DocumentData;
use crate::components::document;
use crate::components::thread_title;
use crate::components::transcript_fragment;
use crate::network::authority;
use crate::network::bind_listeners;
use crate::server_support::canonicalize_cwd;
use crate::server_support::chronological_turns;
use crate::server_support::generate_secret;
use crate::server_support::mcp_callback_url;
use crate::server_support::query_value;
use crate::server_support::thread_with_initial_turns;
use crate::stream;
const MAX_RPC_BODY_BYTES: usize = 64 * 1024 * 1024;
const APP_JS: &str = include_str!("../assets/app.js");
const APP_CSS: &str = include_str!("../assets/app.css");
struct WebState {
secret: String,
instance_id: String,
cwd: String,
bridge: Arc<AppServerBridge>,
events: Arc<stream::LiveEventHub>,
shutdown: CancellationToken,
mcp_callback_port: Option<u16>,
}
#[topcoat::router::path_param]
struct Secret(str);
#[topcoat::router::path_param]
struct ThreadId(str);
#[topcoat::router::path_param]
struct CallbackId(str);
#[topcoat::router::path_param]
struct SessionId(str);
pub async fn run(options: WebOptions) -> std::io::Result<()> {
let cwd = canonicalize_cwd(options.cwd)?;
let secret = generate_secret();
let (listeners, _port) = bind_listeners(options.port).await?;
let addresses = listeners
.iter()
.filter_map(|listener| listener.local_addr().ok())
.collect::<Vec<_>>();
let mut config_overrides = options.config_overrides;
let mcp_callback_port = if let Some(callback_url) = mcp_callback_url(&addresses, &secret) {
let callback_listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await?;
let callback_port = callback_listener.local_addr()?.port();
drop(callback_listener);
config_overrides.push(format!("mcp_oauth_callback_port={callback_port}"));
config_overrides.push(format!("mcp_oauth_callback_url={callback_url:?}"));
Some(callback_port)
} else {
None
};
let bridge = AppServerBridge::start(
&options.codex_executable,
&cwd,
&config_overrides,
options.strict_config,
)
.await?;
let shutdown_token = CancellationToken::new();
let event_hub = stream::LiveEventHub::start(Arc::clone(&bridge)).await;
let state = Arc::new(WebState {
secret: secret.clone(),
instance_id: generate_secret()[..16].to_string(),
cwd: cwd.to_string_lossy().into_owned(),
bridge: Arc::clone(&bridge),
events: event_hub,
shutdown: shutdown_token.clone(),
mcp_callback_port,
});
let router = Router::builder().discover().app_context(state).build();
let service = topcoat::router::RouterService::new(router);
let mut servers = JoinSet::new();
for listener in listeners {
let service = service.clone();
let shutdown_token = shutdown_token.clone();
servers.spawn(async move {
topcoat::serve_until(listener, service, async move {
shutdown_token.cancelled().await;
})
.await
});
}
let urls = addresses
.iter()
.map(|address| format!("http://{}/s/{secret}", authority(*address)))
.collect::<Vec<_>>();
println!("\nCodex Web\n");
for url in &urls {
println!(" {url}");
}
println!("\nAnyone with a link can control this Codex process.\n");
if options.open_browser
&& let Some(url) = urls.iter().find(|url| url.starts_with("http://127.0.0.1:"))
&& let Err(error) = webbrowser::open(url)
{
tracing::warn!(%error, "failed to open Codex Web in a browser");
}
tokio::select! {
signal = tokio::signal::ctrl_c() => signal?,
result = servers.join_next() => {
match result {
Some(Ok(Ok(()))) | None => {}
Some(Ok(Err(error))) => return Err(error),
Some(Err(error)) => return Err(Error::other(error)),
}
}
}
shutdown_token.cancel();
servers.abort_all();
while let Some(result) = servers.join_next().await {
match result {
Ok(Ok(())) => {}
Ok(Err(error)) => tracing::warn!(%error, "Codex Web listener stopped"),
Err(error) if error.is_cancelled() => {}
Err(error) => tracing::warn!(%error, "Codex Web listener task failed"),
}
}
bridge.shutdown().await;
Ok(())
}
#[page("/s/{secret}")]
async fn home(cx: &Cx) -> Result {
let state = authorized_state(cx)?;
render_page(cx, state, None).await
}
#[page("/s/{secret}/thread/{thread_id}")]
async fn thread_page(cx: &Cx) -> Result {
let state = authorized_state(cx)?;
render_page(cx, state, Some(path_param::<ThreadId>(cx))).await
}
#[route(GET "/s/{secret}/assets/app.js")]
async fn app_js(cx: &Cx) -> Result<Response> {
authorized_state(cx)?;
static_asset(cx, "text/javascript; charset=utf-8", APP_JS)
}
#[route(GET "/s/{secret}/assets/app.css")]
async fn app_css(cx: &Cx) -> Result<Response> {
authorized_state(cx)?;
static_asset(cx, "text/css; charset=utf-8", APP_CSS)
}
#[route(GET "/s/{secret}/events")]
async fn event_stream(cx: &Cx) -> Result<Response> {
let state = authorized_state(cx)?;
let last_event_id = headers(cx)
.get("last-event-id")
.and_then(|value| value.to_str().ok())
.and_then(|value| value.parse().ok())
.or_else(|| {
query_value(uri(cx).query().unwrap_or_default(), "after")
.and_then(|value| value.parse().ok())
});
Ok(stream::response(Arc::clone(&state.events), last_event_id).await)
}
#[route(GET "/s/{secret}/oauth/callback/{callback_id}")]
async fn mcp_oauth_callback(cx: &Cx) -> Result<Response> {
let state = authorized_state(cx)?;
let callback_port = state
.mcp_callback_port
.ok_or_else(|| Error::other("remote MCP OAuth callback is not configured"))?;
let callback_id = path_param::<CallbackId>(cx);
let query = uri(cx).query().unwrap_or_default();
let callback_path = format!("/s/{}/oauth/callback/{callback_id}?{query}", state.secret);
let mut stream = TcpStream::connect((Ipv4Addr::LOCALHOST, callback_port))
.await
.map_err(Error::other)?;
stream
.write_all(
format!(
"GET {callback_path} HTTP/1.1\r\nHost: 127.0.0.1:{callback_port}\r\nConnection: close\r\n\r\n"
)
.as_bytes(),
)
.await
.map_err(Error::other)?;
let mut encoded = Vec::new();
stream
.take(64 * 1024)
.read_to_end(&mut encoded)
.await
.map_err(Error::other)?;
let forwarded = String::from_utf8_lossy(&encoded);
let body = forwarded
.split_once("\r\n\r\n")
.map(|(_, body)| body)
.unwrap_or("MCP OAuth callback returned an invalid response")
.to_string();
let mut response = body.into_response(cx)?;
response.headers_mut().insert(
CONTENT_TYPE,
HeaderValue::from_static("text/html; charset=utf-8"),
);
response
.headers_mut()
.insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
Ok(response)
}
#[route(GET "/s/{secret}/api/session/{session_id}")]
async fn session_fragment(cx: &Cx) -> Result<Response> {
let state = authorized_state(cx)?;
let session_id = path_param::<SessionId>(cx);
let base_seq = state.events.watermark();
let active_request = async {
if session_id == "new" {
None
} else {
let paginated = state
.bridge
.request(
"thread/resume",
json!({
"threadId": session_id,
"excludeTurns": true,
"initialTurnsPage": {
"limit": 30,
"sortDirection": "desc",
"itemsView": "full"
}
}),
)
.await;
match paginated {
Ok(response) => Some(response),
Err(_) => state
.bridge
.request("thread/resume", json!({ "threadId": session_id }))
.await
.ok(),
}
}
};
let (active, approvals) =
tokio::join!(active_request, state.bridge.outstanding_server_requests(),);
let active_thread = active.as_ref().and_then(|response| response.get("thread"));
let hydrated_thread = active.as_ref().and_then(thread_with_initial_turns);
let display_thread = hydrated_thread.as_ref().or(active_thread);
let fragment = transcript_fragment(cx, display_thread, &approvals).await?;
let active_turn_id = display_thread
.and_then(|active_thread| active_thread.get("turns"))
.and_then(Value::as_array)
.and_then(|turns| {
turns
.iter()
.rev()
.find(|turn| turn.get("status").and_then(Value::as_str) == Some("inProgress"))
})
.and_then(|turn| turn.get("id"))
.and_then(Value::as_str)
.unwrap_or_default();
let payload = json!({
"html": fragment.render(cx),
"baseSeq": base_seq,
"instanceId": state.instance_id,
"nextCursor": active.as_ref().and_then(|response| response.pointer("/initialTurnsPage/nextCursor")),
"threadId": display_thread
.and_then(|active_thread| active_thread.get("id"))
.and_then(Value::as_str)
.unwrap_or_default(),
"activeTurnId": active_turn_id,
"title": display_thread.map(thread_title).unwrap_or("New task"),
"cwd": display_thread
.and_then(|active_thread| active_thread.get("cwd"))
.and_then(Value::as_str)
.unwrap_or_default(),
"model": active
.as_ref()
.and_then(|response| response.get("model"))
.and_then(Value::as_str),
"reasoningEffort": active
.as_ref()
.and_then(|response| response.get("reasoningEffort"))
.and_then(Value::as_str),
"permissionMode": permission_mode(active.as_ref()),
});
let mut response = serde_json::to_vec(&payload)?.into_response(cx)?;
response
.headers_mut()
.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
response
.headers_mut()
.insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
Ok(response)
}
#[route(GET "/s/{secret}/api/session/{session_id}/turns")]
async fn earlier_turns(cx: &Cx) -> Result<Response> {
let state = authorized_state(cx)?;
let session_id = path_param::<SessionId>(cx);
let cursor = query_value(uri(cx).query().unwrap_or_default(), "cursor");
let page = state
.bridge
.request(
"thread/turns/list",
json!({
"threadId": session_id,
"cursor": cursor,
"limit": 30,
"sortDirection": "desc",
"itemsView": "full"
}),
)
.await
.map_err(Error::other)?;
let thread_value = json!({ "turns": chronological_turns(page.get("data")) });
let fragment = transcript_fragment(cx, Some(&thread_value), &[]).await?;
json_response(
cx,
json!({
"html": fragment.render(cx),
"nextCursor": page.get("nextCursor").cloned().unwrap_or(Value::Null)
}),
)
}
#[route(POST "/s/{secret}/api/shutdown")]
async fn shutdown_process(cx: &Cx) -> Result<Response> {
let state = authorized_state(cx)?;
state.shutdown.cancel();
json_response(cx, json!({ "stopping": true }))
}
#[route(POST "/s/{secret}/api/rpc")]
async fn rpc(cx: &Cx, body: Body) -> Result<Response> {
let state = authorized_state(cx)?;
let bytes = to_bytes(body, MAX_RPC_BODY_BYTES)
.await
.map_err(Error::other)?;
let message: Value = serde_json::from_slice(&bytes)?;
let envelope = if let Some(method) = message.get("method").and_then(Value::as_str) {
state
.bridge
.request_envelope(
method,
message.get("params").cloned().unwrap_or_else(|| json!({})),
)
.await
.map_err(Error::other)?
} else {
state.bridge.respond(message).await.map_err(Error::other)?;
json!({ "result": {} })
};
let mut response = serde_json::to_vec(&envelope)?.into_response(cx)?;
response
.headers_mut()
.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
response
.headers_mut()
.insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
Ok(response)
}
async fn render_page(cx: &Cx, state: &WebState, thread_id: Option<&str>) -> Result {
let base_seq = state.events.watermark();
let threads_request = state.bridge.request(
"thread/list",
json!({ "limit": 100, "sortKey": "recency_at", "sortDirection": "desc" }),
);
let models_request = state.bridge.request("model/list", json!({ "limit": 100 }));
let collaboration_modes_request = state.bridge.request("collaborationMode/list", json!({}));
let mcp_servers_request = state.bridge.request(
"mcpServerStatus/list",
json!({ "limit": 100, "detail": "toolsAndAuthOnly" }),
);
let active_request = async {
match thread_id {
Some(thread_id) => state
.bridge
.request(
"thread/resume",
json!({
"threadId": thread_id,
"excludeTurns": true,
"initialTurnsPage": {
"limit": 30,
"sortDirection": "desc",
"itemsView": "full"
}
}),
)
.await
.ok(),
None => None,
}
};
let approvals_request = state.bridge.outstanding_server_requests();
let (threads, models, collaboration_modes, mcp_servers, active, approvals) = tokio::join!(
threads_request,
models_request,
collaboration_modes_request,
mcp_servers_request,
active_request,
approvals_request,
);
let threads = threads.unwrap_or_else(|_| json!({ "data": [] }));
let models = models.unwrap_or_else(|_| json!({ "data": [] }));
let collaboration_modes = collaboration_modes.unwrap_or_else(|_| json!({ "data": [] }));
let mcp_servers = mcp_servers.unwrap_or_else(|_| json!({ "data": [] }));
let hydrated_thread = active.as_ref().and_then(thread_with_initial_turns);
let active_thread = hydrated_thread
.as_ref()
.or_else(|| active.as_ref().and_then(|response| response.get("thread")));
document(
cx,
DocumentData {
secret: &state.secret,
instance_id: &state.instance_id,
base_seq,
threads: threads
.get("data")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default(),
active_thread,
active_model: active
.as_ref()
.and_then(|response| response.get("model"))
.and_then(Value::as_str)
.unwrap_or_default(),
active_effort: active
.as_ref()
.and_then(|response| response.get("reasoningEffort"))
.and_then(Value::as_str)
.unwrap_or_default(),
active_permission_mode: permission_mode(active.as_ref()),
models: models
.get("data")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default(),
collaboration_modes: collaboration_modes
.get("data")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default(),
approvals: &approvals,
workspace_cwd: &state.cwd,
mcp_servers: mcp_servers
.get("data")
.and_then(Value::as_array)
.map(Vec::as_slice)
.unwrap_or_default(),
initial_next_cursor: active
.as_ref()
.and_then(|response| response.pointer("/initialTurnsPage/nextCursor"))
.and_then(Value::as_str),
},
)
.await
}
fn permission_mode(response: Option<&Value>) -> &'static str {
let profile_id = response
.and_then(|response| response.pointer("/activePermissionProfile/id"))
.and_then(Value::as_str);
match profile_id {
Some(":danger-full-access") => "full-access",
Some(":read-only") => "read-only",
Some(":workspace") => "workspace",
_ => {
let approval = response
.and_then(|response| response.get("approvalPolicy"))
.and_then(Value::as_str);
let sandbox = response
.and_then(|response| response.pointer("/sandbox/type"))
.and_then(Value::as_str);
match (approval, sandbox) {
(Some("never"), Some("danger-full-access")) => "full-access",
(_, Some("read-only")) => "read-only",
_ => "workspace",
}
}
}
}
fn authorized_state(cx: &Cx) -> Result<&WebState> {
let state: &Arc<WebState> = app_context(cx);
if path_param::<Secret>(cx) != state.secret {
return Err(not_found().into());
}
Ok(state)
}
fn static_asset(cx: &Cx, content_type: &'static str, body: &'static str) -> Result<Response> {
let mut response = body.into_response(cx)?;
response
.headers_mut()
.insert(CONTENT_TYPE, HeaderValue::from_static(content_type));
response
.headers_mut()
.insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
Ok(response)
}
fn json_response(cx: &Cx, payload: Value) -> Result<Response> {
let mut response = serde_json::to_vec(&payload)?.into_response(cx)?;
response
.headers_mut()
.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
response
.headers_mut()
.insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
Ok(response)
}
#[cfg(test)]
#[path = "server_tests.rs"]
mod tests;