From 4361cc6c90bd2b25db73160e0235ac261f47a0f1 Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Sun, 17 Aug 2025 14:17:45 -0700 Subject: [PATCH] fix: eliminate ServerOptions.login_timeout and have caller use tokio::time::timeout() instead --- codex-rs/login/src/server.rs | 70 ++----------------- codex-rs/login/tests/login_server_e2e.rs | 2 - .../mcp-server/src/codex_message_processor.rs | 20 ++++-- 3 files changed, 22 insertions(+), 70 deletions(-) diff --git a/codex-rs/login/src/server.rs b/codex-rs/login/src/server.rs index b85a4e0e5a..9eae4be348 100644 --- a/codex-rs/login/src/server.rs +++ b/codex-rs/login/src/server.rs @@ -3,11 +3,7 @@ use std::io::{self}; use std::path::Path; use std::path::PathBuf; use std::sync::Arc; -use std::sync::atomic::AtomicBool; -use std::sync::atomic::Ordering; -use std::sync::mpsc; use std::thread; -use std::time::Duration; use crate::AuthDotJson; use crate::get_auth_file; @@ -32,7 +28,6 @@ pub struct ServerOptions { pub port: u16, pub open_browser: bool, pub force_state: Option, - pub login_timeout: Option, } impl ServerOptions { @@ -44,7 +39,6 @@ impl ServerOptions { port: DEFAULT_PORT, open_browser: true, force_state: None, - login_timeout: None, } } } @@ -121,35 +115,16 @@ pub fn run_login_server( let shutdown_notify: Arc = shutdown_flag.unwrap_or_else(|| Arc::new(tokio::sync::Notify::new())); let shutdown_notify_clone = shutdown_notify.clone(); - let timeout_flag = Arc::new(AtomicBool::new(false)); - - // Channel used to signal completion to timeout watcher. - let (done_tx, done_rx) = mpsc::channel::<()>(); - - if let Some(timeout) = opts.login_timeout { - spawn_timeout_watcher( - done_rx, - timeout, - shutdown_notify.clone(), - timeout_flag.clone(), - server.clone(), - ); - } let (tx, mut rx) = tokio::sync::mpsc::channel::(16); let _server_handle = { let server = server.clone(); thread::spawn(move || -> io::Result<()> { - loop { - match server.recv() { - Ok(request) => tx.blocking_send(request).map_err(|e| { - eprintln!("Failed to send request to channel: {e}"); - io::Error::other("Failed to send request to channel") - })?, - Err(_e) => { - break; - } - }; + while let Ok(request) = server.recv() { + tx.blocking_send(request).map_err(|e| { + eprintln!("Failed to send request to channel: {e}"); + io::Error::other("Failed to send request to channel") + })?; } Ok(()) }) @@ -160,21 +135,11 @@ pub fn run_login_server( loop { tokio::select! { _ = shutdown_notify.notified() => { - let _ = done_tx.send(()); - if timeout_flag.load(Ordering::SeqCst) { - return Err(io::Error::other("Login timed out")); - } else { - return Err(io::Error::other("Login was not completed")); - } + return Err(io::Error::other("Login was not completed")); } maybe_req = rx.recv() => { let Some(req) = maybe_req else { - let _ = done_tx.send(()); - if timeout_flag.load(Ordering::SeqCst) { - return Err(io::Error::other("Login timed out")); - } else { - return Err(io::Error::other("Login was not completed")); - } + return Err(io::Error::other("Login was not completed")); }; let url_raw = req.url().to_string(); @@ -194,7 +159,6 @@ pub fn run_login_server( if is_login_complete { shutdown_notify.notify_waiters(); - let _ = done_tx.send(()); server_for_task.unblock(); return Ok(()); } @@ -316,26 +280,6 @@ async fn process_request( } } -/// Spawns a detached thread that waits for either a completion signal on `done_rx` -/// or the specified `timeout` to elapse. If the timeout elapses first it marks -/// the `shutdown_flag`, records `timeout_flag`, and unblocks the HTTP server so -/// that the main server loop can exit promptly. -fn spawn_timeout_watcher( - done_rx: mpsc::Receiver<()>, - timeout: Duration, - shutdown_notify: Arc, - timeout_flag: Arc, - server: Arc, -) { - thread::spawn(move || { - if done_rx.recv_timeout(timeout).is_err() { - timeout_flag.store(true, Ordering::SeqCst); - shutdown_notify.notify_waiters(); - server.unblock(); - } - }); -} - fn build_authorize_url( issuer: &str, client_id: &str, diff --git a/codex-rs/login/tests/login_server_e2e.rs b/codex-rs/login/tests/login_server_e2e.rs index 09a447d565..ef387f575e 100644 --- a/codex-rs/login/tests/login_server_e2e.rs +++ b/codex-rs/login/tests/login_server_e2e.rs @@ -100,7 +100,6 @@ async fn end_to_end_login_flow_persists_auth_json() { port: 0, open_browser: false, force_state: Some(state), - login_timeout: None, }; let server = run_login_server(opts, None).unwrap(); let login_port = server.actual_port; @@ -159,7 +158,6 @@ async fn creates_missing_codex_home_dir() { port: 0, open_browser: false, force_state: Some(state), - login_timeout: None, }; let server = run_login_server(opts, None).unwrap(); let login_port = server.actual_port; diff --git a/codex-rs/mcp-server/src/codex_message_processor.rs b/codex-rs/mcp-server/src/codex_message_processor.rs index 73a5b40c00..70341a0ab9 100644 --- a/codex-rs/mcp-server/src/codex_message_processor.rs +++ b/codex-rs/mcp-server/src/codex_message_processor.rs @@ -143,7 +143,6 @@ impl CodexMessageProcessor { let opts = LoginServerOptions { open_browser: false, - login_timeout: Some(LOGIN_CHATGPT_TIMEOUT), ..LoginServerOptions::new(config.codex_home.clone(), CLIENT_ID.to_string()) }; @@ -155,6 +154,7 @@ impl CodexMessageProcessor { let reply = match run_login_server(opts, None) { Ok(server) => { let login_id = Uuid::new_v4(); + let shutdown_handle = server.cancel_handle(); // Replace active login if present. { @@ -163,7 +163,7 @@ impl CodexMessageProcessor { existing.drop(); } *guard = Some(ActiveLogin { - shutdown_handle: server.cancel_handle(), + shutdown_handle: shutdown_handle.clone(), login_id, }); } @@ -177,9 +177,19 @@ impl CodexMessageProcessor { let outgoing_clone = self.outgoing.clone(); let active_login = self.active_login.clone(); tokio::spawn(async move { - let (success, error_msg) = match server.block_until_done().await { - Ok(()) => (true, None), - Err(err) => (false, Some(format!("Login server error: {err}"))), + let (success, error_msg) = match tokio::time::timeout( + LOGIN_CHATGPT_TIMEOUT, + server.block_until_done(), + ) + .await + { + Ok(Ok(())) => (true, None), + Ok(Err(err)) => (false, Some(format!("Login server error: {err}"))), + Err(_elapsed) => { + // Timeout: cancel server and report + shutdown_handle.cancel(); + (false, Some("Login timed out".to_string())) + } }; let notification = LoginChatGptCompleteNotification { login_id,