Wire turn item contributors into stream output (#22494)

## Summary
- run registered TurnItemContributor hooks for parsed stream output
items
- plumb the active turn extension store into stream item handling
- preserve existing memory citation parsing as fallback after
contributors run

## Tests
- cargo test -p codex-core stream_events_utils -- --nocapture
- just fmt
- just fix -p codex-core
- git diff --check
This commit is contained in:
jif-oai
2026-05-14 14:48:17 +02:00
committed by GitHub
parent 6d65686313
commit 17cd321c32
7 changed files with 550 additions and 60 deletions
+19 -4
View File
@@ -8,6 +8,7 @@ use std::sync::Arc;
use std::time::Duration;
use std::time::Instant;
use codex_extension_api::ExtensionData;
use futures::future::BoxFuture;
use tokio::select;
use tokio::sync::Notify;
@@ -153,17 +154,25 @@ fn bool_tag(value: bool) -> &'static str {
#[derive(Clone)]
pub(crate) struct SessionTaskContext {
session: Arc<Session>,
turn_extension_data: Arc<ExtensionData>,
}
impl SessionTaskContext {
pub(crate) fn new(session: Arc<Session>) -> Self {
Self { session }
pub(crate) fn new(session: Arc<Session>, turn_extension_data: Arc<ExtensionData>) -> Self {
Self {
session,
turn_extension_data,
}
}
pub(crate) fn clone_session(&self) -> Arc<Session> {
Arc::clone(&self.session)
}
pub(crate) fn turn_extension_data(&self) -> Arc<ExtensionData> {
Arc::clone(&self.turn_extension_data)
}
pub(crate) fn auth_manager(&self) -> Arc<AuthManager> {
Arc::clone(&self.session.services.auth_manager)
}
@@ -362,7 +371,10 @@ impl Session {
let turn = active.get_or_insert_with(ActiveTurn::default);
debug_assert!(turn.tasks.is_empty());
let done_clone = Arc::clone(&done);
let session_ctx = Arc::new(SessionTaskContext::new(Arc::clone(self)));
let session_ctx = Arc::new(SessionTaskContext::new(
Arc::clone(self),
Arc::clone(&turn_extension_data),
));
let ctx = Arc::clone(&turn_context);
let task_for_run = Arc::clone(&task);
let task_cancellation_token = cancellation_token.child_token();
@@ -829,7 +841,10 @@ impl Session {
task.handle.abort();
let session_ctx = Arc::new(SessionTaskContext::new(Arc::clone(self)));
let session_ctx = Arc::new(SessionTaskContext::new(
Arc::clone(self),
Arc::clone(&task.turn_extension_data),
));
session_task
.abort(session_ctx, Arc::clone(&task.turn_context))
.await;
+2
View File
@@ -45,6 +45,7 @@ impl SessionTask for RegularTask {
cancellation_token: CancellationToken,
) -> Option<String> {
let sess = session.clone_session();
let turn_extension_data = session.turn_extension_data();
let run_turn_span = trace_span!("run_turn");
// Regular turns emit `TurnStarted` inline so first-turn lifecycle does
// not wait on startup prewarm resolution.
@@ -72,6 +73,7 @@ impl SessionTask for RegularTask {
let last_agent_message = run_turn(
Arc::clone(&sess),
Arc::clone(&ctx),
Arc::clone(&turn_extension_data),
next_input,
prewarmed_client_session.take(),
cancellation_token.child_token(),