use super::thread_input::ensure_direct_input_allowed; use super::*; use codex_agent_extension::AgentInvocation; use codex_agent_extension::AgentRun; use codex_agent_extension::AgentRunner; use codex_app_server_protocol::ImageReference as V2ImageReference; use codex_protocol::error::CodexErrorDetails; use codex_protocol::models::ContentItem; use codex_protocol::models::FunctionCallOutputBody; use codex_protocol::models::FunctionCallOutputContentItem; use codex_protocol::models::FunctionCallOutputPayload; use codex_protocol::models::ImageReference as CoreImageReference; use codex_protocol::protocol::AdditionalContextEntry as CoreAdditionalContextEntry; use codex_protocol::protocol::AdditionalContextKind as CoreAdditionalContextKind; use codex_protocol::protocol::TurnSettingsUpdate; use codex_protocol::protocol::TurnSettingsUpdateOutcome; use codex_skills::system_cache_root_dir; use crate::image_url::REMOTE_IMAGE_URL_ERROR; use crate::image_url::is_remote_image_url; pub(super) fn validate_user_input_image_urls( input: &[V2UserInput], ) -> Result<(), JSONRPCErrorError> { if input.iter().any(|item| { matches!( item, V2UserInput::Image { image: V2ImageReference::Inline { url }, .. } if is_remote_image_url(url) ) }) { return Err(invalid_request(REMOTE_IMAGE_URL_ERROR)); } Ok(()) } fn validate_response_item_image_urls(items: &[ResponseItem]) -> Result<(), JSONRPCErrorError> { if items.iter().any(|item| match item { ResponseItem::Message { content, .. } => content.iter().any(|item| { matches!( item, ContentItem::InputImage { image: CoreImageReference::Inline { image_url }, .. } if is_remote_image_url(image_url) ) }), ResponseItem::FunctionCallOutput { output, .. } | ResponseItem::CustomToolCallOutput { output, .. } => { output.content_items().is_some_and(|content| { content.iter().any(|item| { matches!( item, FunctionCallOutputContentItem::InputImage { image: CoreImageReference::Inline { image_url }, .. } if is_remote_image_url(image_url) ) }) }) } ResponseItem::Reasoning { .. } | ResponseItem::AgentMessage { .. } | ResponseItem::LocalShellCall { .. } | ResponseItem::FunctionCall { .. } | ResponseItem::ToolSearchCall { .. } | ResponseItem::CustomToolCall { .. } | ResponseItem::ToolSearchOutput { .. } | ResponseItem::WebSearchCall { .. } | ResponseItem::ImageGenerationCall { .. } | ResponseItem::Compaction { .. } | ResponseItem::ConfigurationUpdate { .. } | ResponseItem::CompactionTrigger { .. } | ResponseItem::ContextCompaction { .. } | ResponseItem::AdditionalTools { .. } | ResponseItem::Other => false, }) { return Err(invalid_request(REMOTE_IMAGE_URL_ERROR)); } Ok(()) } #[derive(Clone)] pub(crate) struct TurnRequestProcessor { agent_runner: AgentRunner, auth_manager: Arc, thread_manager: Arc, outgoing: Arc, analytics_events_client: AnalyticsEventsClient, arg0_paths: Arg0DispatchPaths, config: Arc, config_manager: ConfigManager, pending_thread_unloads: Arc>>, thread_state_manager: ThreadStateManager, thread_watch_manager: ThreadWatchManager, skills_watcher: Arc, turn_cost_worker: Option, } fn map_additional_context( additional_context: Option>, ) -> BTreeMap { additional_context .unwrap_or_default() .into_iter() .map(|(key, entry)| { ( key, CoreAdditionalContextEntry { value: entry.value, kind: match entry.kind { AdditionalContextKind::Untrusted => CoreAdditionalContextKind::Untrusted, AdditionalContextKind::Application => { CoreAdditionalContextKind::Application } }, }, ) }) .collect() } #[derive(Default)] struct ThreadEnvironmentOverride { environments: Option, // Only default-environment updates replace the task's separately persisted root selection. runtime_workspace_roots: Option>, } struct ThreadSettingsBuildParams { method: &'static str, disabled_plugin_ids: Option>, environment_override: ThreadEnvironmentOverride, approval_policy: Option, approvals_reviewer: Option, sandbox_policy: Option, permissions: Option, model: Option, service_tier: Option>, effort: Option, summary: Option, collaboration_mode: Option, personality: Option, } impl TurnRequestProcessor { #[allow(clippy::too_many_arguments)] pub(crate) fn new( auth_manager: Arc, thread_manager: Arc, outgoing: Arc, analytics_events_client: AnalyticsEventsClient, arg0_paths: Arg0DispatchPaths, config: Arc, config_manager: ConfigManager, pending_thread_unloads: Arc>>, thread_state_manager: ThreadStateManager, thread_watch_manager: ThreadWatchManager, skills_watcher: Arc, turn_cost_worker: Option, ) -> Self { let agent_runner = AgentRunner::new(Arc::downgrade(&thread_manager)); Self { agent_runner, auth_manager, thread_manager, outgoing, analytics_events_client, arg0_paths, config, config_manager, pending_thread_unloads, thread_state_manager, thread_watch_manager, skills_watcher, turn_cost_worker, } } pub(crate) async fn turn_start( &self, request_id: ConnectionRequestId, params: TurnStartParams, app_server_client_name: Option, app_server_client_version: Option, ) -> Result, JSONRPCErrorError> { validate_user_input_image_urls(¶ms.input)?; self.turn_start_inner( request_id, params, app_server_client_name, app_server_client_version, ) .await .map(|response| Some(response.into())) } pub(crate) async fn thread_inject_items( &self, request_id: &ConnectionRequestId, params: ThreadInjectItemsParams, ) -> Result, JSONRPCErrorError> { self.thread_inject_items_response_inner(request_id, params) .await .map(|response| Some(response.into())) } pub(crate) async fn thread_settings_update( &self, request_id: &ConnectionRequestId, params: ThreadSettingsUpdateParams, ) -> Result, JSONRPCErrorError> { self.thread_settings_update_inner(request_id, params) .await .map(|response| Some(response.into())) } pub(crate) async fn turn_settings_update( &self, request_id: &ConnectionRequestId, params: TurnSettingsUpdateParams, ) -> Result, JSONRPCErrorError> { let (_, thread) = self.load_thread(¶ms.thread_id).await?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; let (reply, outcome) = oneshot::channel(); self.submit_core_op( request_id, &thread, Op::TurnSettings { turn_id: params.turn_id, update: TurnSettingsUpdate { approvals_reviewer: params .approvals_reviewer .map(codex_app_server_protocol::ApprovalsReviewer::to_core), model: params.model, // Match thread/settings/update: public null does not clear effort. effort: params.effort.map(Some), summary: params.summary, service_tier: params.service_tier, }, reply, }, ) .await .map_err(|err| internal_error(format!("failed to submit turn settings: {err}")))?; let outcome = outcome .await .map_err(|_| internal_error("turn settings operation ended before replying"))?; let status = match outcome { TurnSettingsUpdateOutcome::Applied => TurnSettingsUpdateStatus::Applied, TurnSettingsUpdateOutcome::TargetUnavailable => { TurnSettingsUpdateStatus::TargetUnavailable } TurnSettingsUpdateOutcome::Rejected { reason } => return Err(invalid_request(reason)), }; Ok(Some(TurnSettingsUpdateResponse { status }.into())) } pub(crate) async fn turn_steer( &self, request_id: &ConnectionRequestId, params: TurnSteerParams, ) -> Result, JSONRPCErrorError> { validate_user_input_image_urls(¶ms.input)?; self.turn_steer_inner(request_id, params) .await .map(|response| Some(response.into())) } pub(crate) async fn turn_interrupt( &self, request_id: &ConnectionRequestId, params: TurnInterruptParams, ) -> Result, JSONRPCErrorError> { let result = self.turn_interrupt_inner(request_id, params).await; if let Err(error) = &result { self.track_error_response(request_id, error, /*error_type*/ None); } result.map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_start( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeStartParams, ) -> Result, JSONRPCErrorError> { self.thread_realtime_start_inner(request_id, params) .await .map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_append_audio( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendAudioParams, ) -> Result, JSONRPCErrorError> { self.thread_realtime_append_audio_inner(request_id, params) .await .map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_append_text( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendTextParams, ) -> Result, JSONRPCErrorError> { self.thread_realtime_append_text_inner(request_id, params) .await .map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_append_speech( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendSpeechParams, ) -> Result, JSONRPCErrorError> { self.thread_realtime_append_speech_inner(request_id, params) .await .map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_stop( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeStopParams, ) -> Result, JSONRPCErrorError> { self.thread_realtime_stop_inner(request_id, params) .await .map(|response| response.map(Into::into)) } pub(crate) async fn thread_realtime_list_voices( &self, ) -> Result, JSONRPCErrorError> { Ok(Some( ThreadRealtimeListVoicesResponse { voices: RealtimeVoicesList::builtin(), } .into(), )) } pub(crate) async fn review_start( &self, request_id: &ConnectionRequestId, params: ReviewStartParams, ) -> Result, JSONRPCErrorError> { if matches!(params.delivery, Some(ApiReviewDelivery::Detached)) { self.outgoing .send_server_notification_to_connections( &[request_id.connection_id], ServerNotification::DeprecationNotice(DeprecationNoticeNotification { summary: "review/start with delivery \"detached\" is deprecated and will be removed in a future release.".to_string(), details: Some("Use thread/start followed by review/start with delivery \"inline\" for a separate review thread, or thread/fork followed by turn/start with your own review instructions.".to_string()), }), ) .await; } self.review_start_inner(request_id, params) .await .map(|()| None) } fn track_error_response( &self, request_id: &ConnectionRequestId, error: &JSONRPCErrorError, error_type: Option, ) { self.analytics_events_client.track_error_response( request_id.connection_id.0, request_id.request_id.clone(), error.clone(), error_type, ); } async fn load_thread( &self, thread_id: &str, ) -> Result<(ThreadId, Arc), JSONRPCErrorError> { // Resolve the core conversation handle from a v2 thread id string. let thread_id = ThreadId::from_string(thread_id) .map_err(|err| invalid_request(format!("invalid thread id: {err}")))?; let thread = self .thread_manager .get_thread(thread_id) .await .map_err(|_| invalid_request(format!("thread not found: {thread_id}")))?; Ok((thread_id, thread)) } async fn ensure_direct_input_allowed( &self, request_id: &ConnectionRequestId, thread: &CodexThread, ) -> Result<(), JSONRPCErrorError> { ensure_direct_input_allowed(thread) .await .inspect_err(|error| { self.track_error_response(request_id, error, /*error_type*/ None); }) } fn normalize_collaboration_mode( &self, mut collaboration_mode: CollaborationMode, ) -> CollaborationMode { if collaboration_mode.settings.developer_instructions.is_none() && let Some(instructions) = builtin_collaboration_mode_presets() .into_iter() .find(|preset| preset.mode == Some(collaboration_mode.mode)) .and_then(|preset| preset.developer_instructions.flatten()) .filter(|instructions| !instructions.is_empty()) { collaboration_mode.settings.developer_instructions = Some(instructions); } collaboration_mode } fn review_request_from_target( target: ApiReviewTarget, ) -> Result<(ReviewRequest, String, String), JSONRPCErrorError> { let cleaned_target = match target { ApiReviewTarget::UncommittedChanges => ApiReviewTarget::UncommittedChanges, ApiReviewTarget::BaseBranch { branch } => { let branch = branch.trim().to_string(); if branch.is_empty() { return Err(invalid_request("branch must not be empty".to_string())); } ApiReviewTarget::BaseBranch { branch } } ApiReviewTarget::Commit { sha, title } => { let sha = sha.trim().to_string(); if sha.is_empty() { return Err(invalid_request("sha must not be empty".to_string())); } let title = title .map(|t| t.trim().to_string()) .filter(|t| !t.is_empty()); ApiReviewTarget::Commit { sha, title } } ApiReviewTarget::Custom { instructions } => { let trimmed = instructions.trim().to_string(); if trimmed.is_empty() { return Err(invalid_request( "instructions must not be empty".to_string(), )); } ApiReviewTarget::Custom { instructions: trimmed, } } }; let core_target = match cleaned_target { ApiReviewTarget::UncommittedChanges => CoreReviewTarget::UncommittedChanges, ApiReviewTarget::BaseBranch { branch } => CoreReviewTarget::BaseBranch { branch }, ApiReviewTarget::Commit { sha, title } => CoreReviewTarget::Commit { sha, title }, ApiReviewTarget::Custom { instructions } => CoreReviewTarget::Custom { instructions }, }; let target_prompt = match &core_target { CoreReviewTarget::UncommittedChanges => { "Review the current code changes (staged, unstaged, and untracked files)." .to_string() } CoreReviewTarget::BaseBranch { branch } => { format!("Review the code changes against the base branch {branch:?}.") } CoreReviewTarget::Commit { sha, .. } => { format!("Review the changes introduced by commit {sha:?}.") } CoreReviewTarget::Custom { instructions } => instructions.clone(), }; let hint = codex_core::review_prompts::user_facing_hint(&core_target); let review_request = ReviewRequest { target: core_target, user_facing_hint: Some(hint.clone()), }; Ok((review_request, hint, target_prompt)) } async fn request_trace_context( &self, request_id: &ConnectionRequestId, ) -> Option { self.outgoing.request_trace_context(request_id).await } async fn submit_core_op( &self, request_id: &ConnectionRequestId, thread: &CodexThread, op: Op, ) -> CodexResult { thread .submit_with_trace(op, self.request_trace_context(request_id).await) .await } pub(super) fn input_too_large_error(actual_chars: usize) -> JSONRPCErrorError { let mut error = invalid_params(format!( "Input exceeds the maximum length of {MAX_USER_INPUT_TEXT_CHARS} characters." )); error.data = Some(serde_json::json!({ "input_error_code": INPUT_TOO_LARGE_ERROR_CODE, "max_chars": MAX_USER_INPUT_TEXT_CHARS, "actual_chars": actual_chars, })); error } pub(super) fn validate_v2_input_limit(items: &[V2UserInput]) -> Result<(), JSONRPCErrorError> { let actual_chars: usize = items.iter().map(V2UserInput::text_char_count).sum(); if actual_chars > MAX_USER_INPUT_TEXT_CHARS { return Err(Self::input_too_large_error(actual_chars)); } Ok(()) } async fn turn_start_inner( &self, request_id: ConnectionRequestId, params: TurnStartParams, app_server_client_name: Option, app_server_client_version: Option, ) -> Result { let (thread_id, thread) = self.load_thread(¶ms.thread_id) .await .inspect_err(|error| { self.track_error_response(&request_id, error, /*error_type*/ None); })?; self.ensure_direct_input_allowed(&request_id, thread.as_ref()) .await?; self.config_manager .check_thread_model_provider(thread.config().await.as_ref()) .await .map_err(|error| config_load_error(&error))?; if let Some(tool_output) = ¶ms.tool_output { if !params.input.is_empty() { return Err(invalid_request( "`toolOutput` cannot be combined with nonempty `input`", )); } if tool_output.name.is_empty() { return Err(invalid_request("`toolOutput.name` must not be empty")); } } let actual_chars = params .input .iter() .map(V2UserInput::text_char_count) .sum::() + params .tool_output .as_ref() .map_or(0, |output| match &output.output { FunctionCallOutputBody::Text(text) => text.chars().count(), FunctionCallOutputBody::ContentItems(items) => items .iter() .map(|item| match item { FunctionCallOutputContentItem::InputText { text } => { text.chars().count() } _ => 0, }) .sum(), }); if actual_chars > MAX_USER_INPUT_TEXT_CHARS { let error = Self::input_too_large_error(actual_chars); self.track_error_response( &request_id, &error, Some(AnalyticsJsonRpcError::Input(InputError::TooLarge)), ); return Err(error); } Self::set_app_server_client_info( thread.as_ref(), app_server_client_name, app_server_client_version, ) .await .inspect_err(|error| { self.track_error_response(&request_id, error, /*error_type*/ None); })?; let runtime_workspace_roots = params .runtime_workspace_roots .map(resolve_runtime_workspace_roots); let environment_selections = resolve_turn_environment_selections(self.thread_manager.as_ref(), params.environments)?; let additional_context = map_additional_context(params.additional_context); let turn_has_input = !params.input.is_empty(); let input = if let Some(tool_output) = params.tool_output { let item = ResponseItem::FunctionCallOutput { id: None, call_id: None, name: Some(tool_output.name), namespace: tool_output.namespace, output: FunctionCallOutputPayload { body: tool_output.output, success: None, }, internal_chat_message_metadata_passthrough: None, }; validate_response_item_image_urls(std::slice::from_ref(&item))?; TurnInput::ResponseItem(item) } else { TurnInput::UserInput { content: params .input .into_iter() .map(V2UserInput::into_core) .collect(), client_id: params.client_user_message_id, } }; let cwd = resolve_request_cwd(params.cwd)?; let environment_override = self .build_environment_override( thread.as_ref(), cwd, runtime_workspace_roots, environment_selections, ) .await; let thread_settings = self .build_thread_settings_overrides( thread.as_ref(), ThreadSettingsBuildParams { method: "turn/start", disabled_plugin_ids: params.disabled_plugin_ids, environment_override, approval_policy: params.approval_policy, approvals_reviewer: params.approvals_reviewer, sandbox_policy: params.sandbox_policy, permissions: params.permissions, model: params.model, service_tier: params.service_tier, effort: params.effort, summary: params.summary, collaboration_mode: params.collaboration_mode, personality: params.personality, }, ) .await?; let submission = thread .start_or_steer_turn( TurnInputRequest::new(input) .with_thread_settings(thread_settings) .on_start(TurnStartOptions { turn_trigger: params.turn_trigger, final_output_json_schema: params.output_schema, service_tier: params.service_tier_for_turn, cyber_access_program: params.cyber_access_program.map(Into::into), ..Default::default() }) .with_additional_context(additional_context) .with_responses_metadata(params.responsesapi_client_metadata) .with_trace(self.request_trace_context(&request_id).await), ) .await .map_err(|err| { let error = internal_error(format!("failed to submit turn input: {err}")); self.track_error_response(&request_id, &error, /*error_type*/ None); error })?; let (turn_id, started) = match submission { TurnInputSubmission::Started { turn_id } => (turn_id, true), TurnInputSubmission::Steered { turn_id } => (turn_id, false), TurnInputSubmission::NotSubmitted { reason } => { let error = if reason == NotSubmittedReason::ServerDraining { crate::error_code::server_draining_error() } else { internal_error(format!("failed to submit turn input: {reason:?}")) }; self.track_error_response(&request_id, &error, /*error_type*/ None); return Err(error); } }; if turn_has_input && started { let config_snapshot = thread.config_snapshot().await; if config_snapshot.is_primary_environment_configured() { codex_memories_write::start_memories_startup_task( Arc::clone(&self.thread_manager), Arc::clone(&self.auth_manager), thread_id, Arc::clone(&thread), thread.config().await, config_snapshot.permission_profile, &config_snapshot.session_source, ); } } self.outgoing .record_request_turn_id(&request_id, &turn_id) .await; let turn = Turn { id: turn_id, items: vec![], items_view: TurnItemsView::NotLoaded, error: None, status: TurnStatus::InProgress, started_at: None, completed_at: None, duration_ms: None, }; Ok(TurnStartResponse { turn }) } async fn build_environment_override( &self, thread: &CodexThread, cwd: Option, workspace_roots: Option>, environment_selections: Option>, ) -> ThreadEnvironmentOverride { if cwd.is_none() && workspace_roots.is_none() && environment_selections.is_none() { return ThreadEnvironmentOverride::default(); } // Explicit environment selections own their roots and pass through unchanged. Top-level // `runtimeWorkspaceRoots` is only a compatibility input for default environments. if let Some(environment_selections) = environment_selections { let legacy_fallback_cwd = match cwd { Some(cwd) => cwd, None => match environment_selections .iter() .find(|selection| selection.environment_id == LOCAL_ENVIRONMENT_ID) .and_then(|selection| selection.cwd.to_abs_path().ok()) { Some(cwd) => cwd, None => thread.config_snapshot().await.cwd().clone(), }, }; return ThreadEnvironmentOverride { environments: Some(TurnEnvironmentSelections::new( legacy_fallback_cwd, environment_selections, )), ..Default::default() }; } // Default-environment updates retain the task's fallback roots, not its active roots. let snapshot = thread.thread_settings_snapshot().await; let current_cwd = snapshot.cwd; let legacy_fallback_cwd = cwd.unwrap_or_else(|| current_cwd.clone()); let workspace_roots = match workspace_roots { Some(workspace_roots) => workspace_roots, None => path_utils::replace_path_and_deduplicate( snapshot.runtime_workspace_roots.unwrap_or_default(), current_cwd.as_path(), legacy_fallback_cwd.clone(), ), }; let environment_selections = self .thread_manager .default_environment_selections(&legacy_fallback_cwd, &workspace_roots); ThreadEnvironmentOverride { environments: Some(TurnEnvironmentSelections::new( legacy_fallback_cwd, environment_selections, )), runtime_workspace_roots: Some(workspace_roots), } } async fn build_thread_settings_overrides( &self, thread: &CodexThread, params: ThreadSettingsBuildParams, ) -> Result { let ThreadSettingsBuildParams { method, disabled_plugin_ids, environment_override: ThreadEnvironmentOverride { environments, runtime_workspace_roots, }, approval_policy, approvals_reviewer, sandbox_policy, permissions, model, service_tier, effort, summary, collaboration_mode, personality, } = params; if sandbox_policy.is_some() && permissions.is_some() { return Err(invalid_request( "`permissions` cannot be combined with `sandboxPolicy`", )); } let collaboration_mode = collaboration_mode.map(|mode| self.normalize_collaboration_mode(mode)); let has_environment_override = environments.is_some(); // `thread/settings/update` only acknowledges that the update was queued. // Clients that send dependent partial updates should wait for // `thread/settings/updated` or combine the fields in one request. let snapshot = if permissions.is_some() { Some(thread.config_snapshot().await) } else { None }; let has_any_overrides = has_environment_override || disabled_plugin_ids.is_some() || approval_policy.is_some() || approvals_reviewer.is_some() || sandbox_policy.is_some() || permissions.is_some() || model.is_some() || service_tier.is_some() || effort.is_some() || summary.is_some() || collaboration_mode.is_some() || personality.is_some(); let approval_policy = approval_policy.map(codex_app_server_protocol::AskForApproval::to_core); let approvals_reviewer = approvals_reviewer.map(codex_app_server_protocol::ApprovalsReviewer::to_core); let sandbox_policy = sandbox_policy.map(|policy| policy.to_core()); let (permission_profile, active_permission_profile, profile_workspace_roots) = if let Some(permissions) = permissions { let Some(snapshot) = snapshot.as_ref() else { return Err(internal_error(format!( "{method} permission selection missing thread snapshot" ))); }; let overrides = ConfigOverrides { cwd: environments .as_ref() .map(|environments| environments.legacy_fallback_cwd.to_path_buf()), default_permissions: Some(permissions), codex_linux_sandbox_exe: self.arg0_paths.codex_linux_sandbox_exe.clone(), main_execve_wrapper_exe: self.arg0_paths.main_execve_wrapper_exe.clone(), ..Default::default() }; let config = self .config_manager .load_for_cwd( /*request_overrides*/ None, overrides, Some(snapshot.cwd().to_path_buf()), ) .await .map_err(|err| config_load_error(&err))?; // Startup config is allowed to fall back when requirements // disallow a configured profile. An explicit settings update // is different: reject it before accepting the request. if let Some(warning) = config.startup_warnings.iter().find(|warning| { warning.contains("Configured value for `permission_profile` is disallowed") }) { return Err(invalid_request(format!( "invalid thread settings override: {warning}" ))); } ( Some(config.permissions.permission_profile().clone()), config.permissions.active_permission_profile(), Some(config.permissions.profile_workspace_roots().to_vec()), ) } else { (None, None, None) }; let effort = effort.map(Some); if has_any_overrides { thread .preview_thread_settings_overrides(CodexThreadSettingsOverrides { disabled_plugin_ids: disabled_plugin_ids.clone(), environments: environments.clone(), runtime_workspace_roots: runtime_workspace_roots.clone(), approval_policy, approvals_reviewer, sandbox_policy: sandbox_policy.clone(), permission_profile: permission_profile.clone(), active_permission_profile: active_permission_profile.clone(), profile_workspace_roots: profile_workspace_roots.clone(), windows_sandbox_level: None, model: model.clone(), effort: effort.clone(), summary, service_tier: service_tier.clone(), collaboration_mode: collaboration_mode.clone(), personality, }) .await .map_err(|err| { invalid_request(format!("invalid thread settings override: {err}")) })?; } Ok(codex_protocol::protocol::ThreadSettingsOverrides { disabled_plugin_ids, environments, runtime_workspace_roots, profile_workspace_roots, approval_policy, approvals_reviewer, sandbox_policy, permission_profile, active_permission_profile, windows_sandbox_level: None, model, effort, summary, service_tier, collaboration_mode, personality, }) } async fn thread_settings_update_inner( &self, request_id: &ConnectionRequestId, params: ThreadSettingsUpdateParams, ) -> Result { let (_, thread) = self.load_thread(¶ms.thread_id).await?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; let cwd = resolve_request_cwd(params.cwd)?; let environment_override = self .build_environment_override( thread.as_ref(), cwd, /*workspace_roots*/ None, /*environment_selections*/ None, ) .await; let thread_settings = self .build_thread_settings_overrides( thread.as_ref(), ThreadSettingsBuildParams { method: "thread/settings/update", disabled_plugin_ids: params.disabled_plugin_ids, environment_override, approval_policy: params.approval_policy, approvals_reviewer: params.approvals_reviewer, sandbox_policy: params.sandbox_policy, permissions: params.permissions, model: params.model, service_tier: params.service_tier, effort: params.effort, summary: params.summary, collaboration_mode: params.collaboration_mode, personality: params.personality, }, ) .await?; if thread_settings != codex_protocol::protocol::ThreadSettingsOverrides::default() { self.submit_core_op( request_id, thread.as_ref(), Op::ThreadSettings { thread_settings }, ) .await .map_err(|err| internal_error(format!("failed to update thread settings: {err}")))?; } Ok(ThreadSettingsUpdateResponse {}) } async fn thread_inject_items_response_inner( &self, request_id: &ConnectionRequestId, params: ThreadInjectItemsParams, ) -> Result { let (_, thread) = self.load_thread(¶ms.thread_id).await?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; let items = params .items .into_iter() .enumerate() .map(|(index, value)| { serde_json::from_value::(value) .map_err(|err| format!("items[{index}] is not a valid response item: {err}")) }) .collect::, _>>() .map_err(invalid_request)?; validate_response_item_image_urls(&items)?; thread .inject_response_items(items) .await .map_err(|err| match err.details() { CodexErrorDetails::InvalidRequest(message) => invalid_request(message.clone()), _ => internal_error(format!("failed to inject response items: {err}")), })?; Ok(ThreadInjectItemsResponse {}) } async fn set_app_server_client_info( thread: &CodexThread, app_server_client_name: Option, app_server_client_version: Option, ) -> Result<(), JSONRPCErrorError> { let mcp_elicitations_auto_deny = xcode_26_4_mcp_elicitations_auto_deny( app_server_client_name.as_deref(), app_server_client_version.as_deref(), ); thread .set_app_server_client_info( app_server_client_name, app_server_client_version, mcp_elicitations_auto_deny, ) .await .map_err(|err| internal_error(format!("failed to set app server client info: {err}"))) } async fn turn_steer_inner( &self, request_id: &ConnectionRequestId, params: TurnSteerParams, ) -> Result { let (_, thread) = self .load_thread(¶ms.thread_id) .await .inspect_err(|error| { self.track_error_response(request_id, error, /*error_type*/ None); })?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; self.config_manager .check_thread_model_provider(thread.config().await.as_ref()) .await .map_err(|error| config_load_error(&error))?; if params.expected_turn_id.is_empty() { return Err(invalid_request("expectedTurnId must not be empty")); } self.outgoing .record_request_turn_id(request_id, ¶ms.expected_turn_id) .await; if let Err(error) = Self::validate_v2_input_limit(¶ms.input) { self.track_error_response( request_id, &error, Some(AnalyticsJsonRpcError::Input(InputError::TooLarge)), ); return Err(error); } let mapped_items: Vec = params .input .into_iter() .map(V2UserInput::into_core) .collect(); let additional_context = map_additional_context(params.additional_context); let submission = thread .steer_turn( TurnInputRequest::new(TurnInput::UserInput { content: mapped_items, client_id: params.client_user_message_id, }) .with_additional_context(additional_context) .with_responses_metadata(params.responsesapi_client_metadata), params.expected_turn_id, ) .await .map_err(|err| { let error = internal_error(format!("failed to steer turn: {err}")); self.track_error_response(request_id, &error, /*error_type*/ None); error })?; let turn_id = match submission { SteerSubmission::Steered { turn_id } => turn_id, SteerSubmission::NotSubmitted { reason } => { let (message, data, error_type) = match reason { NotSubmittedReason::ServerDraining => { return Err(crate::error_code::server_draining_error()); } NotSubmittedReason::NoActiveTurn | NotSubmittedReason::NotIdle => ( "no active turn to steer".to_string(), None, Some(AnalyticsJsonRpcError::TurnSteer( TurnSteerRequestError::NoActiveTurn, )), ), NotSubmittedReason::ExpectedTurnMismatch { expected, actual } => ( format!("expected active turn id `{expected}` but found `{actual}`"), None, Some(AnalyticsJsonRpcError::TurnSteer( TurnSteerRequestError::ExpectedTurnMismatch, )), ), NotSubmittedReason::ActiveTurnNotSteerable { turn_kind } => { let (message, turn_steer_error) = match turn_kind { codex_protocol::protocol::NonSteerableTurnKind::Review => ( "cannot steer a review turn".to_string(), TurnSteerRequestError::NonSteerableReview, ), codex_protocol::protocol::NonSteerableTurnKind::Compact => ( "cannot steer a compact turn".to_string(), TurnSteerRequestError::NonSteerableCompact, ), }; let error = TurnError { misalignment: None, message: message.clone(), codex_error_info: Some(CodexErrorInfo::ActiveTurnNotSteerable { turn_kind: turn_kind.into(), }), additional_details: None, }; let data = match serde_json::to_value(error) { Ok(data) => Some(data), Err(error) => { tracing::error!( ?error, "failed to serialize active-turn-not-steerable turn error" ); None } }; ( message, data, Some(AnalyticsJsonRpcError::TurnSteer(turn_steer_error)), ) } NotSubmittedReason::EmptyInput => ( "input must not be empty".to_string(), None, Some(AnalyticsJsonRpcError::Input(InputError::Empty)), ), NotSubmittedReason::ActiveTurnOutputSchemaMismatch => ( "active turn uses a different output schema".to_string(), None, None, ), NotSubmittedReason::PendingTriggerTurn | NotSubmittedReason::PlanMode | NotSubmittedReason::Superseded => ( "no active turn to steer".to_string(), None, Some(AnalyticsJsonRpcError::TurnSteer( TurnSteerRequestError::NoActiveTurn, )), ), }; let mut error = invalid_request(message); error.data = data; self.track_error_response(request_id, &error, error_type); return Err(error); } }; Ok(TurnSteerResponse { turn_id }) } async fn prepare_realtime_conversation_thread( &self, request_id: &ConnectionRequestId, thread_id: &str, ) -> Result)>, JSONRPCErrorError> { let (thread_id, thread) = self.load_thread(thread_id).await?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; match self .ensure_conversation_listener( thread_id, request_id.connection_id, /*raw_events_enabled*/ false, ) .await { Ok(EnsureConversationListenerResult::Attached) => {} Ok(EnsureConversationListenerResult::ConnectionClosed) => { return Ok(None); } Err(error) => return Err(error), } Ok(Some((thread_id, thread))) } async fn thread_realtime_start_inner( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeStartParams, ) -> Result, JSONRPCErrorError> { let attaches_existing_call = matches!( ¶ms.transport, Some(ThreadRealtimeStartTransport::ExistingCall { .. }) ); if attaches_existing_call { let unsupported_option = if params.include_startup_context == Some(true) { Some("includeStartupContext") } else if params.prompt.is_some() { Some("prompt") } else if params .initial_items .as_ref() .is_some_and(|items| !items.is_empty()) { Some("initialItems") } else if params.model.is_some() { Some("model") } else if params.voice.is_some() { Some("voice") } else if params.delegation_ack_filler.is_some() { Some("delegationAckFiller") } else { None }; if let Some(option) = unsupported_option { return Err(invalid_request(format!( "existingCall transport does not support {option}" ))); } } let Some((_, thread)) = self .prepare_realtime_conversation_thread(request_id, ¶ms.thread_id) .await? else { return Ok(None); }; self.submit_core_op( request_id, thread.as_ref(), Op::RealtimeConversationStart(ConversationStartParams { client_managed_handoffs: params.client_managed_handoffs.unwrap_or(false), delegation_ack_filler: params.delegation_ack_filler, flush_transcript_tail_on_session_end: params .flush_transcript_tail_on_session_end .unwrap_or(false), codex_responses_as_items: params.codex_responses_as_items.unwrap_or(false), codex_response_item_prefix: params.codex_response_item_prefix, codex_response_handoff_mode: params.codex_response_handoff_mode.unwrap_or_default(), codex_response_handoff_channel_prefixes: params .codex_response_handoff_channel_prefixes, model: params.model, output_modality: params.output_modality, include_startup_context: params .include_startup_context .unwrap_or(!attaches_existing_call), initial_items: params .initial_items .unwrap_or_default() .into_iter() .map(|item| ConversationTextParams { text: item.text, role: item.role, }) .collect(), realtime_start_instructions: params.realtime_start_instructions, realtime_end_instructions: params.realtime_end_instructions, prompt: params.prompt, realtime_session_id: params.realtime_session_id, transport: params.transport.map(|transport| match transport { ThreadRealtimeStartTransport::Websocket => { ConversationStartTransport::Websocket } ThreadRealtimeStartTransport::Webrtc { sdp } => { ConversationStartTransport::Webrtc { sdp } } ThreadRealtimeStartTransport::ExistingCall { call_id } => { ConversationStartTransport::ExistingCall { call_id, sideband_base_url: None, } } }), version: params.version, voice: params.voice, }), ) .await .map_err(|err| internal_error(format!("failed to start realtime conversation: {err}")))?; Ok(Some(ThreadRealtimeStartResponse::default())) } async fn thread_realtime_append_audio_inner( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendAudioParams, ) -> Result, JSONRPCErrorError> { let Some((_, thread)) = self .prepare_realtime_conversation_thread(request_id, ¶ms.thread_id) .await? else { return Ok(None); }; self.submit_core_op( request_id, thread.as_ref(), Op::RealtimeConversationAudio(ConversationAudioParams { frame: params.audio.into(), }), ) .await .map_err(|err| { internal_error(format!( "failed to append realtime conversation audio: {err}" )) })?; Ok(Some(ThreadRealtimeAppendAudioResponse::default())) } async fn thread_realtime_append_text_inner( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendTextParams, ) -> Result, JSONRPCErrorError> { let Some((_, thread)) = self .prepare_realtime_conversation_thread(request_id, ¶ms.thread_id) .await? else { return Ok(None); }; self.submit_core_op( request_id, thread.as_ref(), Op::RealtimeConversationText(ConversationTextParams { text: params.text, role: params.role, }), ) .await .map_err(|err| { internal_error(format!( "failed to append realtime conversation text: {err}" )) })?; Ok(Some(ThreadRealtimeAppendTextResponse::default())) } async fn thread_realtime_append_speech_inner( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeAppendSpeechParams, ) -> Result, JSONRPCErrorError> { let Some((_, thread)) = self .prepare_realtime_conversation_thread(request_id, ¶ms.thread_id) .await? else { return Ok(None); }; self.submit_core_op( request_id, thread.as_ref(), Op::RealtimeConversationSpeech(ConversationSpeechParams { text: params.text }), ) .await .map_err(|err| { internal_error(format!( "failed to append realtime conversation speech: {err}" )) })?; Ok(Some(ThreadRealtimeAppendSpeechResponse::default())) } async fn thread_realtime_stop_inner( &self, request_id: &ConnectionRequestId, params: ThreadRealtimeStopParams, ) -> Result, JSONRPCErrorError> { let Some((_, thread)) = self .prepare_realtime_conversation_thread(request_id, ¶ms.thread_id) .await? else { return Ok(None); }; self.submit_core_op(request_id, thread.as_ref(), Op::RealtimeConversationClose) .await .map_err(|err| { internal_error(format!("failed to stop realtime conversation: {err}")) })?; Ok(Some(ThreadRealtimeStopResponse::default())) } fn build_review_turn(turn_id: String, display_text: &str) -> Turn { let items = if display_text.is_empty() { Vec::new() } else { vec![ThreadItem::UserMessage { id: turn_id.clone(), client_id: None, content: vec![V2UserInput::Text { text: display_text.to_string(), // Review prompt display text is synthesized; no UI element ranges to preserve. text_elements: Vec::new(), }], }] }; Turn { id: turn_id, items, items_view: TurnItemsView::NotLoaded, error: None, status: TurnStatus::InProgress, started_at: None, completed_at: None, duration_ms: None, } } async fn emit_review_started( &self, request_id: &ConnectionRequestId, turn: Turn, review_thread_id: String, ) { let response = ReviewStartResponse { turn, review_thread_id, }; self.outgoing .send_response(request_id.clone(), response) .await; } async fn start_inline_review( &self, request_id: &ConnectionRequestId, parent_thread: Arc, review_request: ReviewRequest, display_text: &str, parent_thread_id: String, ) -> std::result::Result<(), JSONRPCErrorError> { let turn_id = self .submit_core_op( request_id, parent_thread.as_ref(), Op::Review { review_request }, ) .await .map_err(|err| internal_error(format!("failed to start review: {err}")))?; let turn = Self::build_review_turn(turn_id, display_text); self.emit_review_started(request_id, turn, parent_thread_id) .await; Ok(()) } async fn start_detached_review( &self, request_id: &ConnectionRequestId, parent_thread: Arc, prompt: &str, ) -> std::result::Result<(), JSONRPCErrorError> { // AgentRunner::start still delegates to spawn_subagent, which forks from the parent's // full history. Paginated threads only allow bounded model-context reads, so keep this // closed until detached review has a bounded fork path. if matches!( parent_thread.config_snapshot().await.history_mode, codex_protocol::protocol::ThreadHistoryMode::Paginated ) { return Err(invalid_request( "paginated threads do not support detached review", )); } let mut config = parent_thread.config().await.as_ref().clone(); if let Some(review_model) = &config.review_model { config.model = Some(review_model.clone()); } let AgentRun { thread_id, thread: review_thread, turn_id, } = self .agent_runner .start( parent_thread.session_configured().thread_id, AgentInvocation { config, prompt: prompt.to_string(), parent_trace: self.request_trace_context(request_id).await, }, ) .await .map_err(|err| internal_error(format!("failed to start detached review: {err}")))?; let fallback_provider = self.config.model_provider_id.as_str(); let stored_thread = match review_thread .read_thread( /*include_archived*/ true, /*include_history*/ false, ) .await { Ok(stored_thread) => { let (thread, _) = thread_from_stored_thread(stored_thread, fallback_provider, &self.config.cwd); Some(thread) } Err(err) => { tracing::warn!("failed to load summary for review thread {thread_id}: {err}"); None } }; if let Some(mut thread) = stored_thread { let config_snapshot = review_thread.config_snapshot().await; apply_live_thread_settings(&mut thread, &config_snapshot); thread.session_id = review_thread.session_configured().session_id.to_string(); self.thread_watch_manager .upsert_thread_silently(&thread.id) .await; thread.status = resolve_thread_status( self.thread_watch_manager .loaded_status_for_thread(&thread.id) .await, /*has_in_progress_turn*/ false, ); let notif = thread_started_notification(thread); self.outgoing .send_server_notification(ServerNotification::ThreadStarted(notif)) .await; } log_listener_attach_result( self.ensure_conversation_listener( thread_id, request_id.connection_id, /*raw_events_enabled*/ false, ) .await, thread_id, request_id.connection_id, "review thread", ); let turn = Self::build_review_turn(turn_id, prompt); let review_thread_id = thread_id.to_string(); self.emit_review_started(request_id, turn, review_thread_id) .await; Ok(()) } async fn review_start_inner( &self, request_id: &ConnectionRequestId, params: ReviewStartParams, ) -> Result<(), JSONRPCErrorError> { let ReviewStartParams { thread_id, target, delivery, } = params; let (_, parent_thread) = self.load_thread(&thread_id).await?; self.ensure_direct_input_allowed(request_id, parent_thread.as_ref()) .await?; self.config_manager .check_thread_model_provider(parent_thread.config().await.as_ref()) .await .map_err(|error| config_load_error(&error))?; let (review_request, display_text, target_prompt) = Self::review_request_from_target(target)?; match delivery.unwrap_or(ApiReviewDelivery::Inline).to_core() { CoreReviewDelivery::Inline => { self.start_inline_review( request_id, parent_thread, review_request, &display_text, thread_id, ) .await?; } CoreReviewDelivery::Detached => { let review_skill_path = system_cache_root_dir(&self.config.codex_home) .join("review-agent") .join("SKILL.md"); let prompt = format!( "Use [$review-agent]({}) for this review.\n\n{target_prompt}", review_skill_path.display() ); let actual_chars = prompt.chars().count(); if actual_chars > MAX_USER_INPUT_TEXT_CHARS { return Err(Self::input_too_large_error(actual_chars)); } self.start_detached_review(request_id, parent_thread, &prompt) .await?; } } Ok(()) } async fn turn_interrupt_inner( &self, request_id: &ConnectionRequestId, params: TurnInterruptParams, ) -> Result, JSONRPCErrorError> { let TurnInterruptParams { thread_id, turn_id } = params; let is_startup_interrupt = turn_id.is_empty(); let (thread_uuid, thread) = self.load_thread(&thread_id).await?; // Record turn interrupts so we can reply when TurnAborted arrives. Startup // interrupts do not have a turn and are acknowledged after submission. if !is_startup_interrupt { let thread_state = self.thread_state_manager.thread_state(thread_uuid).await; let is_running = matches!(thread.agent_status().await, AgentStatus::Running); { let mut thread_state = thread_state.lock().await; if let Some(active_turn) = thread_state.active_turn_snapshot() { if active_turn.id != turn_id { return Err(invalid_request(format!( "expected active turn id {turn_id} but found {}", active_turn.id ))); } } else if thread_state.last_terminal_turn_id.as_deref() == Some(turn_id.as_str()) || !is_running { return Err(invalid_request("no active turn to interrupt")); } thread_state.pending_interrupts.push(request_id.clone()); } self.outgoing .record_request_turn_id(request_id, &turn_id) .await; } // Submit the interrupt. Turn interrupts respond upon TurnAborted; startup // interrupts respond here because startup cancellation has no turn event. match self .submit_core_op(request_id, thread.as_ref(), Op::Interrupt) .await { Ok(_) if is_startup_interrupt => Ok(Some(TurnInterruptResponse {})), Ok(_) => Ok(None), Err(err) => { if !is_startup_interrupt { let thread_state = self.thread_state_manager.thread_state(thread_uuid).await; let mut thread_state = thread_state.lock().await; thread_state .pending_interrupts .retain(|pending_request_id| pending_request_id != request_id); } let interrupt_target = if is_startup_interrupt { "startup" } else { "turn" }; Err(internal_error(format!( "failed to interrupt {interrupt_target}: {err}" ))) } } } fn listener_task_context(&self) -> ListenerTaskContext { ListenerTaskContext { thread_manager: Arc::clone(&self.thread_manager), thread_state_manager: self.thread_state_manager.clone(), outgoing: Arc::clone(&self.outgoing), pending_thread_unloads: Arc::clone(&self.pending_thread_unloads), thread_watch_manager: self.thread_watch_manager.clone(), codex_home: self.config.codex_home.to_path_buf(), thread_unload_delay: self.config.thread_unload_delay, skills_watcher: Arc::clone(&self.skills_watcher), turn_cost_worker: self.turn_cost_worker.clone(), } } async fn ensure_conversation_listener( &self, conversation_id: ThreadId, connection_id: ConnectionId, raw_events_enabled: bool, ) -> Result { super::thread_lifecycle::ensure_conversation_listener( self.listener_task_context(), conversation_id, connection_id, raw_events_enabled, ) .await } } fn xcode_26_4_mcp_elicitations_auto_deny( client_name: Option<&str>, client_version: Option<&str>, ) -> bool { // Xcode 26.4 shipped before app-server MCP elicitation requests were // client-visible. Keep elicitations auto-denied for that client line. // TODO: Remove this compatibility hack once Xcode 26.4 ages out. client_name == Some("Xcode") && client_version.is_some_and(|version| version.starts_with("26.4")) }