mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
feat: move extension scope ids into ExtensionData (#22490)
## Summary - add a scoped level_id to ExtensionData and expose it through level_id() - remove thread_id/turn_id parameters from extension contributor inputs where the scoped ExtensionData already carries that identity - move turn-scoped extension data onto TurnContext so token usage and lifecycle contributors can share the same turn store ## Testing - cargo check -p codex-extension-api -p codex-core --tests - cargo test -p codex-extension-api - cargo test -p codex-guardian - cargo test -p codex-core --lib record_token_usage_info_notifies_extension_contributors - cargo test -p codex-core --lib submission_loop_channel_close_emits_thread_stop_lifecycle - cargo test -p codex-core --lib submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle - just fix -p codex-extension-api - just fix -p codex-guardian - just fix -p codex-core - just fmt ## Note - Attempted cargo test -p codex-core; it aborted in agent::control::tests::spawn_agent_fork_last_n_turns_keeps_only_recent_turns with the existing stack overflow before the full suite completed.
This commit is contained in:
@@ -635,7 +635,6 @@ async fn shutdown_session_runtime(sess: &Arc<Session>) {
|
||||
fn emit_thread_stop_lifecycle(sess: &Session) {
|
||||
for contributor in sess.services.extensions.thread_lifecycle_contributors() {
|
||||
contributor.on_thread_stop(codex_extension_api::ThreadStopInput {
|
||||
thread_id: sess.conversation_id,
|
||||
session_store: &sess.services.session_extension_data,
|
||||
thread_store: &sess.services.thread_extension_data,
|
||||
});
|
||||
|
||||
@@ -2869,8 +2869,7 @@ impl Session {
|
||||
contributor.on_token_usage(
|
||||
&self.services.session_extension_data,
|
||||
&self.services.thread_extension_data,
|
||||
self.conversation_id,
|
||||
&turn_context.sub_id,
|
||||
turn_context.extension_data.as_ref(),
|
||||
token_info,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -113,7 +113,7 @@ pub(super) async fn spawn_review_thread(
|
||||
));
|
||||
|
||||
let review_turn_context = TurnContext {
|
||||
sub_id: review_turn_id,
|
||||
sub_id: review_turn_id.clone(),
|
||||
trace_id: current_span_trace_id(),
|
||||
realtime_active: parent_turn_context.realtime_active,
|
||||
config: per_turn_config,
|
||||
@@ -149,6 +149,7 @@ pub(super) async fn spawn_review_thread(
|
||||
dynamic_tools: parent_turn_context.dynamic_tools.clone(),
|
||||
truncation_policy: model_info.truncation_policy.into(),
|
||||
turn_metadata_state,
|
||||
extension_data: Arc::new(codex_extension_api::ExtensionData::new(review_turn_id)),
|
||||
turn_skills: TurnSkillsContext::new(parent_turn_context.turn_skills.outcome.clone()),
|
||||
turn_timing_state: Arc::new(TurnTimingState::default()),
|
||||
server_model_warning_emitted: AtomicBool::new(false),
|
||||
|
||||
@@ -811,11 +811,12 @@ impl Session {
|
||||
SessionId::from(thread_id)
|
||||
};
|
||||
let agent_control = agent_control.with_session_id(session_id);
|
||||
let session_extension_data = codex_extension_api::ExtensionData::new();
|
||||
let thread_extension_data = codex_extension_api::ExtensionData::new();
|
||||
let session_extension_data =
|
||||
codex_extension_api::ExtensionData::new(session_id.to_string());
|
||||
let thread_extension_data =
|
||||
codex_extension_api::ExtensionData::new(thread_id.to_string());
|
||||
for contributor in extensions.thread_lifecycle_contributors() {
|
||||
contributor.on_thread_start(codex_extension_api::ThreadStartInput {
|
||||
thread_id,
|
||||
config: config.as_ref(),
|
||||
session_store: &session_extension_data,
|
||||
thread_store: &thread_extension_data,
|
||||
|
||||
@@ -1807,8 +1807,9 @@ async fn record_token_usage_info_notifies_extension_contributors() {
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct RecordedTokenUsage {
|
||||
thread_id: ThreadId,
|
||||
turn_id: String,
|
||||
session_level_id: String,
|
||||
thread_level_id: String,
|
||||
turn_level_id: String,
|
||||
token_usage: TokenUsageInfo,
|
||||
saw_session_store: bool,
|
||||
saw_thread_store: bool,
|
||||
@@ -1823,16 +1824,16 @@ async fn record_token_usage_info_notifies_extension_contributors() {
|
||||
&self,
|
||||
session_store: &codex_extension_api::ExtensionData,
|
||||
thread_store: &codex_extension_api::ExtensionData,
|
||||
thread_id: ThreadId,
|
||||
turn_id: &str,
|
||||
turn_store: &codex_extension_api::ExtensionData,
|
||||
token_usage: &TokenUsageInfo,
|
||||
) {
|
||||
self.records
|
||||
.lock()
|
||||
.expect("token usage records lock")
|
||||
.push(RecordedTokenUsage {
|
||||
thread_id,
|
||||
turn_id: turn_id.to_string(),
|
||||
session_level_id: session_store.level_id().to_string(),
|
||||
thread_level_id: thread_store.level_id().to_string(),
|
||||
turn_level_id: turn_store.level_id().to_string(),
|
||||
token_usage: token_usage.clone(),
|
||||
saw_session_store: session_store.get::<SessionTokenUsageMarker>().is_some(),
|
||||
saw_thread_store: thread_store.get::<ThreadTokenUsageMarker>().is_some(),
|
||||
@@ -1882,8 +1883,9 @@ async fn record_token_usage_info_notifies_extension_contributors() {
|
||||
expected_total_usage.add_assign(&second_usage);
|
||||
let expected = vec![
|
||||
RecordedTokenUsage {
|
||||
thread_id: session.conversation_id,
|
||||
turn_id: turn_context.sub_id.clone(),
|
||||
session_level_id: session.session_id().to_string(),
|
||||
thread_level_id: session.conversation_id.to_string(),
|
||||
turn_level_id: turn_context.sub_id.clone(),
|
||||
token_usage: TokenUsageInfo {
|
||||
total_token_usage: first_usage.clone(),
|
||||
last_token_usage: first_usage,
|
||||
@@ -1893,8 +1895,9 @@ async fn record_token_usage_info_notifies_extension_contributors() {
|
||||
saw_thread_store: true,
|
||||
},
|
||||
RecordedTokenUsage {
|
||||
thread_id: session.conversation_id,
|
||||
turn_id: turn_context.sub_id.clone(),
|
||||
session_level_id: session.session_id().to_string(),
|
||||
thread_level_id: session.conversation_id.to_string(),
|
||||
turn_level_id: turn_context.sub_id.clone(),
|
||||
token_usage: TokenUsageInfo {
|
||||
total_token_usage: expected_total_usage,
|
||||
last_token_usage: second_usage,
|
||||
@@ -3981,8 +3984,10 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) {
|
||||
plugins_manager,
|
||||
mcp_manager,
|
||||
extensions: Arc::new(codex_extension_api::ExtensionRegistryBuilder::new().build()),
|
||||
session_extension_data: codex_extension_api::ExtensionData::new(),
|
||||
thread_extension_data: codex_extension_api::ExtensionData::new(),
|
||||
session_extension_data: codex_extension_api::ExtensionData::new(
|
||||
agent_control.session_id().to_string(),
|
||||
),
|
||||
thread_extension_data: codex_extension_api::ExtensionData::new(thread_id.to_string()),
|
||||
agent_control,
|
||||
network_proxy: None,
|
||||
network_approval: Arc::clone(&network_approval),
|
||||
@@ -5320,7 +5325,10 @@ async fn submission_loop_channel_close_emits_thread_stop_lifecycle() {
|
||||
|
||||
impl codex_extension_api::ThreadLifecycleContributor<crate::config::Config> for ThreadStopRecorder {
|
||||
fn on_thread_stop(&self, input: codex_extension_api::ThreadStopInput<'_>) {
|
||||
assert_eq!(self.expected_thread_id, input.thread_id);
|
||||
assert_eq!(
|
||||
self.expected_thread_id.to_string(),
|
||||
input.thread_store.level_id()
|
||||
);
|
||||
assert!(input.session_store.get::<SessionStopMarker>().is_some());
|
||||
assert!(input.thread_store.get::<ThreadStopMarker>().is_some());
|
||||
self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
@@ -5362,7 +5370,10 @@ async fn submission_loop_channel_close_aborts_active_turn_before_thread_stop_lif
|
||||
|
||||
impl codex_extension_api::ThreadLifecycleContributor<crate::config::Config> for LifecycleRecorder {
|
||||
fn on_thread_stop(&self, input: codex_extension_api::ThreadStopInput<'_>) {
|
||||
assert_eq!(self.expected_thread_id, input.thread_id);
|
||||
assert_eq!(
|
||||
self.expected_thread_id.to_string(),
|
||||
input.thread_store.level_id()
|
||||
);
|
||||
self.calls
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
@@ -5372,8 +5383,11 @@ async fn submission_loop_channel_close_aborts_active_turn_before_thread_stop_lif
|
||||
|
||||
impl codex_extension_api::TurnLifecycleContributor for LifecycleRecorder {
|
||||
fn on_turn_abort(&self, input: codex_extension_api::TurnAbortInput<'_>) {
|
||||
assert_eq!(self.expected_thread_id, input.thread_id);
|
||||
assert_eq!(self.expected_turn_id, input.turn_id);
|
||||
assert_eq!(
|
||||
self.expected_thread_id.to_string(),
|
||||
input.thread_store.level_id()
|
||||
);
|
||||
assert_eq!(self.expected_turn_id, input.turn_store.level_id());
|
||||
assert_eq!(TurnAbortReason::Interrupted, input.reason);
|
||||
self.calls
|
||||
.lock()
|
||||
@@ -5811,8 +5825,10 @@ where
|
||||
plugins_manager,
|
||||
mcp_manager,
|
||||
extensions: Arc::new(codex_extension_api::ExtensionRegistryBuilder::new().build()),
|
||||
session_extension_data: codex_extension_api::ExtensionData::new(),
|
||||
thread_extension_data: codex_extension_api::ExtensionData::new(),
|
||||
session_extension_data: codex_extension_api::ExtensionData::new(
|
||||
agent_control.session_id().to_string(),
|
||||
),
|
||||
thread_extension_data: codex_extension_api::ExtensionData::new(thread_id.to_string()),
|
||||
agent_control,
|
||||
network_proxy: None,
|
||||
network_approval: Arc::clone(&network_approval),
|
||||
|
||||
@@ -92,6 +92,7 @@ pub struct TurnContext {
|
||||
pub(crate) truncation_policy: TruncationPolicy,
|
||||
pub(crate) dynamic_tools: Vec<DynamicToolSpec>,
|
||||
pub(crate) turn_metadata_state: Arc<TurnMetadataState>,
|
||||
pub(crate) extension_data: Arc<codex_extension_api::ExtensionData>,
|
||||
pub(crate) turn_skills: TurnSkillsContext,
|
||||
pub(crate) turn_timing_state: Arc<TurnTimingState>,
|
||||
pub(crate) server_model_warning_emitted: AtomicBool,
|
||||
@@ -274,6 +275,7 @@ impl TurnContext {
|
||||
truncation_policy,
|
||||
dynamic_tools: self.dynamic_tools.clone(),
|
||||
turn_metadata_state: self.turn_metadata_state.clone(),
|
||||
extension_data: Arc::clone(&self.extension_data),
|
||||
turn_skills: self.turn_skills.clone(),
|
||||
turn_timing_state: Arc::clone(&self.turn_timing_state),
|
||||
server_model_warning_emitted: AtomicBool::new(
|
||||
@@ -535,6 +537,7 @@ impl Session {
|
||||
network.is_some(),
|
||||
));
|
||||
let (current_date, timezone) = local_time_context();
|
||||
let extension_data = Arc::new(codex_extension_api::ExtensionData::new(sub_id.clone()));
|
||||
TurnContext {
|
||||
sub_id,
|
||||
trace_id: current_span_trace_id(),
|
||||
@@ -572,6 +575,7 @@ impl Session {
|
||||
truncation_policy: model_info.truncation_policy.into(),
|
||||
dynamic_tools: session_configuration.dynamic_tools.clone(),
|
||||
turn_metadata_state,
|
||||
extension_data,
|
||||
turn_skills: TurnSkillsContext::new(skills_outcome),
|
||||
turn_timing_state: Arc::new(TurnTimingState::default()),
|
||||
server_model_warning_emitted: AtomicBool::new(false),
|
||||
|
||||
Reference in New Issue
Block a user