codex / codex-rs /core /src /session /session.rs
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
52a9af3 verified
Raw
History Blame Contribute Delete
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,
}
#[derive(Clone)]
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());
}
#[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(),
}
}
/// 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(&current_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(),
&current_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(
&current_sandbox_policy,
self.cwd(),
&current_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(
&current_sandbox_policy,
next_configuration.cwd(),
&current_file_system_sandbox_policy,
);
next_configuration
.permission_profile_state
.set_legacy_permission_profile(
PermissionProfile::from_runtime_permissions_with_enforcement(
SandboxEnforcement::from_legacy_sandbox_policy(&current_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,
}
#[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 {
/// 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
}
}
#[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) => {
// 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)
}
}
}
}