Download codex-rs/core/src/session/session.rs from SaylorTwift/codex: direct link, hf CLI and curl.
- Browser
- Download file 89.6 kB
-
https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/core/src/session/session.rs
- Command line
-
hf download hf://SaylorTwift/codex/codex-rs/core/src/session/session.rs
-
curl -L -o session.rs https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/core/src/session/session.rs
89.6 kB
| 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)>>; | |
| /// Context for an initialized model agent | |
| /// | |
| /// A session has at most 1 running task at a time, and can be interrupted by user input. | |
| 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>, | |
| /// Orders accepted settings commits and their persisted events with compaction checkpoints. | |
| /// Keep this separate from `state` so storage I/O does not block runtime state access. | |
| pub(super) thread_settings_persistence: Semaphore, | |
| /// Serializes rebuild/apply cycles for the running proxy; each cycle | |
| /// rebuilds from the current SessionState while holding this lock. | |
| pub(super) managed_network_proxy_refresh_lock: Semaphore, | |
| /// The set of enabled features should be invariant for the lifetime of the | |
| /// session. | |
| 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>, | |
| /// Owns invalidation and serializes refreshes without blocking captured calls. | |
| pub(super) mcp_refresh: McpRefresh, | |
| /// Non-owning lookup for approval data retained by running MCP invocations. | |
| 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, | |
| } | |
| pub(crate) struct SessionConfiguration { | |
| /// Runtime provider and its provider-specific execution policy. | |
| pub(super) provider: SharedModelProvider, | |
| /// Desired configured inputs inherited by future turns. | |
| pub(super) step_settings: Arc<StepSettings>, | |
| /// Explicit startup overrides used when resolving effective model metadata. | |
| pub(super) model_info_overrides: ModelInfoOverrides, | |
| /// Developer instructions that supplement the base instructions. | |
| pub(super) developer_instructions: Option<String>, | |
| /// Base instructions for the session. | |
| pub(super) base_instructions: String, | |
| /// Permission profile state for the session. Keep the constrained profile, | |
| /// active profile id, and profile-defined workspace roots in sync by using | |
| /// the methods below instead of mutating the fields independently. | |
| pub(super) permission_profile_state: PermissionProfileState, | |
| pub(super) allow_login_shell: bool, | |
| pub(super) shell_environment_policy: ShellEnvironmentPolicy, | |
| // TODO(anp): Reconcile these legacy thread defaults with TurnEnvironment::sandbox_context; | |
| // internal sandbox decisions should use the selected environment's configuration. | |
| 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, | |
| /// Legacy thread cwd used when a turn does not select an environment. | |
| pub(super) legacy_fallback_cwd: AbsolutePathBuf, | |
| /// Top-level runtime workspace roots, independent of explicit environment selections. | |
| pub(super) runtime_workspace_roots: Vec<AbsolutePathBuf>, | |
| /// Directory containing all Codex state for this session. | |
| pub(super) codex_home: AbsolutePathBuf, | |
| /// Optional user-facing name for the thread, updated during the session. | |
| pub(super) thread_name: Option<String>, | |
| /// Thread-owned plugin selection inherited by future turns. | |
| pub(super) disabled_plugin_ids: Vec<String>, | |
| // TODO(pakrym): Remove config from here | |
| pub(super) original_config_do_not_use: Arc<Config>, | |
| /// Optional service name tag for session metrics. | |
| pub(super) metrics_service_name: Option<String>, | |
| pub(super) app_server_client_name: Option<String>, | |
| pub(super) app_server_client_version: Option<String>, | |
| /// Guardian reviewer identity is trusted only when established during an in-memory spawn. | |
| pub(super) trusted_guardian_reviewer: bool, | |
| /// Source of the session (cli, vscode, exec, mcp, ...) | |
| pub(super) session_source: SessionSource, | |
| /// Persisted thread history contract selected when this thread was created. | |
| pub(super) history_mode: ThreadHistoryMode, | |
| /// Immediate history source copied into this thread, when this thread was forked. | |
| pub(super) forked_from_thread_id: Option<ThreadId>, | |
| /// Immediate control/spawn parent for this thread, when it has one. | |
| pub(super) parent_thread_id: Option<ThreadId>, | |
| /// Optional analytics source classification for this thread. | |
| pub(super) thread_source: Option<ThreadSource>, | |
| /// Effective originator used for this thread's Responses requests and analytics events. | |
| 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()); | |
| } | |
| 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(), | |
| } | |
| } | |
| /// Captures thread-owned settings for persistence and resume. | |
| 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(), | |
| } | |
| } | |
| /// Captures thread-owned settings and their separately owned environments. | |
| 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 { | |
| // Compatibility projection can resolve filesystem paths. Only compute it | |
| // when a cwd-bound legacy policy might need rebinding. | |
| 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(), | |
| ) { | |
| // Preserve richer split policies across cwd-only updates; only | |
| // rederive when the session is already using a structurally | |
| // cwd-bound legacy bridge. | |
| 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)?; | |
| // Apply step settings last: the proposed permissions and environment | |
| // selections must be complete before deriving their validation constraints. | |
| 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) | |
| } | |
| } | |
| /// The configuration and public snapshot published by one settings commit. | |
| pub(crate) struct SessionSettingsCommit { | |
| pub(crate) configuration: SessionConfiguration, | |
| pub(crate) snapshot: ThreadSettingsSnapshot, | |
| } | |
| 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 { | |
| /// Returns the concrete identity for this thread. | |
| pub(crate) fn thread_id(&self) -> ThreadId { | |
| self.thread_id | |
| } | |
| /// Returns the identity shared by the root thread and all descendant threads. | |
| 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, | |
| ) | |
| } | |
| // TODO(CDXENT-454): Build the compaction request and metadata from the captured execution. | |
| // Remote compaction currently attaches only finalized tool inventory because the rest of the | |
| // request remains turn-backed; local compaction does not have a finalized request inventory. | |
| 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 | |
| } | |
| } | |
| 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) => { | |
| // Both local and CCA thread stores place the resumed thread's | |
| // canonical SessionMeta first. Never inspect inherited metadata: | |
| // an ancestor's history_base describes a different fork boundary. | |
| 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" | |
| )); | |
| } | |
| }; | |
| // Ephemeral forks reuse cache routing, without sharing storage or lifecycle identity. | |
| 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, | |
| }; | |
| // Legacy subagent rollouts synthesized session_id from their own thread ID. | |
| let resumed_session_id = resumed_session_id.filter(|session_id| { | |
| !session_configuration.session_source.is_non_root_agent() | |
| || *session_id != SessionId::from(thread_id) | |
| }); | |
| // session_id is equal to the root thread's 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(), | |
| )); | |
| // Publish the already resolved model before extensions make startup decisions. | |
| // Turn construction refreshes this attachment when the selected model changes. | |
| 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(|| { | |
| // Older reviewer rollouts predate the explicit startup setting. | |
| 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, | |
| ); | |
| // Capture follows the flag; replay selects reviewer policy from the saved checkpoint. | |
| let guardian_context_mode = GuardianContextMode::from_features(&config.features); | |
| thread_extension_data.insert(crate::context::GuardianReviewEvidence::default()); | |
| // Kick off independent async setup tasks in parallel to reduce startup latency. | |
| // | |
| // - initialize thread persistence with new or resumed session info | |
| // - perform default shell discovery | |
| // - load history metadata (skipped for subagents) | |
| let thread_persistence_fut = async { | |
| if config.ephemeral { | |
| Ok::<_, anyhow::Error>((None, LiveThreadInitGuard::new(/*live_thread*/ 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(); | |
| // Fetch the catalog while MCP and plugin/skill initialization continue. | |
| // Context construction still handles filtering and prompt insertion. | |
| 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), | |
| }, | |
| /*ready_selected_capability_roots*/ &[], | |
| /*executor_capability_discovery*/ None, | |
| ) | |
| .await; | |
| (auth, mcp_projection) | |
| } | |
| .instrument(info_span!( | |
| "session_init.auth_mcp", | |
| otel.name = "session_init.auth_mcp", | |
| )); | |
| // Join all independent futures. | |
| 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 { .. }) | |
| ) { | |
| // Spawned child threads are part of their root rollout tree. If the | |
| // parent had no trace bundle, do not create an orphan child bundle | |
| // that looks like an independent rollout. | |
| 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, | |
| /*inc*/ 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; | |
| // TODO(anp): Present AGENTS.md discovery errors more clearly to the user. | |
| 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()); | |
| // The managed proxy can call back into core for allowlist-miss decisions. | |
| 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()), | |
| ); | |
| } | |
| // Hooks and extensions share one stable thread-owned MCP runtime handle. | |
| 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 { | |
| // Start with an empty connection set. The initialized set is | |
| // published after SessionConfigured so MCP events follow it. | |
| 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, | |
| // TODO(jif): extract session to share between sub-agents | |
| 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()), | |
| /*concurrent_reasoning_summaries_enabled*/ 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(/*permits*/ 1), | |
| managed_network_proxy_refresh_lock: Semaphore::new(/*permits*/ 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); | |
| } | |
| // Dispatch the SessionConfiguredEvent first and then report any errors. | |
| // If resuming, include converted initial messages in the payload so UIs can render them immediately. | |
| 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, | |
| ), | |
| }, | |
| /*ready_selected_capability_roots*/ &[], | |
| /*executor_capability_discovery*/ 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 | |
| } | |
| // `thread/resume` with supplied history uses `Forked` internally | |
| // without a fork parent, so it should still report `resume`. | |
| InitialHistory::Resumed(_) | InitialHistory::Forked(_) => { | |
| codex_hooks::SessionStartSource::Resume | |
| } | |
| InitialHistory::New => codex_hooks::SessionStartSource::Startup, | |
| InitialHistory::Cleared => codex_hooks::SessionStartSource::Clear, | |
| }; | |
| // record_initial_history can emit events. We record only after the SessionConfiguredEvent is emitted. | |
| Box::pin(sess.record_initial_history(initial_history)).await; | |
| if restore_child_window { | |
| sess.state.lock().await.restore_auto_compact_window( | |
| /*window_number*/ 0, | |
| initial_auto_compact_window_ids, | |
| ); | |
| } | |
| if matches!(&sess.fork_persistence, ForkPersistence::Referenced { .. }) { | |
| // Keep the source reserved until the child's history reference is durable. | |
| 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) | |
| } | |
| } | |
| } | |
| } | |