mirror of
https://github.com/openai/codex.git
synced 2026-09-29 16:57:06 +08:00
Add turn trigger metadata (#40665)
## What changed - Add an optional `turnTrigger` field to app-server `turn/start` requests and expose it in the generated protocol schemas. - Propagate non-empty trigger values to Responses request metadata as the reserved `turn_trigger` field, while preserving the original value when a request steers an active turn. - Classify turns started by queue dispatch, goal continuation, retry recovery, and realtime handoff. ## Testing - Cover HTTP and WebSocket metadata forwarding, steering behavior, reserved metadata handling, and the built-in trigger classifications. GitOrigin-RevId: c12fe2c522286db21b76081ae6154fbf6bfb3639
This commit is contained in:
@@ -5266,6 +5266,13 @@
|
||||
},
|
||||
"threadId": {
|
||||
"type": "string"
|
||||
},
|
||||
"turnTrigger": {
|
||||
"description": "Optional source classification for the caller that starts this turn. Ignored when this request steers an already-active turn.",
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
]
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
|
||||
+7
@@ -23959,6 +23959,13 @@
|
||||
},
|
||||
"threadId": {
|
||||
"type": "string"
|
||||
},
|
||||
"turnTrigger": {
|
||||
"description": "Optional source classification for the caller that starts this turn. Ignored when this request steers an already-active turn.",
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
]
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
|
||||
+7
@@ -21694,6 +21694,13 @@
|
||||
},
|
||||
"threadId": {
|
||||
"type": "string"
|
||||
},
|
||||
"turnTrigger": {
|
||||
"description": "Optional source classification for the caller that starts this turn. Ignored when this request steers an already-active turn.",
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
]
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
|
||||
@@ -675,6 +675,13 @@
|
||||
},
|
||||
"threadId": {
|
||||
"type": "string"
|
||||
},
|
||||
"turnTrigger": {
|
||||
"description": "Optional source classification for the caller that starts this turn. Ignored when this request steers an already-active turn.",
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
]
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
|
||||
BIN
Binary file not shown.
BIN
Binary file not shown.
@@ -11,6 +11,10 @@ import type { SandboxPolicy } from "./SandboxPolicy";
|
||||
import type { UserInput } from "./UserInput";
|
||||
|
||||
export type TurnStartParams = {threadId: string, clientUserMessageId?: string | null, input: Array<UserInput>, /**
|
||||
* Optional source classification for the caller that starts this turn.
|
||||
* Ignored when this request steers an already-active turn.
|
||||
*/
|
||||
turnTrigger?: string | null, /**
|
||||
* Override the working directory for this turn and subsequent turns.
|
||||
*/
|
||||
cwd?: string | null, /**
|
||||
|
||||
@@ -4714,6 +4714,7 @@ fn turn_start_params_preserve_explicit_null_service_tier() {
|
||||
thread_id: "thread_123".to_string(),
|
||||
client_user_message_id: None,
|
||||
input: vec![],
|
||||
turn_trigger: None,
|
||||
responsesapi_client_metadata: None,
|
||||
additional_context: None,
|
||||
environments: None,
|
||||
|
||||
@@ -121,6 +121,10 @@ pub struct TurnStartParams {
|
||||
#[ts(optional = nullable)]
|
||||
pub client_user_message_id: Option<String>,
|
||||
pub input: Vec<UserInput>,
|
||||
/// Optional source classification for the caller that starts this turn.
|
||||
/// Ignored when this request steers an already-active turn.
|
||||
#[ts(optional = nullable)]
|
||||
pub turn_trigger: Option<String>,
|
||||
/// Optional metadata to enrich Codex's ResponsesAPI turn metadata.
|
||||
///
|
||||
/// Entries are flattened into the JSON string sent as
|
||||
|
||||
@@ -209,7 +209,7 @@ Example with notification opt-out:
|
||||
- `thread/backgroundTerminals/terminate` — terminate one running background terminal by app-server `processId` (experimental; requires `capabilities.experimentalApi`); returns whether a process was terminated.
|
||||
- `thread/rollback` — deprecated and will be removed soon. Drop the last N turns from the agent’s in-memory context and persist a rollback marker in the rollout so future resumes see the pruned history; returns the updated `thread` (with `turns` populated) on success. Paginated threads do not support rollback. Parent-owned Multi-Agent V2 subagents reject direct rollback requests.
|
||||
- `thread/revert` — experimental. Replace a loaded paginated thread's durable history with the prefix strictly before `beforeTurnId` while preserving its thread id. The operation interrupts an active turn if needed, leaves older rollout files immutable, reloads the thread, returns updated thread metadata with empty `turns` plus pagination cursors, and emits `thread/reverted`. It does not revert local file changes. Parent-owned Multi-Agent V2 subagents reject direct revert requests.
|
||||
- `turn/start` — add user input to a thread and begin Codex generation; responds with the initial `turn` object and streams `turn/started`, `item/*`, and `turn/completed` notifications. `clientUserMessageId` is optional; when supplied, the corresponding `userMessage` item echoes it as `clientId`. Experimental `runtimeWorkspaceRoots` supplies the default roots for newly resolved environment selections. Explicit `environments[].runtimeWorkspaceRoots` override that fallback with environment-native absolute paths. Prefer experimental `permissions` profile selection by id for permission overrides; the legacy `sandboxPolicy` field is still accepted but cannot be combined with `permissions`. For `collaborationMode`, `settings.developer_instructions: null` means "use built-in instructions for the selected mode". Deprecated experimental `multiAgentMode` is ignored; Ultra reasoning effort selects proactive behavior. Parent-owned Multi-Agent V2 subagents reject direct turns.
|
||||
- `turn/start` — add user input to a thread and begin Codex generation; responds with the initial `turn` object and streams `turn/started`, `item/*`, and `turn/completed` notifications. Optional `turnTrigger` classifies who or what started a new turn and is sent as `turn_trigger` in Responses request metadata; it is ignored if the request steers an active turn. `clientUserMessageId` is optional; when supplied, the corresponding `userMessage` item echoes it as `clientId`. Experimental `runtimeWorkspaceRoots` supplies the default roots for newly resolved environment selections. Explicit `environments[].runtimeWorkspaceRoots` override that fallback with environment-native absolute paths. Prefer experimental `permissions` profile selection by id for permission overrides; the legacy `sandboxPolicy` field is still accepted but cannot be combined with `permissions`. For `collaborationMode`, `settings.developer_instructions: null` means "use built-in instructions for the selected mode". Deprecated experimental `multiAgentMode` is ignored; Ultra reasoning effort selects proactive behavior. Parent-owned Multi-Agent V2 subagents reject direct turns.
|
||||
- `thread/inject_items` — append raw Responses API items to a loaded thread’s model-visible history without starting a user turn; returns `{}` on success. Parent-owned Multi-Agent V2 subagents reject direct item injection.
|
||||
- `turn/settings/update` — experimental; publish a narrow model-settings patch to the exact live task identified by `threadId` and `turnId`, regardless of task kind. Requires `step_model_switching`; returns `status: "applied"` or `status: "targetUnavailable"`, or a request error if rejected. Future-thread settings and already captured steps are unchanged. Parent-owned Multi-Agent V2 subagents reject direct settings updates.
|
||||
- `turn/steer` — add user input to an already in-flight regular turn without starting a new turn; returns the active `turnId` that accepted the input. `clientUserMessageId` is optional; when supplied, the corresponding `userMessage` item echoes it as `clientId`. Review and manual compaction turns reject `turn/steer`. Parent-owned Multi-Agent V2 subagents reject direct steering.
|
||||
|
||||
@@ -663,6 +663,7 @@ async fn turn_start_jsonrpc_span_parents_core_turn_spans() -> Result<()> {
|
||||
text: "hello".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: None,
|
||||
responsesapi_client_metadata: None,
|
||||
additional_context: None,
|
||||
cwd: None,
|
||||
|
||||
@@ -580,6 +580,7 @@ impl TurnRequestProcessor {
|
||||
})
|
||||
.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,
|
||||
..Default::default()
|
||||
|
||||
@@ -75,6 +75,7 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result
|
||||
"context_window_id".to_string(),
|
||||
"client-supplied".to_string(),
|
||||
),
|
||||
("turn_trigger".to_string(), "client-supplied".to_string()),
|
||||
]);
|
||||
let turn_req = mcp
|
||||
.send_turn_start_request(TurnStartParams {
|
||||
@@ -84,6 +85,7 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result
|
||||
text: "Hello".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: Some("user".to_string()),
|
||||
responsesapi_client_metadata: Some(client_metadata.clone()),
|
||||
..Default::default()
|
||||
})
|
||||
@@ -106,6 +108,7 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result
|
||||
assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123"));
|
||||
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
|
||||
assert_eq!(metadata["thread_source"].as_str(), Some("automation"));
|
||||
assert_eq!(metadata["turn_trigger"].as_str(), Some("user"));
|
||||
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
|
||||
assert!(metadata.get("installation_id").is_some());
|
||||
assert!(metadata.get("session_id").is_some());
|
||||
@@ -437,6 +440,7 @@ async fn turn_steer_updates_client_metadata_on_follow_up_responses_request_v2()
|
||||
text: "Run sleep".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: Some("user".to_string()),
|
||||
responsesapi_client_metadata: Some(start_metadata.clone()),
|
||||
..Default::default()
|
||||
})
|
||||
@@ -490,6 +494,7 @@ async fn turn_steer_updates_client_metadata_on_follow_up_responses_request_v2()
|
||||
Some("fiber-start-123")
|
||||
);
|
||||
assert_eq!(first_metadata["turn_id"].as_str(), Some(turn_id.as_str()));
|
||||
assert_eq!(first_metadata["turn_trigger"].as_str(), Some("user"));
|
||||
|
||||
let second_metadata = requests[1]
|
||||
.header("x-codex-turn-metadata")
|
||||
@@ -502,6 +507,7 @@ async fn turn_steer_updates_client_metadata_on_follow_up_responses_request_v2()
|
||||
);
|
||||
assert_eq!(second_metadata["origin"].as_str(), Some("gaas"));
|
||||
assert_eq!(second_metadata["turn_id"].as_str(), Some(turn_id.as_str()));
|
||||
assert_eq!(second_metadata["turn_trigger"].as_str(), Some("user"));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -555,6 +561,7 @@ async fn turn_start_forwards_client_metadata_to_responses_websocket_request_body
|
||||
text: "Hello".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: Some("user".to_string()),
|
||||
responsesapi_client_metadata: Some(client_metadata),
|
||||
..Default::default()
|
||||
})
|
||||
@@ -589,6 +596,7 @@ async fn turn_start_forwards_client_metadata_to_responses_websocket_request_body
|
||||
assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123"));
|
||||
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
|
||||
assert_eq!(metadata["thread_source"].as_str(), Some("automation"));
|
||||
assert_eq!(metadata["turn_trigger"].as_str(), Some("user"));
|
||||
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
|
||||
assert!(metadata.get("session_id").is_some());
|
||||
assert_eq!(
|
||||
|
||||
@@ -412,6 +412,7 @@ async fn idle_queue_dispatch_preserves_client_id() -> Result<()> {
|
||||
metadata["turn_id"].as_str(),
|
||||
Some(completed.turn.id.as_str())
|
||||
);
|
||||
assert_eq!(metadata["turn_trigger"].as_str(), Some("queue"));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -715,7 +716,7 @@ async fn queue_start_without_id_starts_the_head_when_idle() -> Result<()> {
|
||||
create_final_assistant_message_sse_response("first queued message done")?,
|
||||
create_final_assistant_message_sse_response("second queued message done")?,
|
||||
];
|
||||
let (mut app, _codex_home, _server) = queue_app(responses).await?;
|
||||
let (mut app, _codex_home, server) = queue_app(responses).await?;
|
||||
let thread_id = app
|
||||
.start_thread(ThreadStartParams::default())
|
||||
.await?
|
||||
@@ -765,6 +766,26 @@ async fn queue_start_without_id_starts_the_head_when_idle() -> Result<()> {
|
||||
assert_eq!(completed.turn.status, TurnStatus::Completed);
|
||||
}
|
||||
assert!(list_queue(&mut app, &thread_id).await?.data.is_empty());
|
||||
|
||||
let requests = server
|
||||
.received_requests()
|
||||
.await
|
||||
.context("mock request capture unavailable")?;
|
||||
let response_requests = requests
|
||||
.iter()
|
||||
.filter(|request| request.url.path().ends_with("/responses"))
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(response_requests.len(), 3);
|
||||
for request in &response_requests[1..] {
|
||||
let metadata_header = request
|
||||
.headers
|
||||
.get("x-codex-turn-metadata")
|
||||
.context("queued model request is missing its x-codex-turn-metadata header")?
|
||||
.to_str()
|
||||
.context("queued turn metadata header is not valid ASCII")?;
|
||||
let metadata: Value = serde_json::from_str(metadata_header)?;
|
||||
assert_eq!(metadata["turn_trigger"].as_str(), Some("queue"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -3021,6 +3021,23 @@ async fn thread_goal_lifecycle_emits_analytics_and_clear_deletes_goal() -> Resul
|
||||
.is_some()
|
||||
);
|
||||
|
||||
let requests = server
|
||||
.received_requests()
|
||||
.await
|
||||
.expect("wiremock should record response requests");
|
||||
let response_requests = requests
|
||||
.iter()
|
||||
.filter(|request| request.url.path().ends_with("/responses"))
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(response_requests.len(), 2);
|
||||
let metadata_header = response_requests[1]
|
||||
.headers
|
||||
.get("x-codex-turn-metadata")
|
||||
.expect("goal continuation should include turn metadata")
|
||||
.to_str()?;
|
||||
let metadata: serde_json::Value = serde_json::from_str(metadata_header)?;
|
||||
assert_eq!(metadata["turn_trigger"].as_str(), Some("goal"));
|
||||
|
||||
let status = wait_for_goal_event(
|
||||
&server,
|
||||
DEFAULT_READ_TIMEOUT,
|
||||
|
||||
@@ -357,6 +357,7 @@ async fn turn_start_steers_active_turn_and_returns_active_turn_id() -> Result<()
|
||||
text: "start".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: Some("user".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
})
|
||||
@@ -377,6 +378,7 @@ async fn turn_start_steers_active_turn_and_returns_active_turn_id() -> Result<()
|
||||
text: "steer".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: Some("goal".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
})
|
||||
@@ -391,6 +393,18 @@ async fn turn_start_steers_active_turn_and_returns_active_turn_id() -> Result<()
|
||||
mcp.read_stream_until_notification_message("turn/completed"),
|
||||
)
|
||||
.await??;
|
||||
|
||||
let requests = server.requests().await;
|
||||
assert_eq!(requests.len(), 2);
|
||||
for request in requests {
|
||||
let body: Value = serde_json::from_slice(&request)?;
|
||||
let turn_metadata: Value = serde_json::from_str(
|
||||
body["client_metadata"]["x-codex-turn-metadata"]
|
||||
.as_str()
|
||||
.context("expected x-codex-turn-metadata")?,
|
||||
)?;
|
||||
assert_eq!(turn_metadata["turn_trigger"].as_str(), Some("user"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -2767,6 +2781,7 @@ async fn turn_start_explicit_local_environment_updates_legacy_cwd_between_turns(
|
||||
text: "first turn".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: None,
|
||||
responsesapi_client_metadata: None,
|
||||
additional_context: None,
|
||||
cwd: Some(first_cwd.clone()),
|
||||
@@ -2816,6 +2831,7 @@ async fn turn_start_explicit_local_environment_updates_legacy_cwd_between_turns(
|
||||
text: "second turn".to_string(),
|
||||
text_elements: Vec::new(),
|
||||
}],
|
||||
turn_trigger: None,
|
||||
responsesapi_client_metadata: None,
|
||||
additional_context: None,
|
||||
cwd: None,
|
||||
|
||||
@@ -217,6 +217,7 @@ pub(crate) async fn run_codex_thread_one_shot(
|
||||
service_tier: None,
|
||||
parent_turn_id: Some(parent_turn_id),
|
||||
root_turn_id,
|
||||
..Default::default()
|
||||
}),
|
||||
TurnInputMode::StartIfIdle,
|
||||
)
|
||||
|
||||
@@ -1182,6 +1182,7 @@ async fn run_review_on_session(
|
||||
service_tier: None,
|
||||
parent_turn_id: Some(parent_turn.sub_id.clone()),
|
||||
root_turn_id: parent_turn.turn_metadata_state.root_turn_id(),
|
||||
..Default::default()
|
||||
}),
|
||||
TurnInputMode::StartIfIdle,
|
||||
);
|
||||
|
||||
@@ -43,6 +43,7 @@ pub(crate) const PARENT_TURN_ID_KEY: &str = "parent_turn_id";
|
||||
pub(crate) const ROOT_TURN_ID_KEY: &str = "root_turn_id";
|
||||
pub(crate) const SUBAGENT_KIND_KEY: &str = "subagent_kind";
|
||||
pub(crate) const THREAD_SOURCE_KEY: &str = "thread_source";
|
||||
pub(crate) const TURN_TRIGGER_KEY: &str = "turn_trigger";
|
||||
pub(crate) const SANDBOX_KEY: &str = "sandbox";
|
||||
pub(crate) const SANDBOX_MODE_KEY: &str = "sandbox_mode";
|
||||
pub(crate) const AUTO_REVIEW_ENABLED_KEY: &str = "auto_review_enabled";
|
||||
@@ -76,6 +77,7 @@ const RESERVED_METADATA_KEYS: &[&str] = &[
|
||||
ROOT_TURN_ID_KEY,
|
||||
SUBAGENT_KIND_KEY,
|
||||
THREAD_SOURCE_KEY,
|
||||
TURN_TRIGGER_KEY,
|
||||
SANDBOX_KEY,
|
||||
SANDBOX_MODE_KEY,
|
||||
AUTO_REVIEW_ENABLED_KEY,
|
||||
@@ -220,6 +222,7 @@ pub struct CodexResponsesMetadata {
|
||||
pub(crate) subagent_header: Option<String>,
|
||||
pub(crate) subagent_kind: Option<String>,
|
||||
pub(crate) thread_source: Option<ThreadSource>,
|
||||
pub(crate) turn_trigger: Option<String>,
|
||||
pub(crate) sandbox: Option<String>,
|
||||
pub(crate) sandbox_mode: Option<String>,
|
||||
pub(crate) auto_review_enabled: Option<bool>,
|
||||
@@ -255,6 +258,7 @@ impl CodexResponsesMetadata {
|
||||
subagent_header: None,
|
||||
subagent_kind: None,
|
||||
thread_source: None,
|
||||
turn_trigger: None,
|
||||
sandbox: None,
|
||||
sandbox_mode: None,
|
||||
auto_review_enabled: None,
|
||||
@@ -379,6 +383,7 @@ impl CodexResponsesMetadata {
|
||||
root_turn_id: self.root_turn_id.as_deref(),
|
||||
subagent_kind: self.subagent_kind.as_deref(),
|
||||
thread_source: self.thread_source.as_ref(),
|
||||
turn_trigger: self.turn_trigger.as_deref(),
|
||||
sandbox: self.sandbox.as_deref(),
|
||||
sandbox_mode: self.sandbox_mode.as_deref(),
|
||||
auto_review_enabled: self.auto_review_enabled,
|
||||
@@ -508,6 +513,8 @@ struct CodexTurnMetadataPayload<'a> {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
thread_source: Option<&'a ThreadSource>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
turn_trigger: Option<&'a str>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
sandbox: Option<&'a str>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
sandbox_mode: Option<&'a str>,
|
||||
|
||||
@@ -120,6 +120,7 @@ impl PreparedTurnInputSettings {
|
||||
kind: TurnStartKind,
|
||||
) -> CodexResult<Option<Arc<TurnContext>>> {
|
||||
let TurnStartOptions {
|
||||
turn_trigger,
|
||||
final_output_json_schema,
|
||||
service_tier,
|
||||
parent_turn_id,
|
||||
@@ -147,6 +148,11 @@ impl PreparedTurnInputSettings {
|
||||
let Some((turn_context, settings_snapshot)) = turn_context else {
|
||||
return Ok(None);
|
||||
};
|
||||
if let Some(turn_trigger) = turn_trigger {
|
||||
turn_context
|
||||
.turn_metadata_state
|
||||
.set_turn_trigger(turn_trigger);
|
||||
}
|
||||
if emit_thread_settings_applied {
|
||||
thread_settings::emit_applied(session, submission_id, settings_snapshot).await;
|
||||
}
|
||||
@@ -205,7 +211,12 @@ pub(super) async fn handle_recovery(
|
||||
thread_settings: ThreadSettingsOverrides,
|
||||
submission_id: String,
|
||||
) -> CodexResult<TurnInputSubmission> {
|
||||
let request = TurnInputRequest::user_input(Vec::new()).with_thread_settings(thread_settings);
|
||||
let request = TurnInputRequest::user_input(Vec::new())
|
||||
.with_thread_settings(thread_settings)
|
||||
.on_start(TurnStartOptions {
|
||||
turn_trigger: Some("retry".to_string()),
|
||||
..Default::default()
|
||||
});
|
||||
start_if_idle(session, request, submission_id, TurnStartKind::Recovery).await
|
||||
}
|
||||
|
||||
@@ -482,7 +493,11 @@ impl Session {
|
||||
TurnInputRequest::user_input(vec![UserInput::Text {
|
||||
text,
|
||||
text_elements: Vec::new(),
|
||||
}]),
|
||||
}])
|
||||
.on_start(TurnStartOptions {
|
||||
turn_trigger: Some("realtime".to_string()),
|
||||
..Default::default()
|
||||
}),
|
||||
TurnInputMode::StartOrSteer,
|
||||
submission_id.clone(),
|
||||
)
|
||||
|
||||
@@ -113,6 +113,7 @@ pub(crate) struct TurnMetadataState {
|
||||
subagent_header: Option<String>,
|
||||
subagent_kind: Option<String>,
|
||||
thread_source: Option<ThreadSource>,
|
||||
turn_trigger: OnceLock<String>,
|
||||
turn_id: String,
|
||||
// TODO(anp): Derive this cached tag from TurnEnvironment::sandbox_context
|
||||
// so metadata reflects the selected environment's backend.
|
||||
@@ -184,6 +185,7 @@ impl TurnMetadataState {
|
||||
subagent_header: subagent_header_value(session_source),
|
||||
subagent_kind: subagent_metadata_kind(session_source),
|
||||
thread_source,
|
||||
turn_trigger: OnceLock::new(),
|
||||
turn_id,
|
||||
sandbox,
|
||||
sandbox_mode,
|
||||
@@ -297,6 +299,13 @@ impl TurnMetadataState {
|
||||
let _ = self.root_turn_id.set(root_turn_id);
|
||||
}
|
||||
|
||||
pub(crate) fn set_turn_trigger(&self, turn_trigger: String) {
|
||||
if turn_trigger.trim().is_empty() {
|
||||
return;
|
||||
}
|
||||
let _ = self.turn_trigger.set(turn_trigger);
|
||||
}
|
||||
|
||||
pub(crate) fn root_turn_id(&self) -> Option<String> {
|
||||
self.root_turn_id
|
||||
.get()
|
||||
@@ -390,6 +399,7 @@ impl TurnMetadataState {
|
||||
subagent_header: self.subagent_header.clone(),
|
||||
subagent_kind: self.subagent_kind.clone(),
|
||||
thread_source: self.thread_source.clone(),
|
||||
turn_trigger: self.turn_trigger.get().cloned(),
|
||||
sandbox: self.sandbox.clone(),
|
||||
sandbox_mode: self.sandbox_mode.clone(),
|
||||
auto_review_enabled: Some(self.auto_review_enabled),
|
||||
|
||||
@@ -12,6 +12,7 @@ use crate::responses_metadata::PARENT_TURN_ID_KEY;
|
||||
use crate::responses_metadata::ROOT_TURN_ID_KEY;
|
||||
use crate::responses_metadata::SANDBOX_MODE_KEY;
|
||||
use crate::responses_metadata::TOOL_NAMESPACES_INFO_KEY;
|
||||
use crate::responses_metadata::TURN_TRIGGER_KEY;
|
||||
use crate::responses_metadata::TurnToolFunctionInfo;
|
||||
use crate::responses_metadata::TurnToolNamespaceInfo;
|
||||
use crate::responses_metadata::TurnToolSource;
|
||||
@@ -682,6 +683,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields(
|
||||
)]));
|
||||
state.set_parent_turn_id("parent-turn-a".to_string());
|
||||
state.set_root_turn_id("root-turn-a".to_string());
|
||||
state.set_turn_trigger("goal".to_string());
|
||||
state.set_responsesapi_client_metadata(HashMap::from([
|
||||
(
|
||||
"codex_security_surface".to_string(),
|
||||
@@ -737,6 +739,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields(
|
||||
"client-supplied".to_string(),
|
||||
),
|
||||
("thread_source".to_string(), "client-supplied".to_string()),
|
||||
(TURN_TRIGGER_KEY.to_string(), "client-supplied".to_string()),
|
||||
("request_kind".to_string(), "client-supplied".to_string()),
|
||||
(
|
||||
"turn_started_at_unix_ms".to_string(),
|
||||
@@ -814,6 +817,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields(
|
||||
assert_eq!(json[ROOT_TURN_ID_KEY].as_str(), Some("root-turn-a"));
|
||||
assert_eq!(json["subagent_kind"].as_str(), Some("thread_spawn"));
|
||||
assert_eq!(json["thread_source"].as_str(), Some("automation"));
|
||||
assert_eq!(json[TURN_TRIGGER_KEY].as_str(), Some("goal"));
|
||||
assert_eq!(json["turn_id"].as_str(), Some("turn-a"));
|
||||
assert!(json.get("request_kind").is_none());
|
||||
assert!(json.get(WINDOW_ID_KEY).is_none());
|
||||
@@ -953,7 +957,12 @@ fn turn_metadata_state_overlays_compaction_only_on_compaction_requests() {
|
||||
|
||||
#[test]
|
||||
fn responses_api_metadata_rejects_reserved_keys() {
|
||||
for reserved_key in ["thread_source", WINDOW_ID_KEY, CONTEXT_WINDOW_ID_KEY] {
|
||||
for reserved_key in [
|
||||
"thread_source",
|
||||
TURN_TRIGGER_KEY,
|
||||
WINDOW_ID_KEY,
|
||||
CONTEXT_WINDOW_ID_KEY,
|
||||
] {
|
||||
assert_eq!(
|
||||
validate_extra_metadata(
|
||||
BTreeMap::from([(reserved_key.to_string(), "sdk".to_string())]).iter()
|
||||
|
||||
@@ -4184,6 +4184,13 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> {
|
||||
.await;
|
||||
|
||||
let request = response_mock.single_request();
|
||||
let turn_metadata: Value = serde_json::from_str(
|
||||
request
|
||||
.header("x-codex-turn-metadata")
|
||||
.as_deref()
|
||||
.context("realtime-routed turn should include turn metadata")?,
|
||||
)?;
|
||||
assert_eq!(turn_metadata["turn_trigger"].as_str(), Some("realtime"));
|
||||
let user_texts = request.message_input_texts("user");
|
||||
assert!(user_texts.iter().any(|text| text
|
||||
== "<realtime_delegation>\n <input>text from realtime</input>\n <transcript_delta>user: text from realtime</transcript_delta>\n</realtime_delegation>"));
|
||||
|
||||
@@ -178,9 +178,16 @@ async fn recover_turn_if_idle_preserves_id_and_resumes_plan_mode() {
|
||||
})
|
||||
.await;
|
||||
|
||||
let user_input_groups = response_mock
|
||||
.single_request()
|
||||
.message_input_text_groups("user");
|
||||
let request = response_mock.single_request();
|
||||
let turn_metadata: Value = serde_json::from_str(
|
||||
request
|
||||
.header("x-codex-turn-metadata")
|
||||
.as_deref()
|
||||
.expect("recovered turn should include turn metadata"),
|
||||
)
|
||||
.expect("recovered turn metadata should be valid JSON");
|
||||
assert_eq!(turn_metadata["turn_trigger"].as_str(), Some("retry"));
|
||||
let user_input_groups = request.message_input_text_groups("user");
|
||||
assert_eq!(user_input_groups.len(), 1);
|
||||
assert_eq!(user_input_groups[0].len(), 1);
|
||||
assert!(user_input_groups[0][0].starts_with("<environment_context>"));
|
||||
|
||||
@@ -973,6 +973,7 @@ async fn run_exec_session(args: ExecRunArgs) -> anyhow::Result<()> {
|
||||
request_id: request_ids.next(),
|
||||
params: TurnStartParams {
|
||||
thread_id: primary_thread_id_for_span.clone(),
|
||||
turn_trigger: None,
|
||||
client_user_message_id: None,
|
||||
input: items.into_iter().map(Into::into).collect(),
|
||||
responsesapi_client_metadata: None,
|
||||
|
||||
@@ -7,6 +7,7 @@ use codex_core::StartIfIdleSubmission;
|
||||
use codex_core::ThreadManager;
|
||||
use codex_core::TurnInput;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_core::TurnStartOptions;
|
||||
use codex_protocol::ThreadId;
|
||||
use codex_protocol::models::ResponseItem;
|
||||
use codex_protocol::protocol::ThreadGoal;
|
||||
@@ -406,7 +407,12 @@ impl GoalRuntimeHandle {
|
||||
let item = continuation_steering_item(&protocol_goal_from_state(goal));
|
||||
|
||||
match thread
|
||||
.start_turn_if_idle(TurnInputRequest::new(TurnInput::ResponseItem(item)))
|
||||
.start_turn_if_idle(
|
||||
TurnInputRequest::new(TurnInput::ResponseItem(item)).on_start(TurnStartOptions {
|
||||
turn_trigger: Some("goal".to_string()),
|
||||
..Default::default()
|
||||
}),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(StartIfIdleSubmission::Started { .. }) => {}
|
||||
|
||||
@@ -10,6 +10,7 @@ use codex_core::StartIfIdleSubmission;
|
||||
use codex_core::ThreadManager;
|
||||
use codex_core::TurnInput;
|
||||
use codex_core::TurnInputRequest;
|
||||
use codex_core::TurnStartOptions;
|
||||
use codex_extension_api::ExtensionEventSink;
|
||||
use codex_extension_api::ExtensionFuture;
|
||||
use codex_extension_api::ThreadIdleCause;
|
||||
@@ -388,7 +389,12 @@ impl QueuedItemService {
|
||||
return Err(QueueServiceError::InvalidInput);
|
||||
};
|
||||
let submission = thread
|
||||
.start_turn_if_idle(TurnInputRequest::new(input).with_trace(trace))
|
||||
.start_turn_if_idle(TurnInputRequest::new(input).with_trace(trace).on_start(
|
||||
TurnStartOptions {
|
||||
turn_trigger: Some("queue".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
))
|
||||
.await?;
|
||||
if matches!(submission, StartIfIdleSubmission::Started { .. }) {
|
||||
self.delete_locked(thread_id, queued_item_id).await?;
|
||||
@@ -431,7 +437,10 @@ impl QueuedItemService {
|
||||
}
|
||||
|
||||
match thread
|
||||
.start_turn_if_idle(TurnInputRequest::new(input))
|
||||
.start_turn_if_idle(TurnInputRequest::new(input).on_start(TurnStartOptions {
|
||||
turn_trigger: Some("queue".to_string()),
|
||||
..Default::default()
|
||||
}))
|
||||
.await
|
||||
{
|
||||
Ok(StartIfIdleSubmission::Started { .. }) => {
|
||||
|
||||
@@ -143,6 +143,9 @@ pub enum TurnInputMode {
|
||||
/// child input, Core also compares root lineage to detect ambiguity.
|
||||
#[derive(Clone, Debug, Default)]
|
||||
pub struct TurnStartOptions {
|
||||
/// Source classification for the caller that starts a new turn.
|
||||
/// Ignored when the submitted input steers an active turn.
|
||||
pub turn_trigger: Option<String>,
|
||||
/// Structured-output schema for a new turn. When steering, Core rejects
|
||||
/// the input if the active turn uses a different schema.
|
||||
pub final_output_json_schema: Option<Value>,
|
||||
|
||||
@@ -1201,6 +1201,7 @@ impl AppServerSession {
|
||||
request_id,
|
||||
params: TurnStartParams {
|
||||
thread_id: thread_id.to_string(),
|
||||
turn_trigger: None,
|
||||
client_user_message_id: None,
|
||||
input: items,
|
||||
responsesapi_client_metadata: None,
|
||||
|
||||
Reference in New Issue
Block a user