Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Add opt-in code-mode yielding on new user input (#48135)
## Why

User input queued during a long-running code-mode `exec` or `wait` can remain pending until the call returns. Allow these calls to yield early so Codex can receive the new input while the cell continues running.

## What changed

- Add the `instant_interrupt` feature flag, disabled by default.
- Watch queued user input for each sampling request and pass a shared preemption signal to code-mode `exec` and `wait` calls, including calls received later in the same response.
- Yield running cell IDs without stopping the cells, allowing subsequent `wait` calls to collect their results.

## Testing

Add integration coverage for enabled and disabled behavior, repeated steering during `exec` and `wait`, input deferred during compaction, and later calls in the same response. Verify that direct tool results and queued user messages are preserved, and add a scenario snapshot showing continuation after a cell yields.

GitOrigin-RevId: e33466a980d98e113bfb6c6a45c59de0cc7b3c1d
  • Loading branch information
pakrym-oai authored and copyberry committed Sep 25, 2026
commit 1bf73324cadc72a53ed467edc7d3fd2b145a6166
1 change: 1 addition & 0 deletions codex-rs/code-mode-host/tests/stdio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,7 @@ async fn interrupt_yields_observations_without_stopping_the_cell() {
description: String::new(),
kind: CodeModeToolKind::Function,
input_schema: None,
input_schema_max_bytes: None,
output_schema: None,
}];
let signal = CancellationToken::new();
Expand Down
6 changes: 6 additions & 0 deletions codex-rs/core/config.schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -894,6 +894,9 @@
"in_app_updates": {
"type": "boolean"
},
"instant_interrupt": {
"type": "boolean"
},
"item_ids": {
"type": "boolean"
},
Expand Down Expand Up @@ -6926,6 +6929,9 @@
"in_app_updates": {
"type": "boolean"
},
"instant_interrupt": {
"type": "boolean"
},
"item_ids": {
"type": "boolean"
},
Expand Down
31 changes: 31 additions & 0 deletions codex-rs/core/src/session/input_queue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ use std::collections::VecDeque;
use std::sync::Arc;
use tokio::sync::Mutex;
use tokio::sync::watch;
use tokio_util::sync::CancellationToken;
use tokio_util::task::AbortOnDropHandle;

static PENDING_MAILBOX_MESSAGES: Gauge = Gauge::new("core.mailbox.pending");

Expand Down Expand Up @@ -215,6 +217,29 @@ impl InputQueue {
})
}

/// Signal once a user message is queued for this sampling request.
pub(crate) async fn watch_user_input(
&self,
active_turn: &Mutex<Option<ActiveTurn>>,
sub_id: &str,
interrupt: CancellationToken,
) -> Option<AbortOnDropHandle<()>> {
let turn_state = self.turn_state_for_sub_id(active_turn, sub_id).await?;
// Subscribe before inspecting the queue so an arrival cannot be missed.
let mut activity = self.activity_tx.subscribe();
Some(AbortOnDropHandle::new(tokio::spawn(async move {
loop {
if turn_state.lock().await.pending_input.has_user_input() {
interrupt.cancel();
return;
}
if activity.changed().await.is_err() {
return;
}
}
})))
}

/// Clear any pending waiters and input buffered for the current turn.
pub(crate) async fn clear_pending(&self, active_turn: &ActiveTurn) {
let mut turn_state = active_turn.turn_state.lock().await;
Expand Down Expand Up @@ -363,6 +388,12 @@ impl InputQueue {
}

impl TurnInputQueue {
fn has_user_input(&self) -> bool {
self.items
.iter()
.any(|input| matches!(input, TurnInput::UserInput { .. }))
}

pub(crate) fn is_empty(&self) -> bool {
self.items.is_empty()
}
Expand Down
5 changes: 5 additions & 0 deletions codex-rs/core/src/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3942,6 +3942,11 @@ impl Session {
) = prepared_tools??;
turn_context.extension_data.insert(selected_plugins);
Ok(Arc::new(StepContext {
preempt: turn_context
.config
.features
.enabled(Feature::InstantInterrupt)
.then(CancellationToken::new),
realtime: self.conversation.snapshot().await,
settings,
token_budget,
Expand Down
3 changes: 3 additions & 0 deletions codex-rs/core/src/session/step_context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,13 @@ use codex_mcp::McpBinding;
use codex_otel::SessionTelemetry;
use codex_protocol::items::ModelInvocationContext;
use codex_protocol::protocol::TurnContextItem;
use tokio_util::sync::CancellationToken;

/// Request-scoped state that may change between model sampling requests.
pub(crate) struct StepContext {
pub(crate) turn: Arc<TurnContext>,
/// Yields this request's code-mode observations when a user message arrives.
pub(crate) preempt: Option<CancellationToken>,
/// Realtime call activity and instructions captured for this sampling request.
pub(crate) realtime: RealtimeConversationSnapshot,
/// One immutable settings version captured before request preparation.
Expand Down
5 changes: 5 additions & 0 deletions codex-rs/core/src/session/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,11 @@ impl StepContext {
});
settings.service_tier = turn.config.service_tier.clone();
Arc::new(Self {
preempt: turn
.config
.features
.enabled(Feature::InstantInterrupt)
.then(tokio_util::sync::CancellationToken::new),
token_budget: token_budget::resolve_token_budget(
turn.configured_token_budget.as_ref(),
turn.use_model_token_budget_defaults,
Expand Down
7 changes: 7 additions & 0 deletions codex-rs/core/src/session/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1588,6 +1588,13 @@ async fn run_sampling_request(
cancellation_token: CancellationToken,
) -> CodexResult<(SamplingRequestResult, Vec<ResponseItem>)> {
let turn_context = Arc::clone(&step_context.turn);
let _input_watch = if let Some(preempt) = &step_context.preempt {
sess.input_queue
.watch_user_input(&sess.active_turn, &turn_context.sub_id, preempt.clone())
.await
} else {
None
};
let base_instructions = sess.get_prompt_base_instructions().await;

let tool_runtime = ToolCallRuntime::new(
Expand Down
6 changes: 4 additions & 2 deletions codex-rs/core/src/tools/code_mode/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,21 +133,23 @@ impl CodeModeService {
request
.yield_time_ms
.get_or_insert(self.default_exec_yield_time_ms);
let preempt = step_context.preempt.clone();
let delegate = Arc::new(CodeModeCellDelegate {
broker: Arc::clone(&self.dispatch_broker),
step_context,
});
self.session()
.await?
.execute(request, delegate, /*preempt*/ None)
.execute(request, delegate, preempt)
.await
}

pub(crate) async fn wait(
&self,
request: codex_code_mode::WaitRequest,
preempt: Option<CancellationToken>,
) -> Result<codex_code_mode::WaitOutcome, String> {
self.session().await?.wait(request, /*preempt*/ None).await
self.session().await?.wait(request, preempt).await
}

pub(crate) async fn terminate(
Expand Down
11 changes: 7 additions & 4 deletions codex-rs/core/src/tools/code_mode/wait_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -134,10 +134,13 @@ impl CodeModeWaitHandler {
exec.session
.services
.code_mode_service
.wait(codex_code_mode::WaitRequest {
cell_id,
yield_time_ms: args.yield_time_ms,
})
.wait(
codex_code_mode::WaitRequest {
cell_id,
yield_time_ms: args.yield_time_ms,
},
step_context.preempt.clone(),
)
.await
}
.map_err(|error| {
Expand Down
Loading
Loading