Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Binary file not shown.
9 changes: 9 additions & 0 deletions codex-rs/app-server-protocol/src/protocol/v2/thread.rs
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,15 @@ impl From<ThreadTurnsListResponse> for TurnsPage {
pub struct ThreadForkParams {
pub thread_id: String,

/// Start from the loaded parent's effective config, tools, environments, and
/// persisted history to maximize prompt-cache reuse. Requires a loaded,
/// persisted parent and `ephemeral: true`. Enables `features.reasoning_effort_override`
/// on the fork. Cannot be combined with configuration overrides, `path`, `lastTurnId`,
/// or `beforeTurnId`.
#[experimental("thread/fork.experimentalPredictionMode")]
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub experimental_prediction_mode: bool,

/// Optional last turn id to fork through, inclusive.
///
/// When specified, turns after `last_turn_id` are omitted from the fork.
Expand Down
103 changes: 88 additions & 15 deletions codex-rs/app-server/src/request_processors/thread_processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4890,8 +4890,34 @@ impl ThreadRequestProcessor {
app_server_client_version: Option<String>,
client_mcp_extensions: ClientMcpExtensions,
) -> Result<(), JSONRPCErrorError> {
let has_config_overrides = params.model.is_some()
|| params.model_provider.is_some()
|| params.service_tier.is_some()
|| params.cwd.is_some()
|| params.runtime_workspace_roots.is_some()
|| params.approval_policy.is_some()
|| params.approvals_reviewer.is_some()
|| params.sandbox.is_some()
|| params.permissions.is_some()
|| params.base_instructions.is_some()
|| params.developer_instructions.is_some()
|| params
.config
.as_ref()
.is_some_and(|config| !config.is_empty());
if params.experimental_prediction_mode && has_config_overrides {
return Err(invalid_request(
"`experimentalPredictionMode` cannot be combined with configuration overrides",
));
}
if params.experimental_prediction_mode && !params.ephemeral {
return Err(invalid_request(
"`experimentalPredictionMode` requires `ephemeral: true`",
));
}
let ThreadForkParams {
thread_id,
experimental_prediction_mode,
last_turn_id,
before_turn_id,
path,
Expand All @@ -4918,13 +4944,45 @@ impl ThreadRequestProcessor {
"`permissions` cannot be combined with `sandbox`",
));
}
if experimental_prediction_mode
&& (path.is_some() || last_turn_id.is_some() || before_turn_id.is_some())
{
return Err(invalid_request(
"`experimentalPredictionMode` requires a thread id without a path or turn cutoff",
));
}
let source_thread = self
.read_stored_thread_for_resume(
&thread_id,
path.as_ref(),
/*include_history*/ false,
)
.await?;
let inherited_fork = if experimental_prediction_mode {
// Keep inherited state and its preparation off the fork handler's stack.
let (mut options, settings) = Box::pin(async {
self.thread_manager
.fork_options_from_parent(source_thread.thread_id)
.await
.map(|(options, settings)| (Box::new(options), Box::new(settings)))
})
.await
.map_err(|err| match err.details() {
CodexErrorDetails::ThreadNotFound(_) => {
invalid_request("`experimentalPredictionMode` requires a loaded parent")
}
_ => internal_error(format!("failed to inherit parent state: {err}")),
})?;
options.config.ephemeral = ephemeral;
options
.config
.features
.enable(Feature::ReasoningEffortOverride)
.map_err(|err| invalid_request(err.to_string()))?;
Some((options, settings))
} else {
None
};
let paginated_source = matches!(source_thread.history_mode, ThreadHistoryMode::Paginated);
if last_turn_id.is_some() && before_turn_id.is_some() {
return Err(invalid_request(
Expand Down Expand Up @@ -5011,7 +5069,7 @@ impl ThreadRequestProcessor {

// Persist Windows sandbox mode.
let mut cli_overrides = cli_overrides.unwrap_or_default();
if cfg!(windows) {
if cfg!(windows) && !experimental_prediction_mode {
let mode = self.config.permissions.windows_sandbox_mode.or_else(|| {
match WindowsSandboxLevel::from_config(&self.config) {
WindowsSandboxLevel::Elevated => Some(WindowsSandboxModeToml::Elevated),
Expand Down Expand Up @@ -5122,15 +5180,24 @@ impl ThreadRequestProcessor {
.map(|profile| profile.id);
}
}
// Derive a Config using the same logic as new conversation, honoring overrides if provided.
let config = self
.config_manager
.load_for_cwd(request_overrides, typesafe_overrides, history_cwd)
.await
.map_err(|err| config_load_error(&err))?;
let goals_enabled = config.features.enabled(Feature::Goals);
let (inherited_options, inherited_settings) = inherited_fork.unzip();
let options = if let Some(options) = inherited_options {
*options
} else {
// Derive a Config using the same logic as new conversation, honoring overrides if provided.
let config = self
.config_manager
.load_for_cwd(request_overrides, typesafe_overrides, history_cwd)
.await
.map_err(|err| config_load_error(&err))?;
StartThreadOptions {
client_mcp_extensions,
..StartThreadOptions::new(config)
}
};
let goals_enabled = options.config.features.enabled(Feature::Goals);

let fallback_model_provider = config.model_provider_id.clone();
let fallback_model_provider = options.config.model_provider_id.clone();
let parent_trace = self.request_trace_context(&request_id).await;
let thread_source = thread_source.map(Into::into);

Expand Down Expand Up @@ -5186,7 +5253,7 @@ impl ThreadRequestProcessor {
.then(|| restored_token_usage_turn_id(&history_items, ephemeral_turns.as_slice()));
let token_usage_history_items = paginated_source.then(|| Arc::clone(&history_items));
let inherited_project_id = source_thread.project_id.clone();
let reserved_thread_id = if config.ephemeral {
let reserved_thread_id = if options.config.ephemeral {
None
} else {
stage_pending_thread_metadata(
Expand All @@ -5205,9 +5272,8 @@ impl ThreadRequestProcessor {
let fork_options = StartThreadOptions {
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
..StartThreadOptions::new(config)
..options
};
let new_thread = if let Some(prepared_fork) = prepared_fork {
self.thread_manager
Expand Down Expand Up @@ -5247,6 +5313,13 @@ impl ThreadRequestProcessor {
}
};

if let Some(settings) = inherited_settings {
// Construct the restore future off the fork handler's stack
Box::pin(async { Box::pin(forked_thread.restore_thread_settings(*settings)).await })
.await
.map_err(|err| invalid_request(err.to_string()))?;
}

Self::set_app_server_client_info(
forked_thread.as_ref(),
app_server_client_name,
Expand Down Expand Up @@ -5395,17 +5468,17 @@ impl ThreadRequestProcessor {
let response = ThreadForkResponse {
thread: thread.clone(),
disabled_plugin_ids: config_snapshot.disabled_plugin_ids,
model: session_configured.model,
model: config_snapshot.model.clone(),
model_provider: session_configured.model_provider_id,
service_tier: session_configured.service_tier,
service_tier: config_snapshot.service_tier.clone(),
cwd: session_configured.cwd,
runtime_workspace_roots: config_snapshot.workspace_roots,
instruction_sources,
approval_policy: session_configured.approval_policy.into(),
approvals_reviewer: session_configured.approvals_reviewer.into(),
sandbox,
active_permission_profile,
reasoning_effort: session_configured.reasoning_effort,
reasoning_effort: config_snapshot.reasoning_effort.clone(),
multi_agent_mode: MultiAgentMode::ExplicitRequestOnly,
};

Expand Down
Loading
Loading