From f5ffa4695995bb6664fe8754ecb5ed2973f2801e Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Fri, 25 Sep 2026 20:15:06 +0000 Subject: [PATCH] Preserve queued output for observers during code-mode termination (#48207) ## Why Yield signals and queued runtime events can detach an observer during termination, before the remaining output has been drained. ## What changed Track the active yield signal separately from the observer so termination can disable yielding while keeping the observer attached. Ignore `Started`, `Pending`, and `YieldRequested` events during termination, allowing queued output to reach the observer in `CellEvent::Terminated`. ## Testing Extend the termination regression test with an immediate yield timeout, a canceled yield signal, and queued start, yield, output, and completion events. Assert that both the termination caller and the initial observer receive the queued output in the terminated event. GitOrigin-RevId: 1c08bdec06838123aaa102268da2f50b63ddea3f --- .../code-mode-runtime/src/cell_actor/mod.rs | 17 ++++++++++++-- .../code-mode-runtime/src/cell_actor/tests.rs | 22 ++++++++++++++++--- 2 files changed, 34 insertions(+), 5 deletions(-) diff --git a/codex-rs/code-mode-runtime/src/cell_actor/mod.rs b/codex-rs/code-mode-runtime/src/cell_actor/mod.rs index 7d02cfb915..4480e133da 100644 --- a/codex-rs/code-mode-runtime/src/cell_actor/mod.rs +++ b/codex-rs/code-mode-runtime/src/cell_actor/mod.rs @@ -125,6 +125,7 @@ async fn run_cell( let mut content_items = Vec::new(); let mut pending_tool_call_ids = Vec::new(); let mut pending_frontier_ready = false; + let mut yield_signal = Some(initial_observer.observation.yield_signal.clone()); let mut observer = Some(initial_observer); let mut termination = false; let mut runtime_closed = false; @@ -143,6 +144,7 @@ async fn run_cell( _ = cancellation_token.cancelled(), if !termination => { termination = true; yield_timer = None; + yield_signal = None; drop(command_rx.take()); begin_termination( &runtime_tx, @@ -193,6 +195,7 @@ async fn run_cell( { observer = None; yield_timer = None; + yield_signal = None; } if observer.is_some() || termination { let _ = response_tx.send(Err(CellError::Busy)); @@ -222,6 +225,7 @@ async fn run_cell( } continue; } + yield_signal = Some(observation.yield_signal.clone()); observer = Some(Observer { observation, response_tx }); yield_timer = observer.as_ref().and_then(observer_timer); if runtime_paused && matches!(mode, ObserveMode::YieldAfter(_)) { @@ -245,8 +249,8 @@ async fn run_cell( } } => {} _ = async { - if let Some(observer) = observer.as_ref() { - observer.observation.yield_signal.cancelled().await; + if let Some(yield_signal) = yield_signal.as_ref() { + yield_signal.cancelled().await; } else { std::future::pending::<()>().await; } @@ -254,6 +258,7 @@ async fn run_cell( } } => { yield_timer = None; + yield_signal = None; restore_undelivered_yield( send_observer_event( observer.take(), @@ -273,6 +278,8 @@ async fn run_cell( }, if !yield_deadline_elapsed => { let Some(event) = maybe_event else { runtime_closed = true; + yield_timer = None; + yield_signal = None; if termination || cancellation_token.is_cancelled() { finish_callbacks( &callback_cancellation_token, @@ -341,6 +348,9 @@ async fn run_cell( continue; }; match event { + // Keep the observer attached until termination has drained the output. + RuntimeEvent::Started | RuntimeEvent::Pending | RuntimeEvent::YieldRequested + if termination => {} RuntimeEvent::Started => { yield_timer = observer.as_ref().and_then(observer_timer); } @@ -351,6 +361,7 @@ async fn run_cell( Some(ObserveMode::PendingFrontier) ) { yield_timer = None; + yield_signal = None; pending_frontier_ready = false; match send_observer_event( observer.take(), @@ -388,6 +399,7 @@ async fn run_cell( ); if yield_observer { yield_timer = None; + yield_signal = None; restore_undelivered_yield( send_observer_event( observer.take(), @@ -431,6 +443,7 @@ async fn run_cell( RuntimeEvent::Result { stored_value_writes, error_text } => { runtime_closed = true; yield_timer = None; + yield_signal = None; if termination || cancellation_token.is_cancelled() { finish_callbacks( &callback_cancellation_token, diff --git a/codex-rs/code-mode-runtime/src/cell_actor/tests.rs b/codex-rs/code-mode-runtime/src/cell_actor/tests.rs index a2cf836b35..afb7494cc5 100644 --- a/codex-rs/code-mode-runtime/src/cell_actor/tests.rs +++ b/codex-rs/code-mode-runtime/src/cell_actor/tests.rs @@ -89,6 +89,7 @@ impl CellHost for RecordingHost { struct CellActorHarness { event_tx: mpsc::UnboundedSender, handle: CellHandle, + yield_signal: CancellationToken, initial_event_rx: oneshot::Receiver>, task: tokio::task::JoinHandle<()>, runtime_control_rx: std_mpsc::Receiver, @@ -136,6 +137,7 @@ fn spawn_cell_actor_harness_with_host_and_failure_handler( let (runtime_control_tx, runtime_control_rx) = std_mpsc::channel(); let cell_state = Arc::new(CellState::new(CancellationToken::new())); let handle = CellHandle::new(command_tx, Arc::clone(&cell_state)); + let yield_signal = CancellationToken::new(); let task = tokio::spawn(run_cell( host, CellContext { @@ -147,7 +149,7 @@ fn spawn_cell_actor_harness_with_host_and_failure_handler( event_rx, command_rx, Observer { - observation: initial_observe_mode.into(), + observation: initial_observe_mode.with_yield_signal(yield_signal.clone()), response_tx: initial_event_tx, }, task_failure_handler, @@ -156,6 +158,7 @@ fn spawn_cell_actor_harness_with_host_and_failure_handler( CellActorHarness { event_tx, handle, + yield_signal, initial_event_rx, task, runtime_control_rx, @@ -256,7 +259,17 @@ async fn yield_timer_preempts_buffered_runtime_output() { #[tokio::test] async fn queued_termination_preempts_unobserved_runtime_completion() { - let harness = spawn_cell_actor_harness(ObserveMode::YieldAfter(Duration::from_secs(60))); + let harness = spawn_cell_actor_harness(ObserveMode::YieldAfter(Duration::ZERO)); + harness.event_tx.send(RuntimeEvent::Started).unwrap(); + harness.event_tx.send(RuntimeEvent::YieldRequested).unwrap(); + harness + .event_tx + .send(RuntimeEvent::ContentItem( + FunctionCallOutputContentItem::InputText { + text: "queued output".to_string(), + }, + )) + .unwrap(); harness .event_tx .send(RuntimeEvent::Result { @@ -264,10 +277,13 @@ async fn queued_termination_preempts_unobserved_runtime_completion() { error_text: None, }) .unwrap(); + harness.yield_signal.cancel(); let termination = harness.handle.terminate(); let terminated = Ok(CellEvent::Terminated { - content_items: Vec::new(), + content_items: vec![OutputItem::Text { + text: "queued output".to_string(), + }], }); assert_eq!(termination.await, terminated.clone()); assert_eq!(harness.initial_event_rx.await.unwrap(), terminated);