From 0fa7186cf9f9db8b60e308d369381ab4f076905b Mon Sep 17 00:00:00 2001 From: Anton Panasenko Date: Sat, 6 Jun 2026 22:51:49 -0700 Subject: [PATCH 1/4] fix(app-server): avoid blocking connection cleanup --- codex-rs/app-server/src/connection_cleanup.rs | 49 +++++++++++++++ .../app-server/src/connection_rpc_gate.rs | 63 +++++++++++++++---- codex-rs/app-server/src/lib.rs | 21 +++++-- .../app-server/src/request_serialization.rs | 2 +- 4 files changed, 118 insertions(+), 17 deletions(-) create mode 100644 codex-rs/app-server/src/connection_cleanup.rs diff --git a/codex-rs/app-server/src/connection_cleanup.rs b/codex-rs/app-server/src/connection_cleanup.rs new file mode 100644 index 000000000000..55aff7597377 --- /dev/null +++ b/codex-rs/app-server/src/connection_cleanup.rs @@ -0,0 +1,49 @@ +use std::future::Future; + +use tokio::task::JoinError; +use tokio::task::JoinSet; +use tracing::warn; + +pub(crate) struct ConnectionCleanupTasks { + tasks: JoinSet<()>, +} + +impl ConnectionCleanupTasks { + pub(crate) fn new() -> Self { + Self { + tasks: JoinSet::new(), + } + } + + pub(crate) fn spawn(&mut self, future: impl Future + Send + 'static) { + self.tasks.spawn(future); + } + + pub(crate) fn is_empty(&self) -> bool { + self.tasks.is_empty() + } + + pub(crate) async fn reap_next(&mut self) { + if let Some(result) = self.tasks.join_next().await { + log_cleanup_result(result); + } + } + + pub(crate) async fn drain(&mut self) { + while let Some(result) = self.tasks.join_next().await { + log_cleanup_result(result); + } + } + + pub(crate) fn abort(&mut self) { + self.tasks.abort_all(); + } +} + +fn log_cleanup_result(result: Result<(), JoinError>) { + if let Err(err) = result + && !err.is_cancelled() + { + warn!("connection cleanup task failed: {err}"); + } +} diff --git a/codex-rs/app-server/src/connection_rpc_gate.rs b/codex-rs/app-server/src/connection_rpc_gate.rs index 12fed79b3636..209fc44655c1 100644 --- a/codex-rs/app-server/src/connection_rpc_gate.rs +++ b/codex-rs/app-server/src/connection_rpc_gate.rs @@ -1,6 +1,7 @@ use std::future::Future; +use std::sync::Mutex; +use std::sync::PoisonError; -use tokio::sync::Mutex; use tokio_util::task::TaskTracker; /// Per-connection gate for initialized RPC handler execution. @@ -27,7 +28,10 @@ impl ConnectionRpcGate { F: Future, { let token = { - let accepting = self.accepting.lock().await; + let accepting = self + .accepting + .lock() + .unwrap_or_else(PoisonError::into_inner); if !*accepting { return; } @@ -38,18 +42,26 @@ impl ConnectionRpcGate { drop(token); } + pub(crate) fn close(&self) { + let mut accepting = self + .accepting + .lock() + .unwrap_or_else(PoisonError::into_inner); + *accepting = false; + self.tasks.close(); + } + pub(crate) async fn shutdown(&self) { - { - let mut accepting = self.accepting.lock().await; - *accepting = false; - self.tasks.close(); - } + self.close(); self.tasks.wait().await; } #[cfg(test)] - async fn is_accepting(&self) -> bool { - *self.accepting.lock().await + fn is_accepting(&self) -> bool { + *self + .accepting + .lock() + .unwrap_or_else(PoisonError::into_inner) } #[cfg(test)] @@ -90,9 +102,9 @@ mod tests { } #[tokio::test] - async fn run_drops_future_without_polling_after_shutdown() { + async fn run_drops_future_without_polling_after_close() { let gate = ConnectionRpcGate::new(); - gate.shutdown().await; + gate.close(); let polled = Arc::new(AtomicBool::new(/*v*/ false)); let polled_clone = Arc::clone(&polled); @@ -102,7 +114,34 @@ mod tests { .await; assert!(!polled.load(Ordering::Acquire)); - assert!(!gate.is_accepting().await); + assert!(!gate.is_accepting()); + } + + #[tokio::test] + async fn close_returns_while_started_run_remains_active() { + let gate = Arc::new(ConnectionRpcGate::new()); + let (started_tx, started_rx) = oneshot::channel(); + let (finish_tx, finish_rx) = oneshot::channel(); + let gate_for_run = Arc::clone(&gate); + let run_task = tokio::spawn(async move { + gate_for_run + .run(async move { + started_tx.send(()).expect("receiver should be open"); + let _ = finish_rx.await; + }) + .await; + }); + + started_rx.await.expect("run should start"); + gate.close(); + assert!(!gate.is_accepting()); + assert_eq!(gate.inflight_count(), 1); + + finish_tx + .send(()) + .expect("running future should be waiting"); + run_task.await.expect("run task should complete"); + gate.shutdown().await; } #[tokio::test] diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index 0689a79f006c..7bf45856b731 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -20,6 +20,7 @@ use std::sync::atomic::AtomicBool; use crate::analytics_utils::analytics_events_client_from_config; use crate::config_manager::ConfigManager; +use crate::connection_cleanup::ConnectionCleanupTasks; use crate::message_processor::MessageProcessor; use crate::message_processor::MessageProcessorArgs; use crate::outgoing_message::ConnectionId; @@ -81,6 +82,7 @@ mod command_exec; mod config; mod config_manager; mod config_manager_service; +mod connection_cleanup; mod connection_rpc_gate; mod dynamic_tools; mod error_code; @@ -819,6 +821,7 @@ pub async fn run_main_with_transport_options( let mut thread_created_rx = processor.thread_created_receiver(); let mut running_turn_count_rx = processor.subscribe_running_assistant_turn_count(); let mut connections = HashMap::::new(); + let mut connection_cleanup_tasks = ConnectionCleanupTasks::new(); let mut remote_control_status_rx = remote_control_handle.status_receiver(); let mut remote_control_status = remote_control_status_rx.borrow().clone(); let transport_shutdown_token = transport_shutdown_token.clone(); @@ -906,14 +909,20 @@ pub async fn run_main_with_transport_options( let Some(connection_state) = connections.remove(&connection_id) else { continue; }; - if outbound_control_tx + connection_state.session.rpc_gate.close(); + let outbound_closed = outbound_control_tx .send(OutboundControlEvent::Closed { connection_id }) .await - .is_err() - { + .is_ok(); + let processor = Arc::clone(&processor); + connection_cleanup_tasks.spawn(async move { + processor + .connection_closed(connection_id, &connection_state.session) + .await; + }); + if !outbound_closed { break; } - processor.connection_closed(connection_id, &connection_state.session).await; if shutdown_when_no_connections && connections.is_empty() { break; } @@ -1010,6 +1019,7 @@ pub async fn run_main_with_transport_options( } } } + _ = connection_cleanup_tasks.reap_next(), if !connection_cleanup_tasks.is_empty() => {} changed = remote_control_status_rx.changed() => { if changed.is_err() { continue; @@ -1062,8 +1072,11 @@ pub async fn run_main_with_transport_options( .map(|connection_state| connection_state.session.rpc_gate.shutdown()), ) .await; + connection_cleanup_tasks.drain().await; processor.drain_background_tasks().await; processor.shutdown_threads().await; + } else { + connection_cleanup_tasks.abort(); } info!("processor task exited (channel closed)"); } diff --git a/codex-rs/app-server/src/request_serialization.rs b/codex-rs/app-server/src/request_serialization.rs index 0dd167b74dc2..e405521697eb 100644 --- a/codex-rs/app-server/src/request_serialization.rs +++ b/codex-rs/app-server/src/request_serialization.rs @@ -311,7 +311,7 @@ mod tests { let key = RequestSerializationQueueKey::Global("test"); let live_gate = gate(); let closed_gate = gate(); - closed_gate.shutdown().await; + closed_gate.close(); let (tx, mut rx) = mpsc::unbounded_channel(); let (blocked_tx, blocked_rx) = oneshot::channel::<()>(); From 42fb873e7ca5bc1d7661414db1299490198d52bd Mon Sep 17 00:00:00 2001 From: Anton Panasenko Date: Sat, 6 Jun 2026 22:57:08 -0700 Subject: [PATCH 2/4] refactor(app-server): await RPC gate closure --- .../app-server/src/connection_rpc_gate.rs | 32 +++++++------------ codex-rs/app-server/src/lib.rs | 2 +- .../app-server/src/request_serialization.rs | 2 +- 3 files changed, 13 insertions(+), 23 deletions(-) diff --git a/codex-rs/app-server/src/connection_rpc_gate.rs b/codex-rs/app-server/src/connection_rpc_gate.rs index 209fc44655c1..fb2aedd352bf 100644 --- a/codex-rs/app-server/src/connection_rpc_gate.rs +++ b/codex-rs/app-server/src/connection_rpc_gate.rs @@ -1,7 +1,6 @@ use std::future::Future; -use std::sync::Mutex; -use std::sync::PoisonError; +use tokio::sync::Mutex; use tokio_util::task::TaskTracker; /// Per-connection gate for initialized RPC handler execution. @@ -28,10 +27,7 @@ impl ConnectionRpcGate { F: Future, { let token = { - let accepting = self - .accepting - .lock() - .unwrap_or_else(PoisonError::into_inner); + let accepting = self.accepting.lock().await; if !*accepting { return; } @@ -42,26 +38,20 @@ impl ConnectionRpcGate { drop(token); } - pub(crate) fn close(&self) { - let mut accepting = self - .accepting - .lock() - .unwrap_or_else(PoisonError::into_inner); + pub(crate) async fn close(&self) { + let mut accepting = self.accepting.lock().await; *accepting = false; self.tasks.close(); } pub(crate) async fn shutdown(&self) { - self.close(); + self.close().await; self.tasks.wait().await; } #[cfg(test)] - fn is_accepting(&self) -> bool { - *self - .accepting - .lock() - .unwrap_or_else(PoisonError::into_inner) + async fn is_accepting(&self) -> bool { + *self.accepting.lock().await } #[cfg(test)] @@ -104,7 +94,7 @@ mod tests { #[tokio::test] async fn run_drops_future_without_polling_after_close() { let gate = ConnectionRpcGate::new(); - gate.close(); + gate.close().await; let polled = Arc::new(AtomicBool::new(/*v*/ false)); let polled_clone = Arc::clone(&polled); @@ -114,7 +104,7 @@ mod tests { .await; assert!(!polled.load(Ordering::Acquire)); - assert!(!gate.is_accepting()); + assert!(!gate.is_accepting().await); } #[tokio::test] @@ -133,8 +123,8 @@ mod tests { }); started_rx.await.expect("run should start"); - gate.close(); - assert!(!gate.is_accepting()); + gate.close().await; + assert!(!gate.is_accepting().await); assert_eq!(gate.inflight_count(), 1); finish_tx diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index 7bf45856b731..aa1a990fa8a4 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -909,7 +909,7 @@ pub async fn run_main_with_transport_options( let Some(connection_state) = connections.remove(&connection_id) else { continue; }; - connection_state.session.rpc_gate.close(); + connection_state.session.rpc_gate.close().await; let outbound_closed = outbound_control_tx .send(OutboundControlEvent::Closed { connection_id }) .await diff --git a/codex-rs/app-server/src/request_serialization.rs b/codex-rs/app-server/src/request_serialization.rs index e405521697eb..77ecfc8f56ce 100644 --- a/codex-rs/app-server/src/request_serialization.rs +++ b/codex-rs/app-server/src/request_serialization.rs @@ -311,7 +311,7 @@ mod tests { let key = RequestSerializationQueueKey::Global("test"); let live_gate = gate(); let closed_gate = gate(); - closed_gate.close(); + closed_gate.close().await; let (tx, mut rx) = mpsc::unbounded_channel(); let (blocked_tx, blocked_rx) = oneshot::channel::<()>(); From 6c10ca4be59c4270fcb9f4a03a73baa794b28da2 Mon Sep 17 00:00:00 2001 From: Anton Panasenko Date: Sat, 6 Jun 2026 23:11:42 -0700 Subject: [PATCH 3/4] codex: address PR review feedback (#26852) --- codex-rs/app-server/src/message_processor.rs | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index b7e0ff258d37..65c57bf2668f 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -91,6 +91,7 @@ use tokio_util::sync::CancellationToken; use tracing::Instrument; const EXTERNAL_AUTH_REFRESH_TIMEOUT: Duration = Duration::from_secs(10); +const CONNECTION_RPC_DRAIN_TIMEOUT: Duration = Duration::from_secs(/*secs*/ 30); #[derive(Clone)] struct ExternalAuthRefreshBridge { @@ -723,7 +724,19 @@ impl MessageProcessor { connection_id: ConnectionId, session_state: &ConnectionSessionState, ) { - session_state.rpc_gate.shutdown().await; + if timeout( + CONNECTION_RPC_DRAIN_TIMEOUT, + session_state.rpc_gate.shutdown(), + ) + .await + .is_err() + { + tracing::warn!( + ?connection_id, + timeout_seconds = CONNECTION_RPC_DRAIN_TIMEOUT.as_secs(), + "timed out waiting for connection RPCs to drain" + ); + } self.outgoing.connection_closed(connection_id).await; self.fs_processor.connection_closed(connection_id).await; self.command_exec_processor From 0c649e5492e4ee174f5b4942c9a26ee3e561a322 Mon Sep 17 00:00:00 2001 From: Anton Panasenko Date: Sat, 6 Jun 2026 23:18:24 -0700 Subject: [PATCH 4/4] codex: address PR review feedback (#26852) --- codex-rs/app-server/src/connection_cleanup.rs | 8 ++++---- codex-rs/app-server/src/lib.rs | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/codex-rs/app-server/src/connection_cleanup.rs b/codex-rs/app-server/src/connection_cleanup.rs index 55aff7597377..529020fe3ff9 100644 --- a/codex-rs/app-server/src/connection_cleanup.rs +++ b/codex-rs/app-server/src/connection_cleanup.rs @@ -1,4 +1,5 @@ use std::future::Future; +use std::future::pending; use tokio::task::JoinError; use tokio::task::JoinSet; @@ -19,11 +20,10 @@ impl ConnectionCleanupTasks { self.tasks.spawn(future); } - pub(crate) fn is_empty(&self) -> bool { - self.tasks.is_empty() - } - pub(crate) async fn reap_next(&mut self) { + if self.tasks.is_empty() { + pending::<()>().await; + } if let Some(result) = self.tasks.join_next().await { log_cleanup_result(result); } diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index aa1a990fa8a4..b6f0b8c8597d 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -1019,7 +1019,7 @@ pub async fn run_main_with_transport_options( } } } - _ = connection_cleanup_tasks.reap_next(), if !connection_cleanup_tasks.is_empty() => {} + _ = connection_cleanup_tasks.reap_next() => {} changed = remote_control_status_rx.changed() => { if changed.is_err() { continue;