| use super::input_queue::InputQueue; |
| use super::mcp_refresh::McpRefresh; |
| use super::step_context::StepContext; |
| use super::step_settings::ModelInfoOverrides; |
| use super::step_settings::StepSettings; |
| use super::step_settings::StepSettingsConstraints; |
| use super::step_settings::StepSettingsUpdate; |
| use super::*; |
| use crate::agents_md_manager::AgentsMdManager; |
| use crate::agents_md_manager::SessionInstructions; |
| use crate::config::ConstraintError; |
| use crate::context::GuardianContextMode; |
| use crate::environment_selection::ThreadEnvironments; |
| use crate::environment_selection::TurnEnvironmentSnapshot; |
| use crate::hook_mcp_executor::CoreHookMcpExecutor; |
| use crate::mcp_tool_call::McpToolApprovalMetadata; |
| use crate::responses_metadata::CodexResponsesMetadata; |
| use crate::responses_metadata::CodexResponsesRequestKind; |
| use crate::responses_metadata::CompactionTurnMetadata; |
| use crate::shell_snapshot::ShellSnapshot; |
| use crate::shell_snapshot::SnapshotCredentialBrokerState; |
| use crate::state::ActiveTurn; |
| use crate::turn_metadata::ExecutionMetadata; |
| use codex_attachment_store::AttachmentStore; |
| use codex_extension_api::ExtensionDataInit; |
| use codex_http_client::ClientRouteClass; |
| use codex_http_client::RouteAwareClientPool; |
| use codex_login::auth::AgentIdentityAuthPolicy; |
| use codex_model_provider::SharedModelProvider; |
| use codex_protocol::SessionId; |
| use codex_protocol::capabilities::SelectedCapabilityRoot; |
| use codex_protocol::config_types::ShellEnvironmentPolicy; |
| use codex_protocol::mcp::ClientMcpExtensions; |
| use codex_protocol::models::ProfileWorkspaceRoot; |
| use codex_protocol::permissions::FileSystemPath; |
| use codex_protocol::permissions::FileSystemSpecialPath; |
| use codex_protocol::protocol::EnvironmentConfig; |
| use codex_protocol::protocol::HookCompletedEvent; |
| use codex_protocol::protocol::McpInvocation; |
| use codex_protocol::protocol::MultiAgentVersion; |
| use codex_protocol::protocol::ThreadHistoryMode; |
| use codex_protocol::protocol::ThreadSource; |
| use codex_protocol::protocol::TurnEnvironmentSelections; |
| use codex_sandboxing::SandboxType; |
| use codex_skills::SkillError; |
| use codex_utils_git_discovery::GitRootDiscovery; |
| use codex_utils_path::replace_path_and_deduplicate; |
| use std::sync::OnceLock; |
| use tokio::sync::Semaphore; |
|
|
| type McpToolApprovalMetadataMap = |
| HashMap<(String, String), std::sync::Weak<(Option<McpInvocation>, McpToolApprovalMetadata)>>; |
|
|
| |
| |
| |
| pub(crate) struct Session { |
| pub(crate) thread_id: ThreadId, |
| pub(crate) installation_id: String, |
| pub(super) tx_event: Sender<Event>, |
| pub(super) agent_status: watch::Sender<AgentStatus>, |
| pub(super) state: Mutex<SessionState>, |
| |
| |
| pub(super) thread_settings_persistence: Semaphore, |
| |
| |
| pub(super) managed_network_proxy_refresh_lock: Semaphore, |
| |
| |
| pub(super) features: ManagedFeatures, |
| pub(crate) guardian_context_mode: GuardianContextMode, |
| pub(super) isolation: codex_extension_api::SessionIsolation, |
| pub(crate) allowed_tools: Option<Arc<codex_extension_api::AllowedTools>>, |
| pub(crate) windows_sandbox_proxy_settings_mode: |
| codex_sandboxing::WindowsSandboxProxySettingsMode, |
| pub(super) multi_agent_version: OnceLock<MultiAgentVersion>, |
| |
| pub(super) mcp_refresh: McpRefresh, |
| |
| pub(crate) mcp_tool_approval_metadata: std::sync::Mutex<McpToolApprovalMetadataMap>, |
| pub(super) mcp_elicitation_reviewer_handle: OnceLock<codex_mcp::ElicitationReviewerHandle>, |
| pub(super) mcp_elicitation_lifecycle_handle: OnceLock<codex_mcp::ElicitationLifecycle>, |
| pub(super) mcp_prewarm_tx: async_channel::Sender<()>, |
| pub(super) mcp_prewarm_shutdown: CancellationToken, |
| pub(super) mcp_prewarm_task: std::sync::Mutex<Option<JoinHandle<()>>>, |
| pub(crate) conversation: Arc<RealtimeConversationManager>, |
| pub(crate) realtime_history: Option<Mutex<crate::realtime_history::RealtimeHistoryState>>, |
| pub(crate) active_turn: Mutex<Option<ActiveTurn>>, |
| pub(crate) async_hook_results: async_channel::Receiver<HookCompletedEvent>, |
| pub(crate) input_queue: InputQueue, |
| pub(crate) services: SessionServices, |
| pub(super) git_enrichment_policy: GitEnrichmentPolicy, |
| pub(super) fork_persistence: ForkPersistence, |
| pub(super) forked_from_ordinal_exclusive: Option<u64>, |
| pub(super) next_internal_sub_id: AtomicU64, |
| } |
|
|
| #[derive(Clone)] |
| pub(crate) struct SessionConfiguration { |
| |
| pub(super) provider: SharedModelProvider, |
|
|
| |
| pub(super) step_settings: Arc<StepSettings>, |
| |
| pub(super) model_info_overrides: ModelInfoOverrides, |
|
|
| |
| pub(super) developer_instructions: Option<String>, |
|
|
| |
| pub(super) base_instructions: String, |
|
|
| |
| |
| |
| pub(super) permission_profile_state: PermissionProfileState, |
| pub(super) allow_login_shell: bool, |
| pub(super) shell_environment_policy: ShellEnvironmentPolicy, |
| |
| |
| pub(super) windows_sandbox_level: WindowsSandboxLevel, |
| pub(super) windows_sandbox_type: SandboxType, |
| pub(super) windows_sandbox_private_desktop: bool, |
| pub(super) use_legacy_landlock: bool, |
|
|
| |
| pub(super) legacy_fallback_cwd: AbsolutePathBuf, |
| |
| pub(super) runtime_workspace_roots: Vec<AbsolutePathBuf>, |
| |
| pub(super) codex_home: AbsolutePathBuf, |
| |
| pub(super) thread_name: Option<String>, |
| |
| pub(super) disabled_plugin_ids: Vec<String>, |
|
|
| |
| pub(super) original_config_do_not_use: Arc<Config>, |
| |
| pub(super) metrics_service_name: Option<String>, |
| pub(super) app_server_client_name: Option<String>, |
| pub(super) app_server_client_version: Option<String>, |
| |
| pub(super) trusted_guardian_reviewer: bool, |
| |
| pub(super) session_source: SessionSource, |
| |
| pub(super) history_mode: ThreadHistoryMode, |
| |
| pub(super) forked_from_thread_id: Option<ThreadId>, |
| |
| pub(super) parent_thread_id: Option<ThreadId>, |
| |
| pub(super) thread_source: Option<ThreadSource>, |
| |
| pub(super) originator: String, |
| pub(super) dynamic_tools: Vec<DynamicToolSpec>, |
| pub(super) user_shell_override: Option<shell::Shell>, |
| } |
|
|
| impl SessionConfiguration { |
| pub(super) fn cwd(&self) -> &AbsolutePathBuf { |
| &self.legacy_fallback_cwd |
| } |
|
|
| pub(crate) fn codex_home(&self) -> &AbsolutePathBuf { |
| &self.codex_home |
| } |
|
|
| pub(super) fn inferred_environment_config(&self) -> EnvironmentConfig { |
| EnvironmentConfig { |
| allow_login_shell: self.allow_login_shell, |
| workspace_roots: Vec::new(), |
| permission_profile: self.permission_profile_state.snapshot(), |
| shell_environment_policy: self.shell_environment_policy.clone(), |
| windows_sandbox_level: self.windows_sandbox_level, |
| windows_sandbox_private_desktop: self.windows_sandbox_private_desktop, |
| use_legacy_landlock: self.use_legacy_landlock, |
| exec_policy: None, |
| mcp_policy: None, |
| network_policy: None, |
| selected_capability_roots: Vec::new(), |
| } |
| } |
|
|
| pub(super) fn permission_profile(&self) -> PermissionProfile { |
| self.permission_profile_state.permission_profile().clone() |
| } |
|
|
| fn materialized_permission_profile( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> PermissionProfile { |
| let workspace_roots = ThreadEnvironments::primary_workspace_roots_for(environments); |
| self.permission_profile() |
| .materialize_project_roots_with_workspace_roots(&workspace_roots) |
| } |
|
|
| fn effective_permission_profile( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> PermissionProfile { |
| ThreadEnvironments::primary_config_for(environments) |
| .map(|config| { |
| config |
| .permission_profile |
| .permission_profile() |
| .clone() |
| .materialize_project_roots_with_path_uris(&config.workspace_roots) |
| }) |
| .unwrap_or_else(|| self.materialized_permission_profile(environments)) |
| } |
|
|
| pub(super) fn active_permission_profile(&self) -> Option<ActivePermissionProfile> { |
| self.permission_profile_state.active_permission_profile() |
| } |
|
|
| pub(super) fn apply_permission_profile_to_permissions( |
| &self, |
| permissions: &mut crate::config::Permissions, |
| ) { |
| permissions.set_permission_profile_state(self.permission_profile_state.clone()); |
| } |
|
|
| #[cfg(test)] |
| pub(super) fn set_permission_profile_for_tests( |
| &mut self, |
| permission_profile: PermissionProfile, |
| ) -> ConstraintResult<()> { |
| self.permission_profile_state |
| .set_legacy_permission_profile(permission_profile) |
| } |
|
|
| pub(super) fn sandbox_policy( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> SandboxPolicy { |
| let permission_profile = self.materialized_permission_profile(environments); |
| codex_sandboxing::compatibility_sandbox_policy_for_permission_profile( |
| &permission_profile, |
| self.cwd(), |
| ) |
| } |
|
|
| pub(super) fn file_system_sandbox_policy( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> FileSystemSandboxPolicy { |
| self.materialized_permission_profile(environments) |
| .file_system_sandbox_policy() |
| } |
|
|
| pub(super) fn network_sandbox_policy(&self) -> NetworkSandboxPolicy { |
| self.permission_profile_state |
| .permission_profile() |
| .network_sandbox_policy() |
| } |
|
|
| pub(super) fn thread_config_snapshot( |
| &self, |
| environment_selections: Vec<TurnEnvironmentSelection>, |
| ) -> ThreadConfigSnapshot { |
| let workspace_roots = |
| ThreadEnvironments::primary_workspace_roots_for(&environment_selections); |
| let permission_profile = ThreadEnvironments::primary_config_for(&environment_selections) |
| .map(|config| config.permission_profile.clone()) |
| .unwrap_or_else(|| self.permission_profile_state.snapshot()); |
| ThreadConfigSnapshot { |
| model: self.step_settings.collaboration_mode.model().to_string(), |
| model_provider_id: self.original_config_do_not_use.model_provider_id.clone(), |
| service_tier: self.step_settings.service_tier.clone(), |
| approval_policy: self.step_settings.approval_policy.value(), |
| approvals_reviewer: self.step_settings.approvals_reviewer, |
| permission_profile: self.effective_permission_profile(&environment_selections), |
| full_access: codex_protocol::protocol::has_full_access( |
| self.step_settings.approval_policy.value(), |
| &self.permission_profile(), |
| environment_selections |
| .iter() |
| .map(|environment| &environment.config), |
| ), |
| active_permission_profile: permission_profile.active_permission_profile(), |
| environments: TurnEnvironmentSelections::new( |
| self.legacy_fallback_cwd.clone(), |
| environment_selections, |
| ), |
| workspace_roots, |
| profile_workspace_roots: permission_profile.profile_workspace_roots().to_vec(), |
| ephemeral: self.original_config_do_not_use.ephemeral, |
| reasoning_effort: self.step_settings.collaboration_mode.reasoning_effort(), |
| reasoning_summary: self.step_settings.reasoning_summary, |
| personality: self.step_settings.personality, |
| collaboration_mode: self.step_settings.collaboration_mode.clone(), |
| session_source: self.session_source.clone(), |
| history_mode: self.history_mode, |
| forked_from_thread_id: self.forked_from_thread_id, |
| parent_thread_id: self.parent_thread_id, |
| thread_source: self.thread_source.clone(), |
| originator: self.originator.clone(), |
| disabled_plugin_ids: self.disabled_plugin_ids.clone(), |
| } |
| } |
|
|
| |
| pub(super) fn thread_settings_snapshot( |
| &self, |
| environment_selections: &[TurnEnvironmentSelection], |
| ) -> ThreadSettingsSnapshot { |
| ThreadSettingsSnapshot { |
| model: self.step_settings.collaboration_mode.model().to_string(), |
| model_provider_id: self.original_config_do_not_use.model_provider_id.clone(), |
| service_tier: self.step_settings.service_tier.clone(), |
| approval_policy: self.step_settings.approval_policy.value(), |
| approvals_reviewer: self.step_settings.approvals_reviewer, |
| permission_profile: self.materialized_permission_profile(environment_selections), |
| active_permission_profile: self.active_permission_profile(), |
| cwd: self.legacy_fallback_cwd.clone(), |
| runtime_workspace_roots: Some(self.runtime_workspace_roots.clone()), |
| reasoning_effort: self.step_settings.collaboration_mode.reasoning_effort(), |
| reasoning_summary: self.step_settings.reasoning_summary, |
| personality: self.step_settings.personality, |
| collaboration_mode: self.step_settings.collaboration_mode.clone(), |
| disabled_plugin_ids: self.disabled_plugin_ids.clone(), |
| } |
| } |
|
|
| |
| pub(super) fn restorable_thread_settings( |
| &self, |
| environment_selections: Vec<TurnEnvironmentSelection>, |
| ) -> CodexThreadSettingsOverrides { |
| CodexThreadSettingsOverrides { |
| environments: Some(TurnEnvironmentSelections::new( |
| self.legacy_fallback_cwd.clone(), |
| environment_selections, |
| )), |
| runtime_workspace_roots: Some(self.runtime_workspace_roots.clone()), |
| profile_workspace_roots: Some( |
| self.permission_profile_state |
| .profile_workspace_roots() |
| .to_vec(), |
| ), |
| approval_policy: Some(self.step_settings.approval_policy.value()), |
| approvals_reviewer: Some(self.step_settings.approvals_reviewer), |
| permission_profile: Some(self.permission_profile()), |
| active_permission_profile: self.active_permission_profile(), |
| windows_sandbox_level: Some(self.windows_sandbox_level), |
| summary: self.step_settings.reasoning_summary, |
| service_tier: Some(self.step_settings.service_tier.clone()), |
| collaboration_mode: Some(self.step_settings.collaboration_mode.clone()), |
| personality: self.step_settings.personality, |
| disabled_plugin_ids: Some(self.disabled_plugin_ids.clone()), |
| ..Default::default() |
| } |
| } |
|
|
| pub(super) fn validate( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> ConstraintResult<()> { |
| self.step_settings |
| .validate(&self.step_settings_constraints(environments))?; |
| super::environment::validate_environment_selections(environments) |
| } |
|
|
| pub(super) fn step_settings_constraints( |
| &self, |
| environments: &[TurnEnvironmentSelection], |
| ) -> StepSettingsConstraints<'_> { |
| let permission_profile = self.effective_permission_profile(environments); |
| StepSettingsConstraints { |
| requirements: self |
| .original_config_do_not_use |
| .config_layer_stack |
| .requirements(), |
| guardian_approval_enabled: self |
| .original_config_do_not_use |
| .features |
| .enabled(Feature::GuardianApproval), |
| trusted_guardian_reviewer: self.trusted_guardian_reviewer, |
| has_full_disk_write_access: permission_profile |
| .file_system_sandbox_policy() |
| .has_full_disk_write_access(), |
| } |
| } |
|
|
| pub(super) fn apply( |
| &self, |
| updates: &SessionSettingsUpdate, |
| current_environments: &[TurnEnvironmentSelection], |
| ) -> ConstraintResult<Self> { |
| let mut next_configuration = self.clone(); |
| if let Some(disabled_plugin_ids) = &updates.disabled_plugin_ids { |
| next_configuration.disabled_plugin_ids = disabled_plugin_ids.clone(); |
| } |
| let current_file_system_sandbox_policy = |
| self.file_system_sandbox_policy(current_environments); |
| let file_system_policy_has_rebindable_project_root_write = |
| current_file_system_sandbox_policy |
| .entries |
| .iter() |
| .any(|entry| { |
| entry.access.can_write() |
| && matches!( |
| &entry.path, |
| FileSystemPath::Special { |
| value: FileSystemSpecialPath::ProjectRoots { subpath: None }, |
| } |
| ) |
| }); |
| if let Some(windows_sandbox_level) = updates.windows_sandbox_level { |
| next_configuration.windows_sandbox_level = windows_sandbox_level; |
| } |
|
|
| let current_cwd = self.cwd().clone(); |
| if let Some(environments) = &updates.environments { |
| next_configuration.legacy_fallback_cwd = environments.legacy_fallback_cwd.clone(); |
| } |
| let cwd_changed = next_configuration.legacy_fallback_cwd != current_cwd; |
| if let Some(runtime_workspace_roots) = &updates.runtime_workspace_roots { |
| next_configuration.runtime_workspace_roots = runtime_workspace_roots.clone(); |
| } else if cwd_changed { |
| next_configuration.runtime_workspace_roots = replace_path_and_deduplicate( |
| next_configuration.runtime_workspace_roots, |
| current_cwd.as_path(), |
| next_configuration.legacy_fallback_cwd.clone(), |
| ); |
| } |
|
|
| if let Some(permission_profile) = updates.permission_profile.clone() { |
| let active_permission_profile = |
| updates.active_permission_profile.clone().or_else(|| { |
| if permission_profile == self.permission_profile() { |
| self.active_permission_profile() |
| } else { |
| None |
| } |
| }); |
| next_configuration.set_permission_profile_projection( |
| permission_profile, |
| active_permission_profile, |
| updates.profile_workspace_roots.clone().unwrap_or_default(), |
| Some(¤t_file_system_sandbox_policy), |
| )?; |
| if let Some(active_permission_profile) = next_configuration.active_permission_profile() |
| { |
| let mut config = (*next_configuration.original_config_do_not_use).clone(); |
| let permission_profile = next_configuration.permission_profile(); |
| config.permissions.network = config |
| .network_proxy_spec_for_active_permission_profile( |
| &active_permission_profile, |
| &permission_profile, |
| ) |
| .map_err(|err| ConstraintError::InvalidValue { |
| field_name: "default_permissions", |
| candidate: active_permission_profile.id.clone(), |
| allowed: format!( |
| "configured permission profile with valid network policy ({err})" |
| ), |
| requirement_source: codex_config::RequirementSource::Unknown, |
| })?; |
| config |
| .permissions |
| .set_permission_profile_from_session_snapshot( |
| PermissionProfileSnapshot::active( |
| permission_profile, |
| active_permission_profile, |
| ), |
| )?; |
| next_configuration.original_config_do_not_use = Arc::new(config); |
| } |
| } else if let Some(sandbox_policy) = updates.sandbox_policy.clone() { |
| let file_system_sandbox_policy = |
| FileSystemSandboxPolicy::from_legacy_sandbox_policy_preserving_deny_entries( |
| &sandbox_policy, |
| next_configuration.cwd(), |
| ¤t_file_system_sandbox_policy, |
| ); |
| let network_sandbox_policy = NetworkSandboxPolicy::from(&sandbox_policy); |
| next_configuration |
| .permission_profile_state |
| .set_legacy_permission_profile( |
| PermissionProfile::from_runtime_permissions_with_enforcement( |
| SandboxEnforcement::from_legacy_sandbox_policy(&sandbox_policy), |
| &file_system_sandbox_policy, |
| network_sandbox_policy, |
| ), |
| )?; |
| } else if cwd_changed && file_system_policy_has_rebindable_project_root_write { |
| |
| |
| let current_sandbox_policy = self.sandbox_policy(current_environments); |
| if current_file_system_sandbox_policy.is_semantically_equivalent_to( |
| &FileSystemSandboxPolicy::from_legacy_sandbox_policy_preserving_deny_entries( |
| ¤t_sandbox_policy, |
| self.cwd(), |
| ¤t_file_system_sandbox_policy, |
| ), |
| self.cwd(), |
| ) { |
| |
| |
| |
| let file_system_sandbox_policy = |
| FileSystemSandboxPolicy::from_legacy_sandbox_policy_preserving_deny_entries( |
| ¤t_sandbox_policy, |
| next_configuration.cwd(), |
| ¤t_file_system_sandbox_policy, |
| ); |
| next_configuration |
| .permission_profile_state |
| .set_legacy_permission_profile( |
| PermissionProfile::from_runtime_permissions_with_enforcement( |
| SandboxEnforcement::from_legacy_sandbox_policy(¤t_sandbox_policy), |
| &file_system_sandbox_policy, |
| self.network_sandbox_policy(), |
| ), |
| )?; |
| } |
| } |
| if let Some(app_server_client_name) = updates.app_server_client_name.clone() { |
| next_configuration.app_server_client_name = Some(app_server_client_name); |
| } |
| if let Some(app_server_client_version) = updates.app_server_client_version.clone() { |
| next_configuration.app_server_client_version = Some(app_server_client_version); |
| } |
| let next_environments = updates |
| .environments |
| .as_ref() |
| .map_or(current_environments, |environments| { |
| environments.environments.as_slice() |
| }); |
| super::environment::validate_environment_selections(next_environments)?; |
| |
| |
| next_configuration.step_settings = Arc::new(self.step_settings.apply( |
| &updates.step_settings, |
| &next_configuration.step_settings_constraints(next_environments), |
| )?); |
| Ok(next_configuration) |
| } |
|
|
| fn set_permission_profile_projection( |
| &mut self, |
| permission_profile: PermissionProfile, |
| active_permission_profile: Option<ActivePermissionProfile>, |
| profile_workspace_roots: Vec<ProfileWorkspaceRoot>, |
| preserve_deny_reads_from: Option<&FileSystemSandboxPolicy>, |
| ) -> ConstraintResult<()> { |
| let enforcement = permission_profile.enforcement(); |
| let (mut file_system_sandbox_policy, network_sandbox_policy) = |
| permission_profile.to_runtime_permissions(); |
| if let Some(existing_file_system_policy) = preserve_deny_reads_from { |
| file_system_sandbox_policy |
| .preserve_deny_read_restrictions_from(existing_file_system_policy); |
| } |
| let effective_permission_profile = |
| PermissionProfile::from_runtime_permissions_with_enforcement( |
| enforcement, |
| &file_system_sandbox_policy, |
| network_sandbox_policy, |
| ); |
|
|
| let permission_snapshot = match active_permission_profile { |
| Some(active_permission_profile) => { |
| PermissionProfileSnapshot::active_with_profile_workspace_roots( |
| effective_permission_profile, |
| active_permission_profile, |
| profile_workspace_roots, |
| ) |
| } |
| None => PermissionProfileSnapshot::legacy(effective_permission_profile), |
| }; |
|
|
| self.permission_profile_state |
| .set_permission_profile_snapshot(permission_snapshot) |
| } |
| } |
|
|
| |
| pub(crate) struct SessionSettingsCommit { |
| pub(crate) configuration: SessionConfiguration, |
| pub(crate) snapshot: ThreadSettingsSnapshot, |
| } |
|
|
| #[derive(Default, Clone)] |
| pub(crate) struct SessionSettingsUpdate { |
| pub(crate) step_settings: StepSettingsUpdate, |
| pub(crate) environments: Option<TurnEnvironmentSelections>, |
| pub(crate) runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>, |
| pub(crate) profile_workspace_roots: Option<Vec<ProfileWorkspaceRoot>>, |
| pub(crate) sandbox_policy: Option<SandboxPolicy>, |
| pub(crate) permission_profile: Option<PermissionProfile>, |
| pub(crate) active_permission_profile: Option<ActivePermissionProfile>, |
| pub(crate) windows_sandbox_level: Option<WindowsSandboxLevel>, |
| pub(crate) service_tier_for_turn: Option<String>, |
| pub(crate) app_server_client_name: Option<String>, |
| pub(crate) app_server_client_version: Option<String>, |
| pub(crate) disabled_plugin_ids: Option<Vec<String>>, |
| } |
|
|
| pub(crate) struct AppServerClientMetadata { |
| pub(crate) client_name: Option<String>, |
| pub(crate) client_version: Option<String>, |
| } |
|
|
| async fn warm_plugins_and_skills_for_session_init( |
| config: Arc<Config>, |
| plugins_manager: Arc<PluginsManager>, |
| skills_service: Arc<HostSkillsService>, |
| turn_environments: &TurnEnvironmentSnapshot, |
| extensions: &codex_extension_api::ExtensionRegistry<Config>, |
| ) -> Vec<SkillError> { |
| let plugins_input = config.plugins_config_input(); |
| let plugin_outcome = plugins_manager.plugins_for_config(&plugins_input).await; |
| if config.features.enabled(Feature::SkipHostSkillDiscovery) |
| && !extensions.requires_host_skill_discovery() |
| { |
| return Vec::new(); |
| } |
|
|
| let fs = turn_environments.primary_filesystem(); |
| let effective_skill_roots = plugin_outcome.effective_plugin_skill_roots(); |
| let plugin_skill_snapshots = plugins_manager.plugin_skill_snapshots_for_config(&plugins_input); |
| let skills_input = skills_load_input_from_config(config.as_ref(), effective_skill_roots) |
| .with_plugin_skill_snapshots(plugin_skill_snapshots); |
| skills_service |
| .snapshot_for_config(&skills_input, fs) |
| .await |
| .outcome() |
| .errors |
| .clone() |
| } |
|
|
| impl Session { |
| |
| pub(crate) fn thread_id(&self) -> ThreadId { |
| self.thread_id |
| } |
|
|
| |
| pub(crate) fn session_id(&self) -> SessionId { |
| self.services.agent_control.session_id() |
| } |
|
|
| pub(crate) async fn originator(&self) -> String { |
| let state = self.state.lock().await; |
| state.session_configuration.originator.clone() |
| } |
|
|
| pub(crate) async fn responses_metadata( |
| &self, |
| step_context: &StepContext, |
| request_kind: CodexResponsesRequestKind, |
| ) -> CodexResponsesMetadata { |
| let (window_id, window_number, context_window_id) = self.current_window().await; |
| let mut responses_metadata = step_context.turn.turn_metadata_state.to_responses_metadata( |
| self.installation_id.clone(), |
| window_id, |
| request_kind, |
| ); |
| ExecutionMetadata::from_settings(&step_context.settings).apply_to(&mut responses_metadata); |
| responses_metadata.tool_namespaces_info = if step_context |
| .turn |
| .config |
| .tool_registry |
| .turn_metadata_includes_tool_info |
| && step_context.settings.model_info.use_responses_lite |
| { |
| step_context.tool_router.tool_namespaces_info().cloned() |
| } else { |
| None |
| }; |
| self.with_window_and_fork_metadata( |
| &step_context.turn, |
| responses_metadata, |
| window_number, |
| context_window_id, |
| ) |
| } |
|
|
| |
| |
| |
| pub(crate) async fn compaction_responses_metadata( |
| &self, |
| turn_context: &TurnContext, |
| compaction_metadata: CompactionTurnMetadata, |
| ) -> CodexResponsesMetadata { |
| let (window_id, window_number, context_window_id) = self.current_window().await; |
| let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( |
| self.installation_id.clone(), |
| window_id, |
| CodexResponsesRequestKind::Compaction(compaction_metadata), |
| ); |
| self.with_window_and_fork_metadata( |
| turn_context, |
| responses_metadata, |
| window_number, |
| context_window_id, |
| ) |
| } |
|
|
| fn with_window_and_fork_metadata( |
| &self, |
| turn_context: &TurnContext, |
| responses_metadata: CodexResponsesMetadata, |
| window_number: u64, |
| context_window_id: uuid::Uuid, |
| ) -> CodexResponsesMetadata { |
| CodexResponsesMetadata { |
| window_number: Some(window_number), |
| context_window_id: Some(context_window_id), |
| analytics_enabled: Some(self.services.analytics_events_client.is_enabled()), |
| history_ingest_requested: turn_context |
| .config |
| .token_budget |
| .as_ref() |
| .is_some_and(|config| config.use_history_notes_extension) |
| .then_some(true), |
| forked_from_ordinal_exclusive: self |
| .forked_from_ordinal_exclusive |
| .filter(|_| responses_metadata.forked_from_thread_id.is_some()), |
| ..responses_metadata |
| } |
| } |
|
|
| #[instrument(name = "session_init", level = "info", skip_all)] |
| #[allow(clippy::too_many_arguments)] |
| pub(crate) async fn new( |
| startup: Option<Arc<super::startup::SessionStartup>>, |
| mut session_configuration: SessionConfiguration, |
| environment_selections: &[TurnEnvironmentSelection], |
| config: Arc<Config>, |
| instructions: SessionInstructions, |
| installation_id: String, |
| auth_manager: Arc<AuthManager>, |
| models_manager: SharedModelsManager, |
| git_root_discovery: Arc<GitRootDiscovery>, |
| model_info: ModelInfo, |
| exec_policy: Arc<ExecPolicyManager>, |
| tx_event: Sender<Event>, |
| agent_status: watch::Sender<AgentStatus>, |
| mut initial_history: InitialHistory, |
| fork_persistence: ForkPersistence, |
| session_source: SessionSource, |
| skills_service: Arc<HostSkillsService>, |
| plugins_manager: Arc<PluginsManager>, |
| mcp_manager: Arc<McpManager>, |
| code_mode_session_provider: Arc<dyn codex_code_mode::CodeModeSessionProvider>, |
| extensions: Arc<codex_extension_api::ExtensionRegistry<crate::config::Config>>, |
| mut thread_extension_init: ExtensionDataInit, |
| client_mcp_extensions: ClientMcpExtensions, |
| agent_control: AgentControl, |
| reserved_thread_id: Option<ThreadId>, |
| environment_manager: Arc<EnvironmentManager>, |
| inherited_environments: Option<TurnEnvironmentSnapshot>, |
| analytics_events_client: Option<AnalyticsEventsClient>, |
| image_store: Arc<dyn AttachmentStore>, |
| thread_store: Arc<dyn ThreadStore>, |
| parent_rollout_thread_trace: ThreadTraceContext, |
| attestation_provider: Option<Arc<dyn AttestationProvider>>, |
| external_time_provider: Option<Arc<dyn TimeProvider>>, |
| multi_agent_version: Option<MultiAgentVersion>, |
| git_enrichment_policy: GitEnrichmentPolicy, |
| windows_sandbox_proxy_settings_mode: codex_sandboxing::WindowsSandboxProxySettingsMode, |
| ) -> anyhow::Result<Arc<Self>> { |
| debug!( |
| "Configuring session: model={}; provider={:?}", |
| session_configuration |
| .step_settings |
| .collaboration_mode |
| .model(), |
| session_configuration.provider |
| ); |
| let base_instructions_provenance = if config.base_instructions.is_some() { |
| Some( |
| config |
| .base_instructions_provenance |
| .clone() |
| .unwrap_or(BaseInstructionsProvenance::Custom), |
| ) |
| } else if let Some(inherited_base_instructions) = initial_history.get_base_instructions() { |
| let BaseInstructions { text, provenance } = inherited_base_instructions; |
| provenance.or_else(|| { |
| (text == model_info.get_model_instructions(config.personality)).then(|| { |
| BaseInstructionsProvenance::Model { |
| model: model_info.slug.clone(), |
| } |
| }) |
| }) |
| } else { |
| Some(BaseInstructionsProvenance::Model { |
| model: model_info.slug.clone(), |
| }) |
| }; |
| let forked_from_id = session_configuration |
| .forked_from_thread_id |
| .or_else(|| initial_history.forked_from_id()); |
| session_configuration.forked_from_thread_id = forked_from_id; |
| let forked_from_ordinal_exclusive = match &fork_persistence { |
| ForkPersistence::Referenced { history_base, .. } => { |
| history_base.map(|position| position.end_ordinal_exclusive) |
| } |
| ForkPersistence::Copied => match &initial_history { |
| InitialHistory::Resumed(resumed) => { |
| |
| |
| |
| resumed.history.first().and_then(|item| match item { |
| RolloutItem::SessionMeta(meta) |
| if meta.meta.id == resumed.conversation_id => |
| { |
| codex_rollout::forked_from_ordinal_exclusive( |
| &meta.meta, |
| resumed.rollout_path.as_deref(), |
| ) |
| } |
| _ => None, |
| }) |
| } |
| InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => None, |
| }, |
| } |
| .filter(|_| forked_from_id.is_some()); |
| let parent_thread_id = session_configuration |
| .parent_thread_id |
| .or_else(|| initial_history.get_resumed_parent_thread_id()); |
| session_configuration.parent_thread_id = parent_thread_id; |
| if parent_thread_id.is_none() { |
| agent_control.set_root_service_tier( |
| session_configuration |
| .step_settings |
| .service_tier |
| .clone() |
| .or_else(|| config.service_tier.clone()), |
| ); |
| } |
| let is_paginated_subagent = matches!( |
| session_configuration.history_mode, |
| ThreadHistoryMode::Paginated |
| ) && matches!( |
| session_configuration.thread_source.as_ref(), |
| Some(ThreadSource::Subagent | ThreadSource::GuardianReview) |
| ); |
| if let InitialHistory::Forked(items) = &mut initial_history { |
| Self::assign_missing_rollout_response_item_ids(items); |
| } |
| let multi_agent_version = multi_agent_version.map(OnceLock::from).unwrap_or_default(); |
| let initial_multi_agent_version = multi_agent_version.get().copied(); |
|
|
| let thread_id = match (&initial_history, reserved_thread_id) { |
| ( |
| InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_), |
| Some(thread_id), |
| ) => thread_id, |
| (InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_), None) => { |
| agent_control.generate_thread_id() |
| } |
| (InitialHistory::Resumed(resumed_history), None) => resumed_history.conversation_id, |
| (InitialHistory::Resumed(_), Some(_)) => { |
| return Err(anyhow::anyhow!( |
| "reserved thread ID cannot be used when resuming a thread" |
| )); |
| } |
| }; |
| |
| let fork_cache_key = match &initial_history { |
| InitialHistory::Forked(items) |
| if config.ephemeral |
| && !session_configuration.session_source.is_non_root_agent() => |
| { |
| items.iter().find_map(|item| match item { |
| RolloutItem::SessionMeta(meta) => Some(meta.meta.session_id.to_string()), |
| _ => None, |
| }) |
| } |
| InitialHistory::New |
| | InitialHistory::Cleared |
| | InitialHistory::Resumed(_) |
| | InitialHistory::Forked(_) => None, |
| }; |
| let resumed_session_id = match &initial_history { |
| InitialHistory::Resumed(resumed) => { |
| resumed.history.iter().find_map(|item| match item { |
| RolloutItem::SessionMeta(meta_line) => Some(meta_line.meta.session_id), |
| _ => None, |
| }) |
| } |
| InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => None, |
| }; |
| |
| let resumed_session_id = resumed_session_id.filter(|session_id| { |
| !session_configuration.session_source.is_non_root_agent() |
| || *session_id != SessionId::from(thread_id) |
| }); |
| |
| let session_id = resumed_session_id.unwrap_or_else(|| { |
| if session_configuration.session_source.is_non_root_agent() { |
| agent_control.session_id() |
| } else { |
| SessionId::from(thread_id) |
| } |
| }); |
| let initial_auto_compact_window_ids = AutoCompactWindowIds::new_initial(); |
| let restore_child_window = matches!(&initial_history, InitialHistory::Forked(_)) |
| && session_configuration.session_source.is_non_root_agent() |
| && config.features.enabled(Feature::TokenBudget); |
| if restore_child_window && let InitialHistory::Forked(items) = &mut initial_history { |
| let child_window_id = initial_auto_compact_window_ids.window_id.to_string(); |
| for item in items { |
| if let RolloutItem::Compacted(checkpoint) = item { |
| checkpoint.window_number = Some(0); |
| checkpoint.first_window_id = Some(child_window_id.clone()); |
| checkpoint.previous_window_id = None; |
| checkpoint.window_id = Some(child_window_id.clone()); |
| } |
| } |
| } |
| let agent_control = agent_control.with_session_id( |
| session_id, |
| config |
| .effective_agent_max_threads(MultiAgentVersion::V2) |
| .unwrap_or(usize::MAX), |
| ); |
| let time_provider = crate::current_time::resolve_time_provider( |
| config.current_time_reminder.as_ref(), |
| external_time_provider, |
| )?; |
| let selected_capability_roots = |
| match thread_extension_init.get::<Vec<SelectedCapabilityRoot>>() { |
| Some(roots) => roots.as_ref().clone(), |
| None => { |
| let roots = initial_history.get_selected_capability_roots(); |
| if !roots.is_empty() { |
| thread_extension_init.insert(roots.clone()); |
| } |
| roots |
| } |
| }; |
| thread_extension_init.insert(codex_extension_api::ThreadOriginator( |
| session_configuration.originator.clone(), |
| )); |
| |
| |
| thread_extension_init.insert(model_info); |
| let isolation = thread_extension_init |
| .get::<codex_extension_api::SessionIsolation>() |
| .map(|policy| *policy) |
| .unwrap_or_default(); |
| let allowed_tools = thread_extension_init |
| .get::<codex_extension_api::AllowedTools>() |
| .or_else(|| { |
| |
| crate::guardian::is_basic_session_source(&session_configuration.session_source) |
| .then(|| Arc::new(codex_guardian_reviewer::reviewer_allowed_tools())) |
| }); |
| let mcp_thread_init = thread_extension_init.clone(); |
| let thread_extension_data = codex_extension_api::ExtensionData::new_with_init( |
| thread_id.to_string(), |
| thread_extension_init, |
| ); |
| |
| let guardian_context_mode = GuardianContextMode::from_features(&config.features); |
| thread_extension_data.insert(crate::context::GuardianReviewEvidence::default()); |
| |
| |
| |
| |
| |
| let thread_persistence_fut = async { |
| if config.ephemeral { |
| Ok::<_, anyhow::Error>((None, LiveThreadInitGuard::new( None))) |
| } else { |
| let mut local_guard = LiveThreadInitGuard::default(); |
| let mut managed_guard = match &startup { |
| Some(startup) => Some(startup.persistence.lock().await), |
| None => None, |
| }; |
| let guard = managed_guard.as_deref_mut().unwrap_or(&mut local_guard); |
| let live_thread = match &initial_history { |
| InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => { |
| let params = CreateThreadParams { |
| session_id, |
| thread_id, |
| extra_config: config.extra_config.clone(), |
| forked_from_id, |
| parent_thread_id, |
| source: session_source, |
| thread_source: session_configuration.thread_source.clone(), |
| originator: session_configuration.originator.clone(), |
| base_instructions: BaseInstructions { |
| text: session_configuration.base_instructions.clone(), |
| provenance: base_instructions_provenance.clone(), |
| }, |
| dynamic_tools: session_configuration.dynamic_tools.clone(), |
| selected_capability_roots: selected_capability_roots.clone(), |
| multi_agent_version: initial_multi_agent_version, |
| history_mode: session_configuration.history_mode, |
| history_base: match &fork_persistence { |
| ForkPersistence::Copied => None, |
| ForkPersistence::Referenced { history_base, .. } => *history_base, |
| }, |
| subagent_history_start_ordinal: None, |
| initial_window_id: initial_auto_compact_window_ids |
| .window_id |
| .to_string(), |
| runtime_workspace_roots: Some(config.workspace_roots.clone()), |
| metadata: ThreadPersistenceMetadata { |
| cwd: Some(config.cwd.to_path_buf()), |
| model_provider: config.model_provider_id.clone(), |
| memory_mode: if config.memories.generate_memories { |
| ThreadMemoryMode::Enabled |
| } else { |
| ThreadMemoryMode::Disabled |
| }, |
| }, |
| }; |
| if is_paginated_subagent |
| && matches!(&fork_persistence, ForkPersistence::Copied) |
| && let InitialHistory::Forked(items) = &initial_history |
| { |
| LiveThread::create_with_inherited_model_context( |
| Arc::clone(&thread_store), |
| params, |
| items, |
| guard, |
| ) |
| .await? |
| } else { |
| guard |
| .acquire(LiveThread::create(Arc::clone(&thread_store), params)) |
| .await? |
| } |
| } |
| InitialHistory::Resumed(resumed_history) => { |
| let params = ResumeThreadParams { |
| thread_id: resumed_history.conversation_id, |
| rollout_path: resumed_history.rollout_path.clone(), |
| history: Some(resumed_history.history.clone()), |
| include_archived: true, |
| metadata: ThreadPersistenceMetadata { |
| cwd: Some(config.cwd.to_path_buf()), |
| model_provider: config.model_provider_id.clone(), |
| memory_mode: if config.memories.generate_memories { |
| ThreadMemoryMode::Enabled |
| } else { |
| ThreadMemoryMode::Disabled |
| }, |
| }, |
| }; |
| guard |
| .acquire(LiveThread::resume( |
| Arc::clone(&thread_store), |
| session_configuration.history_mode, |
| params, |
| )) |
| .await? |
| } |
| }; |
| Ok((Some(live_thread), local_guard)) |
| } |
| } |
| .instrument(info_span!( |
| "session_init.thread_persistence", |
| otel.name = "session_init.thread_persistence", |
| session_init.ephemeral = config.ephemeral, |
| )); |
| let state_db_fut = async { |
| if config.ephemeral { |
| None |
| } else if let Some(local_store) = |
| thread_store.as_any().downcast_ref::<LocalThreadStore>() |
| { |
| local_store.state_db().await |
| } else { |
| None |
| } |
| } |
| .instrument(info_span!( |
| "session_init.state_db", |
| otel.name = "session_init.state_db", |
| session_init.ephemeral = config.ephemeral, |
| )); |
|
|
| let mut mcp_auth_changes = auth_manager.auth_change_receiver(); |
| let auth_manager_clone = Arc::clone(&auth_manager); |
| let plugins_manager_for_prewarm = Arc::clone(&plugins_manager); |
| let config_for_mcp = Arc::clone(&config); |
| let mcp_manager_for_mcp = Arc::clone(&mcp_manager); |
| let mcp_thread_init_for_startup = &mcp_thread_init; |
| let thread_extension_data_for_mcp = &thread_extension_data; |
| let mcp_originator = session_configuration.originator.clone(); |
| let mcp_session_source = session_configuration.session_source.clone(); |
| let mcp_disabled_plugin_ids = session_configuration.disabled_plugin_ids.clone(); |
| let mcp_runtime_cwd = environment_selections |
| .first() |
| .and_then(|environment| environment.cwd.to_abs_path().ok()) |
| .map(|cwd| cwd.to_path_buf()) |
| .unwrap_or_else(|| session_configuration.cwd().to_path_buf()); |
| let auth_and_mcp_fut = async move { |
| let auth = auth_manager_clone.auth().await; |
| if config_for_mcp.features.plugin_recommendations_enabled() { |
| let plugins_config = config_for_mcp.plugins_config_input(); |
| let auth_for_prewarm = auth.clone(); |
| |
| |
| tokio::spawn(async move { |
| plugins_manager_for_prewarm |
| .recommended_plugins_mode_for_config( |
| &plugins_config, |
| auth_for_prewarm.as_ref(), |
| ) |
| .await; |
| }); |
| } |
| let mcp_projection = mcp_manager_for_mcp |
| .runtime_config_for_step( |
| &config_for_mcp, |
| mcp_thread_init_for_startup, |
| thread_extension_data_for_mcp, |
| McpThreadIdentity { |
| session_source: &mcp_session_source, |
| originator: &mcp_originator, |
| disabled_plugin_ids: &mcp_disabled_plugin_ids, |
| environments: McpEnvironmentScope::Initial(environment_selections), |
| }, |
| &[], |
| None, |
| ) |
| .await; |
| (auth, mcp_projection) |
| } |
| .instrument(info_span!( |
| "session_init.auth_mcp", |
| otel.name = "session_init.auth_mcp", |
| )); |
|
|
| |
| let (thread_persistence_result, state_db_ctx, (auth, mcp_projection)) = |
| tokio::join!(thread_persistence_fut, state_db_fut, auth_and_mcp_fut); |
|
|
| let (live_thread, mut live_thread_init) = thread_persistence_result.map_err(|e| { |
| error!("failed to initialize thread persistence: {e:#}"); |
| e |
| })?; |
| let session_result: anyhow::Result<Arc<Self>> = async { |
| let rollout_path = if let Some(live_thread) = live_thread.as_ref() { |
| live_thread.local_rollout_path().await? |
| } else { |
| None |
| }; |
| let trace_agent_path = session_configuration |
| .session_source |
| .get_agent_path() |
| .unwrap_or_else(codex_protocol::AgentPath::root); |
| let trace_task_name = |
| (!trace_agent_path.is_root()).then(|| trace_agent_path.name().to_string()); |
| let trace_metadata = ThreadStartedTraceMetadata { |
| thread_id: thread_id.to_string(), |
| agent_path: trace_agent_path.to_string(), |
| task_name: trace_task_name, |
| nickname: session_configuration.session_source.get_nickname(), |
| agent_role: session_configuration.session_source.get_agent_role(), |
| session_source: session_configuration.session_source.clone(), |
| cwd: session_configuration.cwd().to_path_buf(), |
| rollout_path: rollout_path.clone(), |
| model: session_configuration.step_settings.collaboration_mode.model().to_string(), |
| provider_name: config.model_provider_id.clone(), |
| approval_policy: session_configuration.step_settings.approval_policy.value().to_string(), |
| sandbox_policy: format!( |
| "{:?}", |
| session_configuration.sandbox_policy(environment_selections) |
| ), |
| }; |
| let rollout_thread_trace = if matches!( |
| session_configuration.session_source, |
| SessionSource::SubAgent(SubAgentSource::ThreadSpawn { .. }) |
| ) { |
| |
| |
| |
| parent_rollout_thread_trace.start_child_thread_trace_or_disabled(trace_metadata) |
| } else { |
| ThreadTraceContext::start_root_or_disabled(trace_metadata) |
| }; |
|
|
| let mut post_session_configured_events = Vec::<Event>::new(); |
|
|
| for usage in config.features.legacy_feature_usages() { |
| post_session_configured_events.push(Event { |
| id: INITIAL_SUBMIT_ID.to_owned(), |
| msg: EventMsg::DeprecationNotice(DeprecationNoticeEvent { |
| summary: usage.summary.clone(), |
| details: usage.details.clone(), |
| }), |
| }); |
| } |
| for message in &config.startup_warnings { |
| post_session_configured_events.push(Event { |
| id: "".to_owned(), |
| msg: EventMsg::Warning(WarningEvent { |
| message: message.clone(), |
| }), |
| }); |
| } |
| let effective_config = config.config_layer_stack.effective_config(); |
| let config_path = config.codex_home.join(CONFIG_TOML_FILE); |
| if let Some(event) = unstable_features_warning_event( |
| effective_config.get("features").and_then(TomlValue::as_table), |
| config.suppress_unstable_features_warning, |
| &config.features, |
| &config_path.display().to_string(), |
| ) { |
| post_session_configured_events.push(event); |
| } |
| let telemetry_auth = auth.as_ref(); |
| let auth_mode = telemetry_auth |
| .map(CodexAuth::auth_mode) |
| .map(TelemetryAuthMode::from); |
| let account_id = telemetry_auth.and_then(CodexAuth::get_account_id); |
| let account_email = telemetry_auth.and_then(CodexAuth::get_account_email); |
| let originator = session_configuration.originator.clone(); |
| let terminal_type = user_agent(); |
| let session_model = session_configuration.step_settings.collaboration_mode.model().to_string(); |
| let auth_env_telemetry = collect_auth_env_telemetry( |
| session_configuration.provider.info(), |
| auth_manager.codex_api_key_env_enabled(), |
| ); |
| let mut session_telemetry = SessionTelemetry::new( |
| thread_id, |
| session_model.as_str(), |
| session_model.as_str(), |
| account_id.clone(), |
| account_email.clone(), |
| auth_mode, |
| originator.clone(), |
| config.otel.log_user_prompt, |
| terminal_type.clone(), |
| session_configuration.session_source.clone(), |
| ) |
| .with_auth_env(auth_env_telemetry.to_otel_metadata()) |
| .with_tool_result_log_config(config.otel.tool_result); |
| if let Some(service_name) = session_configuration.metrics_service_name.as_deref() { |
| session_telemetry = session_telemetry.with_metrics_service_name(service_name); |
| } |
| let network_proxy_audit_metadata = NetworkProxyAuditMetadata { |
| conversation_id: Some(thread_id.to_string()), |
| app_version: Some(env!("CARGO_PKG_VERSION").to_string()), |
| user_account_id: account_id, |
| auth_mode: auth_mode.map(|mode| mode.to_string()), |
| originator: Some(originator), |
| user_email: account_email, |
| terminal_type: Some(terminal_type), |
| model: Some(session_model.clone()), |
| slug: Some(session_model), |
| }; |
| crate::config::emit_session_start_metrics(config.as_ref(), &session_telemetry); |
| let is_worktree = session_configuration.cwd().canonicalize().ok().and_then(|cwd| { |
| codex_git_utils::repository_identity(&cwd).and_then(|_| { |
| get_git_repo_root(&cwd).map(|root| root.join(".git").is_file()) |
| }) |
| }); |
| let is_worktree_tag = match is_worktree { |
| Some(true) => "true", |
| Some(false) => "false", |
| None => "unknown", |
| }; |
| let is_git_tag = if get_git_repo_root(session_configuration.cwd()).is_some() { |
| "true" |
| } else { |
| "false" |
| }; |
| session_telemetry.counter( |
| THREAD_STARTED_METRIC, |
| 1, |
| &[("is_git", is_git_tag), ("is_worktree", is_worktree_tag)], |
| ); |
|
|
| let mcp_server_names = |
| codex_mcp::effective_mcp_servers( |
| &mcp_projection.config, |
| auth.as_ref(), |
| ) |
| .into_iter() |
| .filter_map(|(name, server)| server.enabled().then_some(name)) |
| .collect::<Vec<_>>(); |
| session_telemetry.conversation_starts( |
| config.model_provider.name.as_str(), |
| session_configuration.step_settings.collaboration_mode.reasoning_effort(), |
| config |
| .model_reasoning_summary |
| .unwrap_or(ReasoningSummaryConfig::Auto), |
| config.model_context_window, |
| config.model_auto_compact_token_limit, |
| config.permissions.approval_policy.value(), |
| config |
| .permissions |
| .legacy_sandbox_policy(session_configuration.cwd().as_path()), |
| mcp_server_names.iter().map(String::as_str).collect(), |
| ); |
|
|
| let use_zsh_fork_shell = config.features.enabled(Feature::ShellZshFork); |
| let default_shell = if let Some(user_shell_override) = |
| session_configuration.user_shell_override.clone() |
| { |
| user_shell_override |
| } else if use_zsh_fork_shell { |
| let zsh_path = config.zsh_path.as_ref().ok_or_else(|| { |
| anyhow::anyhow!( |
| "zsh fork feature enabled, but no packaged zsh fork is available for this install" |
| ) |
| })?; |
| if zsh_path.is_file() { |
| shell::Shell { |
| shell_type: shell::ShellType::Zsh, |
| shell_path: zsh_path.clone(), |
| } |
| } else { |
| shell::get_shell(shell::ShellType::Zsh).ok_or_else(|| { |
| anyhow::anyhow!( |
| "zsh fork feature enabled, but packaged zsh fork `{}` is not usable", |
| zsh_path.display() |
| ) |
| })? |
| } |
| } else { |
| shell::default_user_shell() |
| }; |
| let credential_broker_available = config.features.enabled(Feature::NetworkProxy) |
| && config |
| .config_layer_stack |
| .requirements() |
| .network |
| .as_ref() |
| .is_none_or(|network| network.value.enabled != Some(false)); |
| let credential_broker_configured = credential_broker_available |
| && effective_config |
| .get("features") |
| .and_then(|features| features.get("network_proxy")) |
| .and_then(|network_proxy| network_proxy.get("credential_broker")) |
| .and_then(TomlValue::as_bool) |
| .unwrap_or(false); |
| let credential_broker_active = credential_broker_configured |
| && config |
| .permissions |
| .network |
| .as_ref() |
| .is_some_and(crate::config::NetworkProxySpec::credential_broker_enabled); |
| let prefer_executor_shell_snapshots = config.features.enabled(Feature::ShellSnapshotV2) |
| && config.features.enabled(Feature::ShellTool) |
| && config.features.enabled(Feature::UnifiedExec) |
| && matches!( |
| codex_tools::UnifiedExecShellMode::for_session( |
| config.features.get(), |
| crate::tools::tool_user_shell_type(&default_shell), |
| config.zsh_path.as_ref(), |
| config.main_execve_wrapper_exe.as_ref(), |
| ), |
| codex_tools::UnifiedExecShellMode::Direct |
| ); |
| let use_executor_shell_snapshots = |
| prefer_executor_shell_snapshots && !credential_broker_active; |
| let shell_snapshot = if config.features.enabled(Feature::ShellSnapshot) |
| && (!use_executor_shell_snapshots || credential_broker_available) |
| { |
| let snapshot_credential_broker = credential_broker_available.then(|| { |
| let state = if credential_broker_active { |
| SnapshotCredentialBrokerState::Starting |
| } else { |
| SnapshotCredentialBrokerState::Inactive |
| }; |
| watch::channel(state).0 |
| }); |
| ShellSnapshot::new( |
| config.codex_home.clone(), |
| thread_id, |
| session_telemetry.clone(), |
| state_db_ctx.clone(), |
| snapshot_credential_broker, |
| prefer_executor_shell_snapshots, |
| ) |
| } else { |
| ShellSnapshot::disabled() |
| }; |
| let turn_environments = Arc::new(ThreadEnvironments::new( |
| environment_manager, |
| default_shell.clone(), |
| session_configuration.inferred_environment_config(), |
| shell_snapshot, |
| inherited_environments.unwrap_or_default(), |
| config.features.enabled(Feature::DeferredExecutor), |
| )); |
| turn_environments.update_selections( |
| environment_selections, |
| &session_configuration.inferred_environment_config(), |
| ); |
| let resolved_environments = turn_environments.snapshot().await; |
| let agents_md_manager = Arc::new(AgentsMdManager::new(instructions)); |
| let plugin_skill_warmup = warm_plugins_and_skills_for_session_init( |
| Arc::clone(&config), |
| Arc::clone(&plugins_manager), |
| Arc::clone(&skills_service), |
| &resolved_environments, |
| extensions.as_ref(), |
| ) |
| .instrument(info_span!( |
| "session_init.plugin_skill_warmup", |
| otel.name = "session_init.plugin_skill_warmup", |
| )); |
| let thread_name_lookup = |
| thread_title_from_thread_store(live_thread.as_ref(), &thread_store, thread_id) |
| .instrument(info_span!( |
| "session_init.thread_name_lookup", |
| otel.name = "session_init.thread_name_lookup", |
| )); |
| let (instruction_refresh, plugin_skill_errors, thread_name) = tokio::join!( |
| agents_md_manager.refresh(config.as_ref(), &resolved_environments), |
| plugin_skill_warmup, |
| thread_name_lookup, |
| ); |
| let (agents_md_result, instruction_warnings) = instruction_refresh; |
| |
| agents_md_result?; |
| post_session_configured_events.extend( |
| instruction_warnings.into_iter().map(|message| Event { |
| id: INITIAL_SUBMIT_ID.to_owned(), |
| msg: EventMsg::Warning(WarningEvent { message }), |
| }), |
| ); |
| for err in &plugin_skill_errors { |
| error!( |
| "failed to load skill {}: {}", |
| err.path.display(), |
| err.message |
| ); |
| } |
| session_configuration.thread_name = thread_name.clone(); |
| let mut state = SessionState::new_with_auto_compact_window_ids( |
| session_configuration.clone(), |
| initial_auto_compact_window_ids, |
| ContextManager::with_guardian_context_mode( |
| guardian_context_mode, |
| &session_configuration.session_source, |
| ), |
| ); |
| state.last_started_turn_id = initial_history.get_rollout_items().iter().rev().find_map(|item| { |
| match item { |
| RolloutItem::EventMsg(EventMsg::TurnStarted(event)) => Some(event.turn_id.clone()), |
| _ => None, |
| } |
| }); |
| state.base_instructions_provenance = base_instructions_provenance.clone(); |
| state.active_disabled_plugin_ids = session_configuration.disabled_plugin_ids.clone(); |
| let managed_network_requirements_configured = config |
| .config_layer_stack |
| .requirements_toml() |
| .network |
| .is_some(); |
| let managed_network_requirements_enabled = config.managed_network_requirements_enabled(); |
| let network_approval = Arc::new(NetworkApprovalService::default()); |
| |
| let network_policy_decider_session = if managed_network_requirements_configured { |
| config |
| .permissions |
| .network |
| .as_ref() |
| .map(|_| Arc::new(RwLock::new(std::sync::Weak::<Session>::new()))) |
| } else { |
| None |
| }; |
| let blocked_request_observer = config |
| .permissions |
| .network |
| .as_ref() |
| .map(|_| build_blocked_request_observer(Arc::clone(&network_approval))); |
| let network_policy_decider = |
| network_policy_decider_session |
| .as_ref() |
| .map(|network_policy_decider_session| { |
| build_network_policy_decider( |
| Arc::clone(&network_approval), |
| Arc::clone(network_policy_decider_session), |
| ) |
| }); |
| let (network_proxy, session_network_proxy) = |
| if let Some(spec) = config |
| .permissions |
| .network |
| .as_ref() |
| .filter(|spec| spec.enabled()) |
| { |
| let current_exec_policy = exec_policy.current(); |
| let (network_proxy, session_network_proxy) = Self::start_managed_network_proxy( |
| spec, |
| current_exec_policy.as_ref(), |
| config.permissions.permission_profile(), |
| config.permissions.windows_sandbox_type, |
| network_policy_decider.as_ref().map(Arc::clone), |
| blocked_request_observer.as_ref().map(Arc::clone), |
| managed_network_requirements_configured, |
| network_proxy_audit_metadata.clone(), |
| ) |
| .instrument(info_span!( |
| "session_init.network_proxy", |
| otel.name = "session_init.network_proxy", |
| session_init.managed_network_requirements_enabled = |
| managed_network_requirements_enabled, |
| )) |
| .await?; |
| (Some(network_proxy), Some(session_network_proxy)) |
| } else { |
| (None, None) |
| }; |
| if let Some(network_proxy) = network_proxy.as_ref() |
| && config |
| .permissions |
| .network |
| .as_ref() |
| .is_some_and(crate::config::NetworkProxySpec::credential_broker_enabled) |
| { |
| turn_environments.set_snapshot_credential_broker( |
| SnapshotCredentialBrokerState::Ready(network_proxy.proxy()), |
| ); |
| } |
|
|
| |
| let mcp_runtime = Arc::new(McpRuntime::empty( |
| mcp_projection.config.prefix_mcp_tool_names, |
| )); |
| let hooks_config = build_hooks_config( |
| &config, |
| plugins_manager.as_ref(), |
| resolved_environments.single_local_environment(), |
| &session_configuration.disabled_plugin_ids, |
| ) |
| .await; |
| let (hooks, async_hook_results) = Hooks::new( |
| hooks_config, |
| thread_id, |
| Arc::new(CoreHookMcpExecutor { |
| runtime: Arc::clone(&mcp_runtime), |
| thread_id, |
| }), |
| )?; |
| for warning in hooks.startup_warnings() { |
| post_session_configured_events.push(Event { |
| id: INITIAL_SUBMIT_ID.to_owned(), |
| msg: EventMsg::Warning(WarningEvent { |
| message: warning.clone(), |
| }), |
| }); |
| } |
|
|
| let analytics_events_client = if config.analytics_enabled == Some(false) { |
| AnalyticsEventsClient::disabled() |
| } else { |
| analytics_events_client.unwrap_or_else(|| { |
| AnalyticsEventsClient::new( |
| Arc::clone(&auth_manager), |
| config.chatgpt_base_url.trim_end_matches('/').to_string(), |
| config.analytics_enabled, |
| ) |
| }) |
| }; |
| for item in initial_history.get_rollout_items() { |
| match item { |
| RolloutItem::Compacted(compacted) => { |
| if let Some(checkpoint) = &compacted.mcp_resource_origins { |
| mcp_runtime.restore_resource_origin_checkpoint(checkpoint); |
| } |
| } |
| RolloutItem::EventMsg(event) => mcp_runtime.observe_event(event), |
| RolloutItem::SessionMeta(_) |
| | RolloutItem::ResponseItem(_) |
| | RolloutItem::InterAgentCommunication(_) |
| | RolloutItem::InterAgentCommunicationMetadata { .. } |
| | RolloutItem::TurnContext(_) |
| | RolloutItem::WorldState(_) |
| | RolloutItem::RealtimeItem(_) |
| | RolloutItem::TokenUsageRecord(_) |
| | RolloutItem::RetainedContext(_) |
| | RolloutItem::SecurityRiskScore(_) => {} |
| } |
| } |
| let session_extension_data = |
| codex_extension_api::ExtensionData::new(session_id.to_string()); |
| session_extension_data.insert(analytics_events_client.clone()); |
| let mcp_resource_client = Arc::new(McpResourceClient::new(Arc::clone(&mcp_runtime))); |
| let extension_metrics = |
| extension_metrics::from_session_telemetry(session_telemetry.clone()); |
| for contributor in extensions.thread_lifecycle_contributors() { |
| contributor.on_thread_start(codex_extension_api::ThreadStartInput { |
| config: config.as_ref(), |
| session_source: &session_configuration.session_source, |
| persistent_thread_state_available: state_db_ctx.is_some(), |
| environments: environment_selections, |
| mcp_resource_client: Some(Arc::clone(&mcp_resource_client)), |
| extension_metrics: Some(Arc::clone(&extension_metrics)), |
| session_store: &session_extension_data, |
| thread_store: &thread_extension_data, |
| }).await; |
| } |
|
|
| let executed_tool_calls = crate::state::ExecutedToolCalls::new( |
| &config.features, |
| &initial_history, |
| ); |
| let codex_responses_headers = thread_extension_data.get::<crate::CodexResponsesHeaders>(); |
| let services = SessionServices { |
| |
| |
| mcp_runtime, |
| mcp_handler_cache: Default::default(), |
| unified_exec_manager: UnifiedExecProcessManager::new( |
| config.background_terminal_max_timeout, |
| ), |
| elicitations: crate::elicitation::ElicitationService::new(), |
| shell_zsh_path: config.zsh_path.clone(), |
| main_execve_wrapper_exe: config.main_execve_wrapper_exe.clone(), |
| analytics_events_client, |
| hooks: arc_swap::ArcSwap::from_pointee(hooks), |
| rollout_thread_trace, |
| user_shell: Arc::new(default_shell), |
| show_raw_agent_reasoning: config.show_raw_agent_reasoning, |
| exec_policy, |
| auth_manager: Arc::clone(&auth_manager), |
| openai_file_upload_client_pool: RouteAwareClientPool::new_without_request_logging( |
| config.http_client_factory(), |
| ClientRouteClass::Api, |
| ) |
| .with_legacy_custom_ca_fallback(), |
| session_telemetry, |
| models_manager: Arc::clone(&models_manager), |
| git_root_discovery, |
| tool_approvals: Mutex::new(ApprovalStore::default()), |
| runtime_handle: tokio::runtime::Handle::current(), |
| skills_service, |
| agents_md_manager, |
| plugins_manager: Arc::clone(&plugins_manager), |
| mcp_manager: Arc::clone(&mcp_manager), |
| extensions, |
| |
| session_extension_data, |
| thread_extension_data, |
| selected_capability_roots, |
| mcp_thread_init, |
| client_mcp_extensions, |
| agent_control, |
| network_proxy: arc_swap::ArcSwapOption::from(network_proxy.map(Arc::new)), |
| network_proxy_audit_metadata, |
| managed_network_requirements_configured, |
| network_approval: Arc::clone(&network_approval), |
| state_db: state_db_ctx.clone(), |
| live_thread: live_thread.clone(), |
| image_store, |
| thread_store: Arc::clone(&thread_store), |
| attestation_provider: attestation_provider.clone(), |
| time_provider, |
| model_client: ModelClient::new( |
| Some(Arc::clone(&auth_manager)), |
| if config.features.enabled(Feature::UseAgentIdentity) { |
| AgentIdentityAuthPolicy::ChatGptAuth |
| } else { |
| AgentIdentityAuthPolicy::JwtOnly |
| }, |
| thread_id, |
| session_configuration.provider.info().clone(), |
| session_configuration.session_source.clone(), |
| session_configuration.originator.clone(), |
| config.model_verbosity, |
| config.features.enabled(Feature::ContentItemKinds), |
| config.features.enabled(Feature::EnableRequestCompression), |
| config.features.enabled(Feature::RuntimeMetrics), |
| Self::build_model_client_beta_features_header(config.as_ref()), |
| config |
| .features |
| .enabled(Feature::ConcurrentReasoningSummaries), |
| attestation_provider, |
| config.http_client_factory(), |
| config.workspace_routing_context(), |
| ) |
| .with_session_context( |
| crate::guardian::prompt_cache_key_override_for_review_session( |
| &session_configuration.session_source, |
| session_configuration.parent_thread_id, |
| ) |
| .or(fork_cache_key), |
| tx_event.clone(), |
| codex_responses_headers, |
| ), |
| executed_tool_calls: executed_tool_calls.clone(), |
| code_mode_service: crate::tools::code_mode::CodeModeService::new( |
| thread_id, |
| Arc::clone(&code_mode_session_provider), |
| &config.code_mode, |
| executed_tool_calls, |
| ), |
| tool_search_handler_cache: Default::default(), |
| turn_environments: Arc::clone(&turn_environments), |
| }; |
| let (mcp_prewarm_tx, mcp_prewarm_rx) = async_channel::bounded(1); |
| let sess = Arc::new(Session { |
| thread_id, |
| installation_id, |
| tx_event: tx_event.clone(), |
| agent_status, |
| state: Mutex::new(state), |
| thread_settings_persistence: Semaphore::new( 1), |
| managed_network_proxy_refresh_lock: Semaphore::new( 1), |
| features: config.features.clone(), |
| guardian_context_mode, |
| isolation, |
| allowed_tools, |
| windows_sandbox_proxy_settings_mode, |
| multi_agent_version, |
| mcp_refresh: McpRefresh::new(), |
| mcp_tool_approval_metadata: Default::default(), |
| mcp_elicitation_reviewer_handle: OnceLock::new(), |
| mcp_elicitation_lifecycle_handle: OnceLock::new(), |
| mcp_prewarm_tx, |
| mcp_prewarm_shutdown: CancellationToken::new(), |
| mcp_prewarm_task: std::sync::Mutex::new(None), |
| conversation: Arc::new(RealtimeConversationManager::new()), |
| realtime_history: (session_configuration.history_mode == ThreadHistoryMode::Paginated |
| && services.live_thread.is_some()) |
| .then(|| Mutex::new(Default::default())), |
| active_turn: Mutex::new(None), |
| async_hook_results, |
| input_queue: InputQueue::new(), |
| services, |
| git_enrichment_policy, |
| fork_persistence, |
| forked_from_ordinal_exclusive, |
| next_internal_sub_id: AtomicU64::new(0), |
| }); |
| if let Some(startup) = &startup { |
| let _ = startup.session.set(Arc::clone(&sess)); |
| } |
| if let Some(network_policy_decider_session) = network_policy_decider_session { |
| let mut guard = network_policy_decider_session.write().await; |
| *guard = Arc::downgrade(&sess); |
| } |
| |
| |
| let initial_messages = initial_history.get_event_msgs(); |
| let thread_config = |
| session_configuration.thread_config_snapshot(turn_environments.selections()); |
| let events = std::iter::once(Event { |
| id: INITIAL_SUBMIT_ID.to_owned(), |
| msg: EventMsg::SessionConfigured(SessionConfiguredEvent { |
| session_id, |
| thread_id, |
| cwd: thread_config.cwd().clone(), |
| forked_from_id: thread_config.forked_from_thread_id, |
| parent_thread_id: thread_config.parent_thread_id, |
| thread_source: thread_config.thread_source, |
| thread_name: session_configuration.thread_name.clone(), |
| model: thread_config.model, |
| model_provider_id: thread_config.model_provider_id, |
| service_tier: thread_config.service_tier, |
| approval_policy: thread_config.approval_policy, |
| approvals_reviewer: thread_config.approvals_reviewer, |
| network_proxy: session_network_proxy.filter(|_| { |
| Self::managed_network_proxy_active_for_permission_profile( |
| &thread_config.permission_profile, |
| ) |
| }), |
| permission_profile: thread_config.permission_profile, |
| active_permission_profile: thread_config.active_permission_profile, |
| reasoning_effort: thread_config.reasoning_effort, |
| initial_messages, |
| rollout_path, |
| }), |
| }) |
| .chain(post_session_configured_events.into_iter()); |
| for event in events { |
| sess.send_event_raw(event).await; |
| } |
| turn_environments.start_connection_event_forwarding(tx_event.clone()); |
|
|
| let startup_auth_changed = mcp_auth_changes.has_changed().unwrap_or(false); |
| if startup_auth_changed { |
| mcp_auth_changes.mark_unchanged(); |
| } |
| let latest_auth = sess.services.auth_manager.auth().await; |
| let mcp_projection = if startup_auth_changed |
| || mcp_auth_changes.has_changed().unwrap_or(false) |
| { |
| sess.services |
| .mcp_manager |
| .runtime_config_for_step( |
| config.as_ref(), |
| &sess.services.mcp_thread_init, |
| &sess.services.thread_extension_data, |
| McpThreadIdentity { |
| session_source: &session_configuration.session_source, |
| originator: &session_configuration.originator, |
| disabled_plugin_ids: &session_configuration.disabled_plugin_ids, |
| environments: McpEnvironmentScope::Live( |
| &sess.services.turn_environments, |
| ), |
| }, |
| &[], |
| None, |
| ) |
| .await |
| } else { |
| mcp_projection |
| }; |
| sess.install_initial_mcp_runtime( |
| &session_configuration, |
| latest_auth, |
| mcp_projection, |
| &resolved_environments, |
| mcp_runtime_cwd, |
| ) |
| .await?; |
| sess.start_mcp_prewarm_worker(mcp_prewarm_rx, mcp_auth_changes); |
| sess.schedule_startup_prewarm(sess.get_prompt_base_instructions().await.text) |
| .await; |
| let session_start_source = match &initial_history { |
| InitialHistory::Forked(_) if forked_from_id.is_some() => { |
| codex_hooks::SessionStartSource::Fork |
| } |
| |
| |
| InitialHistory::Resumed(_) | InitialHistory::Forked(_) => { |
| codex_hooks::SessionStartSource::Resume |
| } |
| InitialHistory::New => codex_hooks::SessionStartSource::Startup, |
| InitialHistory::Cleared => codex_hooks::SessionStartSource::Clear, |
| }; |
|
|
| |
| Box::pin(sess.record_initial_history(initial_history)).await; |
| if restore_child_window { |
| sess.state.lock().await.restore_auto_compact_window( |
| 0, |
| initial_auto_compact_window_ids, |
| ); |
| } |
| if matches!(&sess.fork_persistence, ForkPersistence::Referenced { .. }) { |
| |
| sess.try_ensure_rollout_materialized(PersistContext::Standard) |
| .await?; |
| } |
| { |
| let mut state = sess.state.lock().await; |
| state.queue_pending_session_start_source(session_start_source); |
| } |
| Ok(sess) |
| } |
| .await; |
| match session_result { |
| Ok(sess) => { |
| live_thread_init.commit(); |
| Ok(sess) |
| } |
| Err(err) => { |
| live_thread_init.discard().await; |
| Err(err) |
| } |
| } |
| } |
| } |
|
|