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
This commit is contained in:
pakrym-oai
2026-09-25 20:34:52 +00:00
committed by copyberry
parent 4193f1acfb
commit f5ffa46959
2 changed files with 34 additions and 5 deletions
@@ -125,6 +125,7 @@ async fn run_cell<H: CellHost>(
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<H: CellHost>(
_ = 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<H: CellHost>(
{
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<H: CellHost>(
}
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<H: CellHost>(
}
} => {}
_ = 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<H: CellHost>(
}
} => {
yield_timer = None;
yield_signal = None;
restore_undelivered_yield(
send_observer_event(
observer.take(),
@@ -273,6 +278,8 @@ async fn run_cell<H: CellHost>(
}, 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<H: CellHost>(
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<H: CellHost>(
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<H: CellHost>(
);
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<H: CellHost>(
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,
@@ -89,6 +89,7 @@ impl CellHost for RecordingHost {
struct CellActorHarness {
event_tx: mpsc::UnboundedSender<RuntimeEvent>,
handle: CellHandle,
yield_signal: CancellationToken,
initial_event_rx: oneshot::Receiver<Result<CellEvent, CellError>>,
task: tokio::task::JoinHandle<()>,
runtime_control_rx: std_mpsc::Receiver<RuntimeControlCommand>,
@@ -136,6 +137,7 @@ fn spawn_cell_actor_harness_with_host_and_failure_handler<H: CellHost>(
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<H: CellHost>(
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<H: CellHost>(
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);