//! Session-wide mutable state. use codex_protocol::models::AdditionalPermissionProfile; use codex_protocol::models::BaseInstructionsProvenance; use codex_protocol::models::ResponseItem; use codex_protocol::openai_models::ReasoningEffort; use codex_sandboxing::policy_transforms::merge_permission_profiles; use std::collections::HashMap; use std::collections::HashSet; use std::collections::VecDeque; use super::AdditionalContextStore; use super::auto_compact_window::AutoCompactWindow; use super::auto_compact_window::AutoCompactWindowIds; use super::auto_compact_window::AutoCompactWindowSnapshot; use crate::context_manager::ContextManager; use crate::context_manager::HistoryReplacement; use crate::session::PreviousTurnSettings; use crate::session::session::SessionConfiguration; use crate::session::time_reminder::CurrentTimeReminderState; use crate::session_startup_prewarm::SessionStartupPrewarmHandle; use codex_history::ResponseItemEnvelope; use codex_protocol::SessionId; use codex_protocol::ThreadId; use codex_protocol::protocol::RateLimitSnapshot; use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::TokenUsageInfo; use codex_protocol::protocol::TokenUsageRecord; use codex_protocol::protocol::TurnContextItem; use codex_utils_output_truncation::TruncationPolicy; use tokio_util::sync::CancellationToken; use tokio_util::task::AbortOnDropHandle; /// Runtime request effort, initially unset and established by prewarm or sampling. /// Successful compaction allows a fresh baseline without an override. pub(crate) enum ReasoningEffortPin { Unset, Compacted, Active { model: String, effort: ReasoningEffort, }, } impl ReasoningEffortPin { pub(crate) fn get(&self, model: &str) -> Option { match self { Self::Active { model: pinned_model, effort, } if pinned_model == model => Some(effort.clone()), Self::Unset | Self::Compacted | Self::Active { .. } => None, } } pub(crate) fn pin(&mut self, model: &str, effort: ReasoningEffort) -> ReasoningEffort { if let Some(pinned) = self.get(model) { return pinned; } *self = Self::Active { model: model.to_owned(), effort: effort.clone(), }; effort } } /// Persistent, session-scoped state previously stored directly on `Session`. pub(crate) struct SessionState { pub(crate) session_configuration: SessionConfiguration, /// Plugin selection of the last admitted task; settings updates take effect on the next task. pub(crate) active_disabled_plugin_ids: Vec, /// Persisted origin of the session base instructions, when known. pub(crate) base_instructions_provenance: Option, pub(crate) history: ContextManager, /// Cancels work bound to discarded history or a superseded Guardian evidence policy. pub(crate) history_reset: CancellationToken, pub(crate) latest_rate_limits: Option, pub(crate) latest_token_usage_record: Option, pub(crate) server_reasoning_included: bool, pub(crate) mcp_dependency_prompted: HashSet, pub(crate) additional_context: AdditionalContextStore, /// Settings used by the latest regular user turn, used for turn-to-turn /// model/realtime handling on subsequent regular turns (including full-context /// reinjection after resume or `/compact`). previous_turn_settings: Option, /// Latest task admitted in this runtime, retained across completion and history edits. /// Cleared by standalone settings changes to invalidate pending continuation. pub(crate) last_started_turn_id: Option, /// Runtime accounting state for the active auto-compaction window. auto_compact_window: AutoCompactWindow, /// Original request effort for the current model while configuration updates remain active. pub(crate) reasoning_effort_pin: ReasoningEffortPin, /// Startup prewarmed session prepared during session initialization. pub(crate) startup_prewarm: Option, /// Retained after completion so later turns do not repeat speculative captures. pub(crate) shell_snapshot_prewarm: Option>, pub(crate) current_time_reminder: CurrentTimeReminderState, pub(crate) active_connector_selection: HashSet, pub(crate) pending_session_start_sources: VecDeque, granted_permissions_by_environment_id: HashMap, next_turn_is_first: bool, } impl SessionState { /// Create a new session state mirroring previous `State::default()` semantics. #[cfg(test)] pub(crate) fn new(session_configuration: SessionConfiguration) -> Self { Self::new_with_auto_compact_window_ids( session_configuration, AutoCompactWindowIds::new_initial(), ContextManager::new(), ) } pub(crate) fn new_with_auto_compact_window_ids( session_configuration: SessionConfiguration, auto_compact_window_ids: AutoCompactWindowIds, history: ContextManager, ) -> Self { Self { active_disabled_plugin_ids: Vec::new(), session_configuration, base_instructions_provenance: None, history, history_reset: CancellationToken::new(), latest_rate_limits: None, latest_token_usage_record: None, server_reasoning_included: false, mcp_dependency_prompted: HashSet::new(), additional_context: AdditionalContextStore::default(), previous_turn_settings: None, last_started_turn_id: None, auto_compact_window: AutoCompactWindow::new_with_ids(auto_compact_window_ids), reasoning_effort_pin: ReasoningEffortPin::Unset, startup_prewarm: None, shell_snapshot_prewarm: None, current_time_reminder: CurrentTimeReminderState::default(), active_connector_selection: HashSet::new(), pending_session_start_sources: VecDeque::new(), granted_permissions_by_environment_id: HashMap::new(), next_turn_is_first: true, } } // History helpers pub(crate) fn record_items(&mut self, items: I, policy: TruncationPolicy) where I: IntoIterator, I::Item: std::ops::Deref, { self.history.record_items(items, policy); } pub(crate) fn previous_turn_settings(&self) -> Option { self.previous_turn_settings.clone() } pub(crate) fn set_previous_turn_settings( &mut self, previous_turn_settings: Option, ) { self.previous_turn_settings = previous_turn_settings; } pub(crate) fn set_next_turn_is_first(&mut self, value: bool) { self.next_turn_is_first = value; } pub(crate) fn take_next_turn_is_first(&mut self) -> bool { let is_first_turn = self.next_turn_is_first; self.next_turn_is_first = false; is_first_turn } pub(crate) fn clone_history(&self) -> ContextManager { self.history.clone() } #[cfg(test)] pub(crate) fn replace_history( &mut self, items: Vec, reference_context_item: Option, ) { self.replace_annotated_history( items.into_iter().map(ResponseItemEnvelope::new).collect(), reference_context_item, HistoryReplacement::Reset, ); } pub(crate) fn replace_annotated_history( &mut self, items: Vec, reference_context_item: Option, replacement: HistoryReplacement, ) { let invalidate_reviews = match replacement { HistoryReplacement::Compaction { reviewer_compaction_hash, } => self .history .replace_compacted(items, reviewer_compaction_hash.as_deref()), HistoryReplacement::Reset => { self.history.replace_annotated(items); true } }; if invalidate_reviews { std::mem::take(&mut self.history_reset).cancel(); } self.history .set_reference_context_item(reference_context_item); self.auto_compact_window.clear_prefill(); } pub(crate) fn set_token_info(&mut self, info: Option) { self.history.set_token_info(info); } pub(crate) fn record_token_usage( &mut self, thread_id: ThreadId, turn_id: &str, session_id: SessionId, root_turn_id: String, response_id: String, usage: &TokenUsage, ) -> TokenUsageRecord { let mut turn_token_usage = self .latest_token_usage_record .as_ref() .filter(|record| record.turn_id == turn_id) .map_or_else(TokenUsage::default, |record| { record.turn_token_usage.clone() }); turn_token_usage.add_assign(usage); let mut thread_token_usage = self .latest_token_usage_record .as_ref() .map_or_else(TokenUsage::default, |record| { record.thread_token_usage.clone() }); thread_token_usage.add_assign(usage); let record = TokenUsageRecord { thread_id, turn_id: turn_id.to_string(), session_id, root_turn_id, response_id, usage: usage.clone(), turn_token_usage, thread_token_usage, }; self.latest_token_usage_record = Some(record.clone()); record } pub(crate) fn set_reference_context_item(&mut self, item: Option) { self.history.set_reference_context_item(item); } pub(crate) fn reference_context_item(&self) -> Option { self.history.reference_context_item() } // Token/rate limit helpers pub(crate) fn update_token_info_from_usage( &mut self, usage: &TokenUsage, model_context_window: Option, ) { self.history.update_token_info(usage, model_context_window); } pub(crate) fn ensure_auto_compact_window_server_prefill_from_usage( &mut self, usage: &TokenUsage, ) { self.auto_compact_window .ensure_server_observed_prefill_from_usage(usage); } pub(crate) fn set_auto_compact_window_estimated_prefill(&mut self, tokens: i64) { self.auto_compact_window.set_estimated_prefill(tokens); } pub(crate) fn auto_compact_window_snapshot(&self) -> AutoCompactWindowSnapshot { self.auto_compact_window.snapshot() } pub(crate) fn claim_token_budget_reminder(&mut self) -> bool { self.auto_compact_window.claim_token_budget_reminder() } pub(crate) fn claim_auto_compact_fallback(&mut self) -> bool { self.auto_compact_window.claim_auto_compact_fallback() } pub(crate) fn auto_compact_window_number(&self) -> u64 { self.auto_compact_window.window_number() } pub(crate) fn auto_compact_window_ids(&self) -> AutoCompactWindowIds { self.auto_compact_window.ids() } pub(crate) fn restore_auto_compact_window( &mut self, window_number: u64, ids: AutoCompactWindowIds, ) { self.auto_compact_window.restore(window_number, ids); } pub(crate) fn advance_auto_compact_window(&mut self) -> (u64, AutoCompactWindowIds) { self.auto_compact_window.advance() } pub(crate) fn request_new_context_window(&mut self) { self.auto_compact_window.request_new_context_window(); } pub(crate) fn take_new_context_window_request(&mut self) -> bool { self.auto_compact_window.take_new_context_window_request() } pub(crate) fn start_new_context_window(&mut self) -> (u64, AutoCompactWindowIds) { let window = self.auto_compact_window.advance(); self.auto_compact_window.clear_prefill(); window } pub(crate) fn token_info(&self) -> Option { self.history.token_info() } pub(crate) fn set_rate_limits(&mut self, snapshot: RateLimitSnapshot) { self.latest_rate_limits = Some(merge_rate_limit_fields( self.latest_rate_limits.as_ref(), snapshot, )); } pub(crate) fn token_info_and_rate_limits( &self, ) -> (Option, Option) { (self.token_info(), self.latest_rate_limits.clone()) } pub(crate) fn set_token_usage_full(&mut self, context_window: i64) { self.history.set_token_usage_full(context_window); } pub(crate) fn get_total_token_usage(&self, server_reasoning_included: bool) -> i64 { self.history .get_total_token_usage(server_reasoning_included) } pub(crate) fn set_server_reasoning_included(&mut self, included: bool) { self.server_reasoning_included = included; } pub(crate) fn server_reasoning_included(&self) -> bool { self.server_reasoning_included } pub(crate) fn record_mcp_dependency_prompted(&mut self, names: I) where I: IntoIterator, { self.mcp_dependency_prompted.extend(names); } pub(crate) fn mcp_dependency_prompted(&self) -> HashSet { self.mcp_dependency_prompted.clone() } pub(crate) fn set_session_startup_prewarm( &mut self, startup_prewarm: SessionStartupPrewarmHandle, ) { self.startup_prewarm = Some(startup_prewarm); } pub(crate) fn take_session_startup_prewarm(&mut self) -> Option { self.startup_prewarm.take() } // Adds connector IDs to the active set and returns the merged selection. pub(crate) fn merge_connector_selection(&mut self, connector_ids: I) -> HashSet where I: IntoIterator, { self.active_connector_selection.extend(connector_ids); self.active_connector_selection.clone() } // Returns the current connector selection tracked on session state. pub(crate) fn get_connector_selection(&self) -> HashSet { self.active_connector_selection.clone() } // Removes all currently tracked connector selections. pub(crate) fn clear_connector_selection(&mut self) { self.active_connector_selection.clear(); } pub(crate) fn queue_pending_session_start_source( &mut self, value: codex_hooks::SessionStartSource, ) { self.pending_session_start_sources.push_back(value); } pub(crate) fn take_pending_session_start_source( &mut self, ) -> Option { self.pending_session_start_sources.pop_front() } pub(crate) fn record_granted_permissions( &mut self, environment_id: &str, permissions: AdditionalPermissionProfile, ) { let granted_permissions = merge_permission_profiles( self.granted_permissions_by_environment_id .get(environment_id), Some(&permissions), ); if let Some(granted_permissions) = granted_permissions { self.granted_permissions_by_environment_id .insert(environment_id.to_string(), granted_permissions); } } pub(crate) fn granted_permissions( &self, environment_id: &str, ) -> Option { self.granted_permissions_by_environment_id .get(environment_id) .cloned() } } // Sometimes new snapshots don't include credits or plan information. // Preserve those from the previous snapshot when missing. For `limit_id`, treat // missing values as the default `"codex"` bucket. fn merge_rate_limit_fields( previous: Option<&RateLimitSnapshot>, mut snapshot: RateLimitSnapshot, ) -> RateLimitSnapshot { if snapshot.limit_id.is_none() { snapshot.limit_id = Some("codex".to_string()); } if snapshot.credits.is_none() { snapshot.credits = previous.and_then(|prior| prior.credits.clone()); } if snapshot.individual_limit.is_none() { snapshot.individual_limit = previous.and_then(|prior| prior.individual_limit.clone()); } if snapshot.spend_control_reached.is_none() { snapshot.spend_control_reached = previous.and_then(|prior| prior.spend_control_reached); } if snapshot.plan_type.is_none() { snapshot.plan_type = previous.and_then(|prior| prior.plan_type); } snapshot } #[cfg(test)] #[path = "session_tests.rs"] mod tests;