Forward app-server turn clientMetadata to Responses (#16009)

## Summary
App-server v2 already receives turn-scoped `clientMetadata`, but the
Rust app-server was dropping it before the outbound Responses request.
This change keeps the fix lightweight by threading that metadata through
the existing turn-metadata path rather than inventing a new transport.

## What we're trying to do and why
We want turn-scoped metadata from the app-server protocol layer,
especially fields like Hermes/GAAS run IDs, to survive all the way to
the actual Responses API request so it is visible in downstream
websocket request logging and analytics.

The specific bug was:
- app-server protocol uses camelCase `clientMetadata`
- Responses transport already has an existing turn metadata carrier:
`x-codex-turn-metadata`
- websocket transport already rewrites that header into
`request.request_body.client_metadata["x-codex-turn-metadata"]`
- but the Rust app-server never parsed or stored `clientMetadata`, so
nothing from the app-server request was making it into that existing
path

This PR fixes that without adding a new header or a second metadata
channel.

## How we did it
### Protocol surface
- Add optional `clientMetadata` to v2 `TurnStartParams` and
`TurnSteerParams`
- Regenerate the JSON schema / TypeScript fixtures
- Update app-server docs to describe the field and its behavior

### Runtime plumbing
- Add a dedicated core op for app-server user input carrying turn-scoped
metadata: `Op::UserInputWithClientMetadata`
- Wire `turn/start` and `turn/steer` through that op / signature path
instead of dropping the metadata at the message-processor boundary
- Store the metadata in `TurnMetadataState`

### Transport behavior
- Reuse the existing serialized `x-codex-turn-metadata` payload
- Merge the new app-server `clientMetadata` into that JSON additively
- Do **not** replace built-in reserved fields already present in the
turn metadata payload
- Keep websocket behavior unchanged at the outer shape level: it still
sends only `client_metadata["x-codex-turn-metadata"]`, but that JSON
string now contains the merged fields
- Keep HTTP fallback behavior unchanged except that the existing
`x-codex-turn-metadata` header now includes the merged fields too

### Request shape before / after
Before, a websocket `response.create` looked like:
```json
{
  "type": "response.create",
  "client_metadata": {
    "x-codex-turn-metadata": "{\"session_id\":\"...\",\"turn_id\":\"...\"}"
  }
}
```
Even if the app-server caller supplied `clientMetadata`, it was not
represented there.

After, the same request shape is preserved, but the serialized payload
now includes the new turn-scoped fields:
```json
{
  "type": "response.create",
  "client_metadata": {
    "x-codex-turn-metadata": "{\"session_id\":\"...\",\"turn_id\":\"...\",\"fiber_run_id\":\"fiber-start-123\",\"origin\":\"gaas\"}"
  }
}
```

## Validation
### Targeted tests added / updated
- protocol round-trip coverage for `clientMetadata` on `turn/start` and
`turn/steer`
- protocol round-trip coverage for `Op::UserInputWithClientMetadata`
- `TurnMetadataState` merge test proving client metadata is added
without overwriting reserved built-in fields
- websocket request-shape test proving outbound `response.create`
contains merged metadata inside
`client_metadata["x-codex-turn-metadata"]`
- app-server integration tests proving:
- `turn/start` forwards `clientMetadata` into the outbound Responses
request path
  - websocket warmup + real turn request both behave correctly
  - `turn/steer` updates the follow-up request metadata

### Commands run
- `just write-app-server-schema`
- `cargo test -p codex-app-server-protocol`
- `cargo test -p codex-protocol`
- `cargo test -p codex-core
turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields
--lib`
- `cargo test -p codex-core --test all
responses_websocket_preserves_custom_turn_metadata_fields`
- `cargo test -p codex-app-server --test all client_metadata`
- `cargo test -p codex-app-server --test all
turn_start_forwards_client_metadata_to_responses_websocket_request_body_v2
-- --nocapture`
- `just fmt`
- `just fix -p codex-core -p codex-protocol -p codex-app-server-protocol
-p codex-app-server`
- `just fix -p codex-exec -p codex-tui-app-server`
- `just argument-comment-lint`

### Full suite note
`cargo test` in `codex-rs` still fails in:
-
`suite::v2::turn_interrupt::turn_interrupt_resolves_pending_command_approval_request`

I verified that same failure on a clean detached `HEAD` worktree with an
isolated `CARGO_TARGET_DIR`, so it is not caused by this patch.
This commit is contained in:
neil-oai
2026-04-09 14:52:37 -04:00
committed by GitHub
Unverified
parent 244b15c95d
commit a92a5085bd
51 changed files with 867 additions and 45 deletions
@@ -3,9 +3,8 @@
// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually.
import type { UserInput } from "./UserInput";
export type TurnSteerParams = { threadId: string, input: Array<UserInput>,
/**
export type TurnSteerParams = {threadId: string, input: Array<UserInput>, /**
* Required active turn id precondition. The request fails when it does not
* match the currently active turn.
*/
expectedTurnId: string, };
expectedTurnId: string};
@@ -388,6 +388,7 @@ client_request_definitions! {
},
TurnSteer => "turn/steer" {
params: v2::TurnSteerParams,
inspect_params: true,
response: v2::TurnSteerResponse,
},
TurnInterrupt => "turn/interrupt" {
@@ -4034,6 +4034,10 @@ pub enum TurnStatus {
pub struct TurnStartParams {
pub thread_id: String,
pub input: Vec<UserInput>,
/// Optional turn-scoped Responses API client metadata.
#[experimental("turn/start.responsesapiClientMetadata")]
#[ts(optional = nullable)]
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
/// Override the working directory for this turn and subsequent turns.
#[ts(optional = nullable)]
pub cwd: Option<PathBuf>,
@@ -4144,12 +4148,18 @@ pub struct TurnStartResponse {
pub turn: Turn,
}
#[derive(Serialize, Deserialize, Debug, Default, Clone, PartialEq, JsonSchema, TS)]
#[derive(
Serialize, Deserialize, Debug, Default, Clone, PartialEq, JsonSchema, TS, ExperimentalApi,
)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct TurnSteerParams {
pub thread_id: String,
pub input: Vec<UserInput>,
/// Optional turn-scoped Responses API client metadata.
#[experimental("turn/steer.responsesapiClientMetadata")]
#[ts(optional = nullable)]
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
/// Required active turn id precondition. The request fails when it does not
/// match the currently active turn.
pub expected_turn_id: String,
@@ -8422,6 +8432,7 @@ mod tests {
let without_override = TurnStartParams {
thread_id: "thread_123".to_string(),
input: vec![],
responsesapi_client_metadata: None,
cwd: None,
approval_policy: None,
approvals_reviewer: None,
@@ -6666,6 +6666,7 @@ impl CodexMessageProcessor {
Op::UserInput {
items: mapped_items,
final_output_json_schema: params.output_schema,
responsesapi_client_metadata: params.responsesapi_client_metadata,
},
)
.await;
@@ -6746,7 +6747,11 @@ impl CodexMessageProcessor {
.collect();
match thread
.steer_input(mapped_items, Some(&params.expected_turn_id))
.steer_input(
mapped_items,
Some(&params.expected_turn_id),
params.responsesapi_client_metadata,
)
.await
{
Ok(turn_id) => {
@@ -606,6 +606,7 @@ async fn turn_start_jsonrpc_span_parents_core_turn_spans() -> Result<()> {
text: "hello".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
cwd: None,
approval_policy: None,
sandbox_policy: None,
@@ -106,13 +106,10 @@ async fn remote_control_state_runtime(codex_home: &TempDir) -> Arc<StateRuntime>
}
fn remote_control_url_for_listener(listener: &TcpListener) -> String {
format!(
"http://localhost:{}/backend-api/",
listener
.local_addr()
.expect("listener should have a local addr")
.port()
)
let addr = listener
.local_addr()
.expect("listener should have a local addr");
format!("http://{addr}/backend-api/")
}
#[tokio::test]
@@ -931,13 +931,10 @@ mod tests {
}
fn remote_control_url_for_listener(listener: &TcpListener) -> String {
format!(
"http://localhost:{}/backend-api/",
listener
.local_addr()
.expect("listener should have a local addr")
.port()
)
let addr = listener
.local_addr()
.expect("listener should have a local addr");
format!("http://{addr}/backend-api/")
}
fn remote_control_auth_dot_json(access_token: &str) -> AuthDotJson {
@@ -1041,14 +1038,6 @@ mod tests {
let remote_control_url = remote_control_url_for_listener(&listener);
let remote_control_target =
normalize_remote_control_url(&remote_control_url).expect("target should parse");
let server_task = tokio::spawn(async move {
let (stream, request_line) = accept_http_request(&listener).await;
assert_eq!(
request_line,
"GET /backend-api/wham/remote/control/server HTTP/1.1"
);
respond_with_status_and_headers(stream, "401 Unauthorized", &[], "unauthorized").await;
});
let codex_home = TempDir::new().expect("temp dir should create");
save_auth(
codex_home.path(),
@@ -1076,6 +1065,15 @@ mod tests {
)
.expect("fresh auth should save");
let server_task = tokio::spawn(async move {
let (stream, request_line) = accept_http_request(&listener).await;
assert_eq!(
request_line,
"GET /backend-api/wham/remote/control/server HTTP/1.1"
);
respond_with_status_and_headers(stream, "401 Unauthorized", &[], "unauthorized").await;
});
let err = connect_remote_control_websocket(
&remote_control_target,
Some(state_db.as_ref()),
@@ -0,0 +1,366 @@
use anyhow::Result;
use app_test_support::McpProcess;
use app_test_support::to_response;
use codex_app_server_protocol::JSONRPCResponse;
use codex_app_server_protocol::RequestId;
use codex_app_server_protocol::ThreadStartParams;
use codex_app_server_protocol::ThreadStartResponse;
use codex_app_server_protocol::TurnStartParams;
use codex_app_server_protocol::TurnStartResponse;
use codex_app_server_protocol::TurnSteerParams;
use codex_app_server_protocol::TurnSteerResponse;
use codex_app_server_protocol::UserInput as V2UserInput;
use core_test_support::responses;
use core_test_support::skip_if_no_network;
use pretty_assertions::assert_eq;
use std::collections::HashMap;
use std::path::Path;
use tempfile::TempDir;
use tokio::time::timeout;
const DEFAULT_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
#[tokio::test]
async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result<()> {
skip_if_no_network!(Ok(()));
let server = responses::start_mock_server().await;
let response_mock = responses::mount_sse_once(
&server,
responses::sse(vec![
responses::ev_response_created("resp-1"),
responses::ev_assistant_message("msg-1", "Done"),
responses::ev_completed("resp-1"),
]),
)
.await;
let codex_home = TempDir::new()?;
create_config_toml(
codex_home.path(),
&server.uri(),
/*supports_websockets*/ false,
)?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??;
let thread_req = mcp
.send_thread_start_request(ThreadStartParams::default())
.await?;
let thread_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(thread_req)),
)
.await??;
let ThreadStartResponse { thread, .. } = to_response::<ThreadStartResponse>(thread_resp)?;
let client_metadata = HashMap::from([
("fiber_run_id".to_string(), "fiber-start-123".to_string()),
("origin".to_string(), "gaas".to_string()),
]);
let turn_req = mcp
.send_turn_start_request(TurnStartParams {
thread_id: thread.id,
input: vec![V2UserInput::Text {
text: "Hello".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: Some(client_metadata.clone()),
..Default::default()
})
.await?;
let turn_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(turn_req)),
)
.await??;
let TurnStartResponse { turn } = to_response::<TurnStartResponse>(turn_resp)?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
let request = response_mock.single_request();
let metadata = request
.header("x-codex-turn-metadata")
.as_deref()
.map(parse_json_header)
.unwrap_or_else(|| panic!("missing x-codex-turn-metadata header"));
assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123"));
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
assert!(metadata.get("session_id").is_some());
Ok(())
}
#[tokio::test]
async fn turn_steer_updates_client_metadata_on_follow_up_responses_request_v2() -> Result<()> {
skip_if_no_network!(Ok(()));
let codex_home = TempDir::new()?;
let server = responses::start_mock_server().await;
let first_response = responses::sse_response(responses::sse(vec![
responses::ev_response_created("resp-1"),
responses::ev_assistant_message("msg-1", "Working"),
responses::ev_completed("resp-1"),
]))
.set_delay(std::time::Duration::from_secs(2));
let second_response = responses::sse_response(responses::sse(vec![
responses::ev_response_created("resp-2"),
responses::ev_assistant_message("msg-2", "Done"),
responses::ev_completed("resp-2"),
]));
let request_log =
responses::mount_response_sequence(&server, vec![first_response, second_response]).await;
create_config_toml(
codex_home.path(),
&server.uri(),
/*supports_websockets*/ false,
)?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??;
let thread_req = mcp
.send_thread_start_request(ThreadStartParams::default())
.await?;
let thread_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(thread_req)),
)
.await??;
let ThreadStartResponse { thread, .. } = to_response::<ThreadStartResponse>(thread_resp)?;
let start_metadata =
HashMap::from([("fiber_run_id".to_string(), "fiber-start-123".to_string())]);
let turn_req = mcp
.send_turn_start_request(TurnStartParams {
thread_id: thread.id.clone(),
input: vec![V2UserInput::Text {
text: "Run sleep".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: Some(start_metadata.clone()),
..Default::default()
})
.await?;
let turn_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(turn_req)),
)
.await??;
let TurnStartResponse { turn } = to_response::<TurnStartResponse>(turn_resp)?;
let turn_id = turn.id.clone();
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/started"),
)
.await??;
wait_for_request_count(&request_log, /*expected*/ 1).await?;
let steer_metadata = HashMap::from([
("fiber_run_id".to_string(), "fiber-steer-456".to_string()),
("origin".to_string(), "gaas".to_string()),
]);
let steer_req = mcp
.send_turn_steer_request(TurnSteerParams {
thread_id: thread.id.clone(),
input: vec![V2UserInput::Text {
text: "Focus on the failure".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: Some(steer_metadata.clone()),
expected_turn_id: turn_id.clone(),
})
.await?;
let steer_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(steer_req)),
)
.await??;
let _turn: TurnSteerResponse = to_response::<TurnSteerResponse>(steer_resp)?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
let requests = request_log.requests();
assert_eq!(requests.len(), 2);
let first_metadata = requests[0]
.header("x-codex-turn-metadata")
.as_deref()
.map(parse_json_header)
.unwrap_or_else(|| panic!("missing first x-codex-turn-metadata header"));
assert_eq!(
first_metadata["fiber_run_id"].as_str(),
Some("fiber-start-123")
);
assert_eq!(first_metadata["turn_id"].as_str(), Some(turn_id.as_str()));
let second_metadata = requests[1]
.header("x-codex-turn-metadata")
.as_deref()
.map(parse_json_header)
.unwrap_or_else(|| panic!("missing second x-codex-turn-metadata header"));
assert_eq!(
second_metadata["fiber_run_id"].as_str(),
Some("fiber-steer-456")
);
assert_eq!(second_metadata["origin"].as_str(), Some("gaas"));
assert_eq!(second_metadata["turn_id"].as_str(), Some(turn_id.as_str()));
Ok(())
}
#[tokio::test]
async fn turn_start_forwards_client_metadata_to_responses_websocket_request_body_v2() -> Result<()>
{
skip_if_no_network!(Ok(()));
let websocket_server = responses::start_websocket_server(vec![vec![
vec![
responses::ev_response_created("warm-1"),
responses::ev_completed("warm-1"),
],
vec![
responses::ev_response_created("resp-1"),
responses::ev_assistant_message("msg-1", "Done"),
responses::ev_completed("resp-1"),
],
]])
.await;
let codex_home = TempDir::new()?;
create_config_toml(
codex_home.path(),
&websocket_server.uri().replacen("ws://", "http://", 1),
/*supports_websockets*/ true,
)?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??;
let thread_req = mcp
.send_thread_start_request(ThreadStartParams::default())
.await?;
let thread_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(thread_req)),
)
.await??;
let ThreadStartResponse { thread, .. } = to_response::<ThreadStartResponse>(thread_resp)?;
let client_metadata = HashMap::from([
("fiber_run_id".to_string(), "fiber-start-123".to_string()),
("origin".to_string(), "gaas".to_string()),
]);
let turn_req = mcp
.send_turn_start_request(TurnStartParams {
thread_id: thread.id,
input: vec![V2UserInput::Text {
text: "Hello".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: Some(client_metadata),
..Default::default()
})
.await?;
let turn_resp: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(turn_req)),
)
.await??;
let TurnStartResponse { turn } = to_response::<TurnStartResponse>(turn_resp)?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
let warmup = websocket_server
.wait_for_request(/*connection_index*/ 0, /*request_index*/ 0)
.await
.body_json();
let request = websocket_server
.wait_for_request(/*connection_index*/ 0, /*request_index*/ 1)
.await
.body_json();
assert_eq!(warmup["type"].as_str(), Some("response.create"));
assert_eq!(warmup["generate"].as_bool(), Some(false));
assert_eq!(request["type"].as_str(), Some("response.create"));
assert_eq!(request["previous_response_id"].as_str(), Some("warm-1"));
let metadata = request["client_metadata"]["x-codex-turn-metadata"]
.as_str()
.map(parse_json_header)
.unwrap_or_else(|| panic!("missing websocket x-codex-turn-metadata client metadata"));
assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123"));
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
assert!(metadata.get("session_id").is_some());
websocket_server.shutdown().await;
Ok(())
}
fn create_config_toml(
codex_home: &Path,
server_uri: &str,
supports_websockets: bool,
) -> std::io::Result<()> {
let config_toml = codex_home.join("config.toml");
std::fs::write(
config_toml,
format!(
r#"
model = "mock-model"
approval_policy = "never"
sandbox_mode = "read-only"
model_provider = "mock_provider"
[model_providers.mock_provider]
name = "Mock provider for test"
base_url = "{server_uri}/v1"
wire_api = "responses"
request_max_retries = 0
stream_max_retries = 0
supports_websockets = {supports_websockets}
"#
),
)
}
fn parse_json_header(value: &str) -> serde_json::Value {
match serde_json::from_str(value) {
Ok(value) => value,
Err(err) => panic!("metadata header should be valid json: {err}"),
}
}
async fn wait_for_request_count(
request_log: &core_test_support::responses::ResponseMock,
expected: usize,
) -> Result<()> {
timeout(DEFAULT_READ_TIMEOUT, async {
loop {
if request_log.requests().len() >= expected {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await?;
Ok(())
}
@@ -1,6 +1,7 @@
mod account;
mod analytics;
mod app_list;
mod client_metadata;
mod collaboration_mode_list;
#[cfg(unix)]
mod command_exec;
@@ -1377,6 +1377,7 @@ async fn turn_start_updates_sandbox_and_cwd_between_turns_v2() -> Result<()> {
text: "first turn".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
cwd: Some(first_cwd.clone()),
approval_policy: Some(codex_app_server_protocol::AskForApproval::Never),
approvals_reviewer: None,
@@ -1416,6 +1417,7 @@ async fn turn_start_updates_sandbox_and_cwd_between_turns_v2() -> Result<()> {
text: "second turn".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
cwd: Some(second_cwd.clone()),
approval_policy: Some(codex_app_server_protocol::AskForApproval::Never),
approvals_reviewer: None,
@@ -57,6 +57,7 @@ async fn turn_steer_requires_active_turn() -> Result<()> {
text: "steer".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
expected_turn_id: "turn-does-not-exist".to_string(),
})
.await?;
@@ -145,6 +146,7 @@ async fn turn_steer_rejects_oversized_text_input() -> Result<()> {
text: oversized_input.clone(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
expected_turn_id: turn.id.clone(),
})
.await?;
@@ -247,6 +249,7 @@ async fn turn_steer_returns_active_turn_id() -> Result<()> {
text: "steer".to_string(),
text_elements: Vec::new(),
}],
responsesapi_client_metadata: None,
expected_turn_id: turn.id.clone(),
})
.await?;
+3
View File
@@ -425,6 +425,7 @@ async fn send_input_submits_user_message() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
},
);
let captured = harness
@@ -571,6 +572,7 @@ async fn spawn_agent_creates_thread_and_sends_prompt() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
},
);
let captured = harness
@@ -678,6 +680,7 @@ async fn spawn_agent_can_fork_parent_thread_history_with_sanitized_items() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
},
);
let captured = harness
+29 -3
View File
@@ -763,8 +763,11 @@ impl Codex {
&self,
input: Vec<UserInput>,
expected_turn_id: Option<&str>,
responsesapi_client_metadata: Option<HashMap<String, String>>,
) -> Result<String, SteerInputError> {
self.session.steer_input(input, expected_turn_id).await
self.session
.steer_input(input, expected_turn_id, responsesapi_client_metadata)
.await
}
pub(crate) async fn set_app_server_client_info(
@@ -2264,6 +2267,7 @@ impl Session {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
},
)
.await;
@@ -4136,6 +4140,7 @@ impl Session {
&self,
input: Vec<UserInput>,
expected_turn_id: Option<&str>,
responsesapi_client_metadata: Option<HashMap<String, String>>,
) -> Result<String, SteerInputError> {
if input.is_empty() {
return Err(SteerInputError::EmptyInput);
@@ -4174,6 +4179,15 @@ impl Session {
None => return Err(SteerInputError::NoActiveTurn(input)),
}
if let Some(responsesapi_client_metadata) = responsesapi_client_metadata
&& let Some((_, active_task)) = active_turn.tasks.first()
{
active_task
.turn_context
.turn_metadata_state
.set_responsesapi_client_metadata(responsesapi_client_metadata);
}
let mut turn_state = active_turn.turn_state.lock().await;
turn_state.push_pending_input(input.into());
turn_state.accept_mailbox_delivery_for_current_turn();
@@ -4942,7 +4956,7 @@ mod handlers {
}
pub async fn user_input_or_turn(sess: &Arc<Session>, sub_id: String, op: Op) {
let (items, updates) = match op {
let (items, updates, responsesapi_client_metadata) = match op {
Op::UserTurn {
cwd,
approval_policy,
@@ -4983,17 +4997,20 @@ mod handlers {
app_server_client_name: None,
app_server_client_version: None,
},
None,
)
}
Op::UserInput {
items,
final_output_json_schema,
responsesapi_client_metadata,
} => (
items,
SessionSettingsUpdate {
final_output_json_schema: Some(final_output_json_schema),
..Default::default()
},
responsesapi_client_metadata,
),
_ => unreachable!(),
};
@@ -5005,11 +5022,20 @@ mod handlers {
sess.maybe_emit_unknown_model_warning_for_turn(current_context.as_ref())
.await;
match sess
.steer_input(items.clone(), /*expected_turn_id*/ None)
.steer_input(
items.clone(),
/*expected_turn_id*/ None,
responsesapi_client_metadata.clone(),
)
.await
{
Ok(_) => current_context.session_telemetry.user_prompt(&items),
Err(SteerInputError::NoActiveTurn(items)) => {
if let Some(responsesapi_client_metadata) = responsesapi_client_metadata {
current_context
.turn_metadata_state
.set_responsesapi_client_metadata(responsesapi_client_metadata);
}
current_context.session_telemetry.user_prompt(&items);
sess.refresh_mcp_servers_if_requested(&current_context)
.await;
+1
View File
@@ -185,6 +185,7 @@ pub(crate) async fn run_codex_thread_one_shot(
io.submit(Op::UserInput {
items: input,
final_output_json_schema,
responsesapi_client_metadata: None,
})
.await?;
+23 -4
View File
@@ -1331,6 +1331,7 @@ async fn fork_startup_context_then_first_turn_diff_snapshot() -> anyhow::Result<
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1386,6 +1387,7 @@ async fn fork_startup_context_then_first_turn_diff_snapshot() -> anyhow::Result<
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&forked.thread, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -3279,6 +3281,7 @@ fn op_kind_distinguishes_turn_ops() {
Op::UserInput {
items: vec![],
final_output_json_schema: None,
responsesapi_client_metadata: None,
}
.kind(),
"user_input"
@@ -4857,7 +4860,9 @@ async fn steer_input_requires_active_turn() {
}];
let err = sess
.steer_input(input, /*expected_turn_id*/ None)
.steer_input(
input, /*expected_turn_id*/ None, /*responsesapi_client_metadata*/ None,
)
.await
.expect_err("steering without active turn should fail");
@@ -4886,7 +4891,11 @@ async fn steer_input_enforces_expected_turn_id() {
text_elements: Vec::new(),
}];
let err = sess
.steer_input(steer_input, Some("different-turn-id"))
.steer_input(
steer_input,
Some("different-turn-id"),
/*responsesapi_client_metadata*/ None,
)
.await
.expect_err("mismatched expected turn id should fail");
@@ -4928,7 +4937,11 @@ async fn steer_input_rejects_non_regular_turns() {
text_elements: Vec::new(),
}];
let err = sess
.steer_input(steer_input, /*expected_turn_id*/ None)
.steer_input(
steer_input,
/*expected_turn_id*/ None,
/*responsesapi_client_metadata*/ None,
)
.await
.expect_err("steering a non-regular turn should fail");
@@ -4960,7 +4973,11 @@ async fn steer_input_returns_active_turn_id() {
text_elements: Vec::new(),
}];
let turn_id = sess
.steer_input(steer_input, Some(&tc.sub_id))
.steer_input(
steer_input,
Some(&tc.sub_id),
/*responsesapi_client_metadata*/ None,
)
.await
.expect("steering with matching expected turn id should succeed");
@@ -5147,6 +5164,7 @@ async fn steered_input_reopens_mailbox_delivery_for_current_turn() {
text_elements: Vec::new(),
}],
Some(&tc.sub_id),
/*responsesapi_client_metadata*/ None,
)
.await
.expect("steered input should be accepted");
@@ -5191,6 +5209,7 @@ async fn stale_defer_mailbox_delivery_does_not_override_steered_input() {
text_elements: Vec::new(),
}],
Some(&tc.sub_id),
/*responsesapi_client_metadata*/ None,
)
.await
.expect("steered input should be accepted");
+5 -1
View File
@@ -23,6 +23,7 @@ use codex_protocol::protocol::TokenUsage;
use codex_protocol::protocol::W3cTraceContext;
use codex_protocol::user_input::UserInput;
use rmcp::model::ReadResourceRequestParams;
use std::collections::HashMap;
use std::path::PathBuf;
use tokio::sync::Mutex;
use tokio::sync::watch;
@@ -97,8 +98,11 @@ impl CodexThread {
&self,
input: Vec<UserInput>,
expected_turn_id: Option<&str>,
responsesapi_client_metadata: Option<HashMap<String, String>>,
) -> Result<String, SteerInputError> {
self.codex.steer_input(input, expected_turn_id).await
self.codex
.steer_input(input, expected_turn_id, responsesapi_client_metadata)
.await
}
pub async fn set_app_server_client_info(
@@ -2004,6 +2004,7 @@ async fn send_input_accepts_structured_items() {
},
],
final_output_json_schema: None,
responsesapi_client_metadata: None,
};
let captured = manager
.captured_ops()
+41 -4
View File
@@ -1,4 +1,5 @@
use std::collections::BTreeMap;
use std::collections::HashMap;
use std::path::Path;
use std::path::PathBuf;
use std::sync::Arc;
@@ -6,6 +7,7 @@ use std::sync::Mutex;
use std::sync::RwLock;
use serde::Serialize;
use serde_json::Value;
use tokio::task::JoinHandle;
use crate::sandbox_tags::sandbox_tag;
@@ -69,6 +71,20 @@ impl TurnMetadataBag {
}
}
fn merge_responsesapi_client_metadata(
header: &str,
responsesapi_client_metadata: Option<&HashMap<String, String>>,
) -> Option<String> {
let responsesapi_client_metadata = responsesapi_client_metadata?;
let mut metadata = serde_json::from_str::<serde_json::Map<String, Value>>(header).ok()?;
for (key, value) in responsesapi_client_metadata {
metadata
.entry(key.clone())
.or_insert_with(|| Value::String(value.clone()));
}
serde_json::to_string(&metadata).ok()
}
fn build_turn_metadata_bag(
session_id: Option<String>,
turn_id: Option<String>,
@@ -129,6 +145,7 @@ pub(crate) struct TurnMetadataState {
base_metadata: TurnMetadataBag,
base_header: String,
enriched_header: Arc<RwLock<Option<String>>>,
responsesapi_client_metadata: Arc<RwLock<Option<HashMap<String, String>>>>,
enrichment_task: Arc<Mutex<Option<JoinHandle<()>>>>,
}
@@ -159,21 +176,30 @@ impl TurnMetadataState {
base_metadata,
base_header,
enriched_header: Arc::new(RwLock::new(None)),
responsesapi_client_metadata: Arc::new(RwLock::new(None)),
enrichment_task: Arc::new(Mutex::new(None)),
}
}
pub(crate) fn current_header_value(&self) -> Option<String> {
if let Some(header) = self
let header = if let Some(header) = self
.enriched_header
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
.cloned()
{
return Some(header);
}
Some(self.base_header.clone())
header
} else {
self.base_header.clone()
};
let responsesapi_client_metadata = self
.responsesapi_client_metadata
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
merge_responsesapi_client_metadata(&header, responsesapi_client_metadata.as_ref())
.or(Some(header))
}
pub(crate) fn current_meta_value(&self) -> Option<serde_json::Value> {
@@ -181,6 +207,17 @@ impl TurnMetadataState {
.and_then(|header| serde_json::from_str(&header).ok())
}
pub(crate) fn set_responsesapi_client_metadata(
&self,
responsesapi_client_metadata: HashMap<String, String>,
) {
*self
.responsesapi_client_metadata
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(responsesapi_client_metadata);
}
pub(crate) fn spawn_git_enrichment_task(&self) {
if self.repo_root.is_none() {
return;
+27
View File
@@ -1,6 +1,7 @@
use super::*;
use serde_json::Value;
use std::collections::HashMap;
use tempfile::TempDir;
use tokio::process::Command;
@@ -83,3 +84,29 @@ fn turn_metadata_state_uses_platform_sandbox_tag() {
assert_eq!(sandbox_name, Some(expected_sandbox));
assert_eq!(session_id, Some("session-a"));
}
#[test]
fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields() {
let temp_dir = TempDir::new().expect("temp dir");
let cwd = temp_dir.path().to_path_buf();
let sandbox_policy = SandboxPolicy::new_read_only_policy();
let state = TurnMetadataState::new(
"session-a".to_string(),
"turn-a".to_string(),
cwd,
&sandbox_policy,
WindowsSandboxLevel::Disabled,
);
state.set_responsesapi_client_metadata(HashMap::from([
("fiber_run_id".to_string(), "fiber-123".to_string()),
("session_id".to_string(), "client-supplied".to_string()),
]));
let header = state.current_header_value().expect("header");
let json: Value = serde_json::from_str(&header).expect("json");
assert_eq!(json["fiber_run_id"].as_str(), Some("fiber-123"));
assert_eq!(json["session_id"].as_str(), Some("session-a"));
assert_eq!(json["turn_id"].as_str(), Some("turn-a"));
}
+5
View File
@@ -51,6 +51,7 @@ async fn interrupt_long_running_tool_emits_turn_aborted() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -105,6 +106,7 @@ async fn interrupt_tool_records_history_entries() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -123,6 +125,7 @@ async fn interrupt_tool_records_history_entries() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -203,6 +206,7 @@ async fn interrupt_persists_turn_aborted_marker_in_next_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -221,6 +225,7 @@ async fn interrupt_persists_turn_aborted_marker_in_next_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
+31
View File
@@ -387,6 +387,7 @@ async fn resume_includes_initial_messages_and_sends_prior_items() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -749,6 +750,7 @@ async fn includes_conversation_id_and_model_headers_in_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -947,6 +949,7 @@ async fn includes_base_instructions_override_in_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1000,6 +1003,7 @@ async fn chatgpt_auth_sends_correct_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1112,6 +1116,7 @@ async fn prefers_apikey_when_config_prefers_apikey_even_with_chatgpt_tokens() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1148,6 +1153,7 @@ async fn includes_user_instructions_message_in_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1233,6 +1239,7 @@ async fn includes_apps_guidance_as_developer_message_for_chatgpt_auth() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1293,6 +1300,7 @@ async fn omits_apps_guidance_for_api_key_auth_even_when_feature_enabled() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1349,6 +1357,7 @@ async fn omits_apps_guidance_when_configured_off() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1388,6 +1397,7 @@ async fn omits_environment_context_when_configured_off() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1442,6 +1452,7 @@ async fn skills_append_to_developer_message() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1493,6 +1504,7 @@ async fn includes_configured_effort_in_request() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1535,6 +1547,7 @@ async fn includes_no_effort_in_request() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1575,6 +1588,7 @@ async fn includes_default_reasoning_effort_in_request_when_defined_by_model_info
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1687,6 +1701,7 @@ async fn configured_reasoning_summary_is_sent() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1802,6 +1817,7 @@ async fn reasoning_summary_is_omitted_when_disabled() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1858,6 +1874,7 @@ async fn reasoning_summary_none_overrides_model_catalog_default() -> anyhow::Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1894,6 +1911,7 @@ async fn includes_default_verbosity_in_request() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1939,6 +1957,7 @@ async fn configured_verbosity_not_sent_for_models_without_support() -> anyhow::R
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1983,6 +2002,7 @@ async fn configured_verbosity_is_sent() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2032,6 +2052,7 @@ async fn includes_developer_instructions_message_in_request() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2322,6 +2343,7 @@ async fn token_count_includes_rate_limits_snapshot() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2487,6 +2509,7 @@ async fn usage_limit_error_emits_rate_limit_event() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submission should succeed while emitting usage limit error events");
@@ -2561,6 +2584,7 @@ async fn context_window_error_sets_total_tokens_to_model_window() -> anyhow::Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -2573,6 +2597,7 @@ async fn context_window_error_sets_total_tokens_to_model_window() -> anyhow::Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -2655,6 +2680,7 @@ async fn incomplete_response_emits_content_filter_error_message() -> anyhow::Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -2762,6 +2788,7 @@ async fn azure_overrides_assign_properties_used_for_responses_url() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2847,6 +2874,7 @@ async fn env_var_overrides_loaded_auth() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2908,6 +2936,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2921,6 +2950,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2934,6 +2964,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1003,6 +1003,7 @@ async fn responses_websocket_usage_limit_error_emits_rate_limit_event() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submission should succeed while emitting usage limit error events");
@@ -1088,6 +1089,7 @@ async fn responses_websocket_invalid_request_error_with_status_is_forwarded() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submission should succeed while emitting invalid request events");
@@ -1265,6 +1267,56 @@ async fn responses_websocket_forwards_turn_metadata_on_initial_and_incremental_c
server.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn responses_websocket_preserves_custom_turn_metadata_fields() {
skip_if_no_network!();
let server = start_websocket_server(vec![vec![vec![
ev_response_created("resp-1"),
ev_completed("resp-1"),
]]])
.await;
let harness = websocket_harness(&server).await;
let mut client_session = harness.client.new_session();
let prompt = prompt_with_input(vec![message_item("hello")]);
let turn_metadata = json!({
"turn_id": "turn-123",
"fiber_run_id": "fiber-123",
"origin": "app-server",
})
.to_string();
stream_until_complete_with_turn_metadata(
&mut client_session,
&harness,
&prompt,
/*service_tier*/ None,
Some(&turn_metadata),
)
.await;
let body = server
.single_connection()
.first()
.expect("missing request")
.body_json();
assert_eq!(body["type"].as_str(), Some("response.create"));
assert_eq!(
body["client_metadata"]["x-codex-turn-metadata"]
.as_str()
.map(|value| serde_json::from_str::<serde_json::Value>(value).expect("valid json")),
Some(json!({
"turn_id": "turn-123",
"fiber_run_id": "fiber-123",
"origin": "app-server",
}))
);
server.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn responses_websocket_uses_previous_response_id_when_prefix_after_completed() {
skip_if_no_network!();
@@ -1817,6 +1869,23 @@ async fn stream_until_complete_with_turn_metadata(
prompt: &Prompt,
service_tier: Option<ServiceTier>,
turn_metadata_header: Option<&str>,
) {
stream_until_complete_with_request_metadata(
client_session,
harness,
prompt,
service_tier,
turn_metadata_header,
)
.await;
}
async fn stream_until_complete_with_request_metadata(
client_session: &mut ModelClientSession,
harness: &WebsocketTestHarness,
prompt: &Prompt,
service_tier: Option<ServiceTier>,
turn_metadata_header: Option<&str>,
) {
let mut stream = client_session
.stream(
@@ -84,6 +84,7 @@ async fn no_collaboration_instructions_by_default() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -143,6 +144,7 @@ async fn user_input_includes_collaboration_instructions_after_override() -> Resu
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -241,6 +243,7 @@ async fn override_then_next_turn_uses_updated_collaboration_instructions() -> Re
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -364,6 +367,7 @@ async fn collaboration_mode_update_emits_new_instruction_message() -> Result<()>
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -391,6 +395,7 @@ async fn collaboration_mode_update_emits_new_instruction_message() -> Result<()>
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -447,6 +452,7 @@ async fn collaboration_mode_update_noop_does_not_append() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -474,6 +480,7 @@ async fn collaboration_mode_update_noop_does_not_append() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -532,6 +539,7 @@ async fn collaboration_mode_update_emits_new_instruction_message_when_mode_chang
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -562,6 +570,7 @@ async fn collaboration_mode_update_emits_new_instruction_message_when_mode_chang
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -621,6 +630,7 @@ async fn collaboration_mode_update_noop_does_not_append_when_mode_is_unchanged()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -651,6 +661,7 @@ async fn collaboration_mode_update_noop_does_not_append_when_mode_is_unchanged()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -714,6 +725,7 @@ async fn resume_replays_collaboration_instructions() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -727,6 +739,7 @@ async fn resume_replays_collaboration_instructions() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -783,6 +796,7 @@ async fn empty_collaboration_instructions_are_ignored() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
+29
View File
@@ -245,6 +245,7 @@ async fn summarize_context_three_requests_and_instructions() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -267,6 +268,7 @@ async fn summarize_context_three_requests_and_instructions() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -443,6 +445,7 @@ async fn manual_compact_uses_custom_prompt() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit first user turn");
@@ -587,6 +590,7 @@ async fn manual_compact_emits_context_compaction_items() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -750,6 +754,7 @@ async fn multiple_auto_compact_per_task_runs_after_token_limit_hit() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit user input");
@@ -1249,6 +1254,7 @@ async fn auto_compact_runs_after_token_limit_hit() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1262,6 +1268,7 @@ async fn auto_compact_runs_after_token_limit_hit() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1275,6 +1282,7 @@ async fn auto_compact_runs_after_token_limit_hit() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1443,6 +1451,7 @@ async fn auto_compact_emits_context_compaction_items() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1521,6 +1530,7 @@ async fn auto_compact_starts_after_turn_started() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1533,6 +1543,7 @@ async fn auto_compact_starts_after_turn_started() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1545,6 +1556,7 @@ async fn auto_compact_starts_after_turn_started() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2037,6 +2049,7 @@ async fn auto_compact_persists_rollout_entries() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2049,6 +2062,7 @@ async fn auto_compact_persists_rollout_entries() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2061,6 +2075,7 @@ async fn auto_compact_persists_rollout_entries() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2147,6 +2162,7 @@ async fn manual_compact_retries_after_context_window_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2258,6 +2274,7 @@ async fn manual_compact_non_context_failure_retries_then_emits_task_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit user input");
@@ -2350,6 +2367,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2365,6 +2383,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2380,6 +2399,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2541,6 +2561,7 @@ async fn auto_compact_allows_multiple_attempts_when_interleaved_with_other_turn_
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2643,6 +2664,7 @@ async fn snapshot_request_shape_mid_turn_continuation_compaction() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2839,6 +2861,7 @@ async fn auto_compact_counts_encrypted_reasoning_before_last_user() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2956,6 +2979,7 @@ async fn auto_compact_runs_when_reasoning_header_clears_between_turns() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -3015,6 +3039,7 @@ async fn snapshot_request_shape_pre_turn_compaction_including_incoming_user_mess
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit user input");
@@ -3050,6 +3075,7 @@ async fn snapshot_request_shape_pre_turn_compaction_including_incoming_user_mess
},
],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit user input");
@@ -3260,6 +3286,7 @@ async fn snapshot_request_shape_pre_turn_compaction_context_window_exceeded() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit first user");
@@ -3272,6 +3299,7 @@ async fn snapshot_request_shape_pre_turn_compaction_context_window_exceeded() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit second user");
@@ -3342,6 +3370,7 @@ async fn snapshot_request_shape_manual_compact_without_previous_user_messages()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit follow-up user input");
@@ -237,6 +237,7 @@ async fn remote_compact_replaces_history_for_followups() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -251,6 +252,7 @@ async fn remote_compact_replaces_history_for_followups() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -381,6 +383,7 @@ async fn remote_compact_runs_automatically() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -455,6 +458,7 @@ async fn remote_compact_trims_function_call_history_to_fit_context_window() -> R
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -466,6 +470,7 @@ async fn remote_compact_trims_function_call_history_to_fit_context_window() -> R
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -581,6 +586,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -592,6 +598,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -609,6 +616,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -706,6 +714,7 @@ async fn auto_remote_compact_failure_stops_agent_loop() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -717,6 +726,7 @@ async fn auto_remote_compact_failure_stops_agent_loop() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -808,6 +818,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&baseline_codex, |event| {
@@ -822,6 +833,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&baseline_codex, |event| {
@@ -907,6 +919,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&override_codex, |event| {
@@ -921,6 +934,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&override_codex, |event| {
@@ -989,6 +1003,7 @@ async fn remote_manual_compact_emits_context_compaction_items() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -1067,6 +1082,7 @@ async fn remote_manual_compact_failure_emits_task_error_event() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -1149,6 +1165,7 @@ async fn remote_compact_persists_replacement_history_in_rollout() -> Result<()>
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1290,6 +1307,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1305,6 +1323,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1327,6 +1346,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1421,6 +1441,7 @@ async fn remote_compact_refreshes_stale_developer_instructions_without_resume()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1435,6 +1456,7 @@ async fn remote_compact_refreshes_stale_developer_instructions_without_resume()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1504,6 +1526,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_sta
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1515,6 +1538,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_sta
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1579,6 +1603,7 @@ async fn remote_request_uses_custom_experimental_realtime_start_instructions() -
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1637,6 +1662,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_end
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1650,6 +1676,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_end
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1722,6 +1749,7 @@ async fn snapshot_request_shape_remote_manual_compact_restates_realtime_start()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1736,6 +1764,7 @@ async fn snapshot_request_shape_remote_manual_compact_restates_realtime_start()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1816,6 +1845,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_does_not_restate_real
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1829,6 +1859,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_does_not_restate_real
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1917,6 +1948,7 @@ async fn snapshot_request_shape_remote_compact_resume_restates_realtime_end() ->
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -1944,6 +1976,7 @@ async fn snapshot_request_shape_remote_compact_resume_restates_realtime_end() ->
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2037,6 +2070,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_including_incoming_us
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2121,6 +2155,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_strips_incoming_model
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2147,6 +2182,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_strips_incoming_model
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2263,6 +2299,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_context_window_exceed
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2274,6 +2311,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_context_window_exceed
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
let error_message = wait_for_event_match(&codex, |event| match event {
@@ -2356,6 +2394,7 @@ async fn snapshot_request_shape_remote_mid_turn_continuation_compaction() -> Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2431,6 +2470,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_summary_only_reinject
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2514,6 +2554,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_multi_summary_reinjec
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2528,6 +2569,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_multi_summary_reinjec
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -2607,6 +2649,7 @@ async fn snapshot_request_shape_remote_manual_compact_without_previous_user_mess
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -808,6 +808,7 @@ async fn user_turn(conversation: &Arc<CodexThread>, text: &str) {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.expect("submit user turn");
+1
View File
@@ -53,6 +53,7 @@ async fn fork_thread_twice_drops_to_first_message() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
+2
View File
@@ -904,6 +904,7 @@ async fn blocked_queued_prompt_does_not_strand_earlier_accepted_prompt() -> Resu
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -920,6 +921,7 @@ async fn blocked_queued_prompt_does_not_strand_earlier_accepted_prompt() -> Resu
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
}
+9
View File
@@ -86,6 +86,7 @@ async fn user_message_item_is_emitted() -> anyhow::Result<()> {
.submit(Op::UserInput {
items: vec![expected_input.clone()],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -142,6 +143,7 @@ async fn assistant_message_item_is_emitted() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -200,6 +202,7 @@ async fn reasoning_item_is_emitted() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -259,6 +262,7 @@ async fn web_search_item_is_emitted() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -323,6 +327,7 @@ async fn image_generation_call_event_is_emitted() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -385,6 +390,7 @@ async fn image_generation_call_event_is_emitted_when_image_save_fails() -> anyho
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -438,6 +444,7 @@ async fn agent_message_content_delta_has_item_metadata() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -1085,6 +1092,7 @@ async fn reasoning_content_delta_has_item_metadata() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -1144,6 +1152,7 @@ async fn reasoning_raw_content_delta_respects_flag() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -322,6 +322,7 @@ async fn snapshot_model_visible_layout_resume_with_personality_change() -> Resul
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -421,6 +422,7 @@ async fn snapshot_model_visible_layout_resume_override_matches_rollout_model() -
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -466,6 +468,7 @@ async fn snapshot_model_visible_layout_resume_override_matches_rollout_model() -
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |event| {
+22
View File
@@ -107,6 +107,7 @@ async fn responses_api_emits_api_request_event() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -150,6 +151,7 @@ async fn process_sse_emits_tracing_for_output_item() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -193,6 +195,7 @@ async fn process_sse_emits_failed_event_on_parse_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -237,6 +240,7 @@ async fn process_sse_records_failed_event_when_stream_closes_without_completed()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -301,6 +305,7 @@ async fn process_sse_failed_event_records_response_error_message() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -363,6 +368,7 @@ async fn process_sse_failed_event_logs_parse_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -412,6 +418,7 @@ async fn process_sse_failed_event_logs_missing_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -470,6 +477,7 @@ async fn process_sse_failed_event_logs_response_completed_parse_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -522,6 +530,7 @@ async fn process_sse_emits_completed_telemetry() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -594,6 +603,7 @@ async fn handle_responses_span_records_response_kind_and_tool_name() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -678,6 +688,7 @@ async fn record_responses_sets_span_fields_for_response_events() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -762,6 +773,7 @@ async fn handle_response_item_records_tool_result_for_custom_tool_call() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -835,6 +847,7 @@ async fn handle_response_item_records_tool_result_for_function_call() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -918,6 +931,7 @@ async fn handle_response_item_records_tool_result_for_local_shell_missing_ids()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -986,6 +1000,7 @@ async fn handle_response_item_records_tool_result_for_local_shell_call() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1094,6 +1109,7 @@ async fn handle_container_exec_autoapprove_from_config_records_tool_decision() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1145,6 +1161,7 @@ async fn handle_container_exec_user_approved_records_tool_decision() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1211,6 +1228,7 @@ async fn handle_container_exec_user_approved_for_session_records_tool_decision()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1277,6 +1295,7 @@ async fn handle_sandbox_error_user_approves_retry_records_tool_decision() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1343,6 +1362,7 @@ async fn handle_container_exec_user_denies_records_tool_decision() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1409,6 +1429,7 @@ async fn handle_sandbox_error_user_approves_for_session_records_tool_decision()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -1476,6 +1497,7 @@ async fn handle_sandbox_error_user_denies_records_tool_decision() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -100,6 +100,7 @@ async fn submit_user_input(codex: &CodexThread, text: &str) {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap_or_else(|err| panic!("submit user input: {err}"));
@@ -136,6 +137,7 @@ async fn steer_user_input(codex: &CodexThread, text: &str) {
text_elements: Vec::new(),
}],
/*expected_turn_id*/ None,
/*responsesapi_client_metadata*/ None,
)
.await
.unwrap_or_else(|err| panic!("steer user input: {err:?}"));
@@ -267,6 +269,7 @@ async fn injected_user_input_triggers_follow_up_request_with_deltas() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -283,6 +286,7 @@ async fn injected_user_input_triggers_follow_up_request_with_deltas() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -53,6 +53,7 @@ async fn permissions_message_sent_once_on_start() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -90,6 +91,7 @@ async fn permissions_message_added_on_override_change() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -117,6 +119,7 @@ async fn permissions_message_added_on_override_change() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -160,6 +163,7 @@ async fn permissions_message_not_added_when_no_change() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -171,6 +175,7 @@ async fn permissions_message_not_added_when_no_change() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -214,6 +219,7 @@ async fn permissions_message_omitted_when_disabled() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -241,6 +247,7 @@ async fn permissions_message_omitted_when_disabled() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -297,6 +304,7 @@ async fn resume_replays_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -326,6 +334,7 @@ async fn resume_replays_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -339,6 +348,7 @@ async fn resume_replays_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -396,6 +406,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -425,6 +436,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -444,6 +456,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -476,6 +489,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&forked.thread, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -528,6 +542,7 @@ async fn permissions_message_includes_writable_roots() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
+3
View File
@@ -213,6 +213,7 @@ async fn capability_sections_render_in_developer_message_in_order() -> Result<()
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -288,6 +289,7 @@ async fn explicit_plugin_mentions_inject_plugin_guidance() -> Result<()> {
path: format!("plugin://{SAMPLE_PLUGIN_CONFIG_NAME}"),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -360,6 +362,7 @@ async fn explicit_plugin_mentions_track_plugin_used_analytics() -> Result<()> {
path: format!("plugin://{SAMPLE_PLUGIN_CONFIG_NAME}"),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -151,6 +151,7 @@ async fn prompt_tools_are_consistent_across_requests() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -162,6 +163,7 @@ async fn prompt_tools_are_consistent_across_requests() -> anyhow::Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -247,6 +249,7 @@ async fn gpt_5_tools_without_apply_patch_append_apply_patch_instructions() -> an
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -258,6 +261,7 @@ async fn gpt_5_tools_without_apply_patch_append_apply_patch_instructions() -> an
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -320,6 +324,7 @@ async fn prefixes_context_and_instructions_once_and_consistently_across_requests
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -331,6 +336,7 @@ async fn prefixes_context_and_instructions_once_and_consistently_across_requests
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -412,6 +418,7 @@ async fn overrides_turn_context_but_keeps_cached_prefix_and_key_constant() -> an
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -448,6 +455,7 @@ async fn overrides_turn_context_but_keeps_cached_prefix_and_key_constant() -> an
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -530,6 +538,7 @@ async fn override_before_first_turn_emits_environment_context() -> anyhow::Resul
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -681,6 +690,7 @@ async fn per_turn_overrides_keep_cached_prefix_and_key_constant() -> anyhow::Res
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -46,6 +46,7 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -2650,6 +2650,7 @@ async fn inbound_handoff_request_steers_active_turn() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -45,6 +45,7 @@ async fn request_body_is_zstd_compressed_for_codex_backend_when_enabled() -> any
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -92,6 +93,7 @@ async fn request_body_is_not_compressed_for_api_key_auth_even_when_enabled() ->
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
+7
View File
@@ -91,6 +91,7 @@ async fn resume_includes_initial_messages_from_rollout_events() -> Result<()> {
text_elements: text_elements.clone(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -176,6 +177,7 @@ async fn resume_includes_initial_messages_from_reasoning_events() -> Result<()>
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
@@ -265,6 +267,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -305,6 +308,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |event| {
@@ -320,6 +324,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |event| {
@@ -390,6 +395,7 @@ async fn resume_model_switch_is_not_duplicated_after_pre_turn_override() -> Resu
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
@@ -433,6 +439,7 @@ async fn resume_model_switch_is_not_duplicated_after_pre_turn_override() -> Resu
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&resumed.codex, |event| {
+1
View File
@@ -726,6 +726,7 @@ async fn review_history_surfaces_in_parent_session() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
+1
View File
@@ -427,6 +427,7 @@ async fn tool_search_returns_deferred_tools_without_follow_up_tool_injection() -
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
+7 -6
View File
@@ -463,16 +463,17 @@ async fn tool_call_logs_include_thread_id() -> Result<()> {
let db = test.codex.state_db().expect("state db enabled");
let expected_thread_id = test.session_configured.session_id.to_string();
let subscriber = tracing_subscriber::registry().with(codex_state::log_db::start(db.clone()));
let dispatch = tracing::Dispatch::new(subscriber);
let _guard = tracing::dispatcher::set_default(&dispatch);
test.submit_turn("run a shell command").await?;
{
let log_db_layer = codex_state::log_db::start(db.clone());
let subscriber = tracing_subscriber::registry().with(log_db_layer.clone());
let dispatch = tracing::Dispatch::new(subscriber);
tracing::dispatcher::with_default(&dispatch, || {
let span = tracing::info_span!("test_log_span", thread_id = %expected_thread_id);
let _entered = span.enter();
tracing::info!("ToolCall: shell_command {{\"command\":\"echo hello\"}}");
}
});
log_db_layer.flush().await;
let mut found = None;
for _ in 0..80 {
@@ -98,6 +98,7 @@ async fn continue_after_stream_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -117,6 +118,7 @@ async fn continue_after_stream_error() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -82,6 +82,7 @@ async fn retries_on_early_close() {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
.unwrap();
@@ -62,6 +62,7 @@ mv "${tmp_path}" "${payload_path}""#,
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;
@@ -109,6 +109,7 @@ async fn submit_user_turn(codex: &Arc<CodexThread>, text: &str) -> Result<()> {
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await?;
wait_for_event(codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
+1
View File
@@ -717,6 +717,7 @@ async fn run_exec_session(args: ExecRunArgs) -> anyhow::Result<()> {
params: TurnStartParams {
thread_id: primary_thread_id_for_span.clone(),
input: items.into_iter().map(Into::into).collect(),
responsesapi_client_metadata: None,
cwd: Some(default_cwd),
approval_policy: Some(default_approval_policy.into()),
approvals_reviewer: None,
@@ -114,6 +114,7 @@ pub async fn run_codex_tool_session(
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
},
trace: None,
};
@@ -161,6 +162,7 @@ pub async fn run_codex_tool_session_reply(
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
})
.await
{
+34
View File
@@ -381,6 +381,9 @@ pub enum Op {
/// Optional JSON Schema used to constrain the final assistant message for this turn.
#[serde(skip_serializing_if = "Option::is_none")]
final_output_json_schema: Option<Value>,
/// Optional turn-scoped Responses API `client_metadata`.
#[serde(default, skip_serializing_if = "Option::is_none")]
responsesapi_client_metadata: Option<HashMap<String, String>>,
},
/// Similar to [`Op::UserInput`], but contains additional context required
@@ -654,6 +657,7 @@ impl From<Vec<UserInput>> for Op {
Op::UserInput {
items: value,
final_output_json_schema: None,
responsesapi_client_metadata: None,
}
}
}
@@ -4704,6 +4708,7 @@ mod tests {
let op = Op::UserInput {
items: Vec::new(),
final_output_json_schema: None,
responsesapi_client_metadata: None,
};
let json_op = serde_json::to_value(op)?;
@@ -4721,6 +4726,7 @@ mod tests {
Op::UserInput {
items: Vec::new(),
final_output_json_schema: None,
responsesapi_client_metadata: None,
}
);
@@ -4740,6 +4746,7 @@ mod tests {
let op = Op::UserInput {
items: Vec::new(),
final_output_json_schema: Some(schema.clone()),
responsesapi_client_metadata: None,
};
let json_op = serde_json::to_value(op)?;
@@ -4755,6 +4762,33 @@ mod tests {
Ok(())
}
#[test]
fn user_input_with_responsesapi_client_metadata_round_trips() -> Result<()> {
let op = Op::UserInput {
items: Vec::new(),
final_output_json_schema: None,
responsesapi_client_metadata: Some(HashMap::from([(
"fiber_run_id".to_string(),
"fiber-123".to_string(),
)])),
};
let json_op = serde_json::to_value(&op)?;
assert_eq!(
json_op,
json!({
"type": "user_input",
"items": [],
"responsesapi_client_metadata": {
"fiber_run_id": "fiber-123",
}
})
);
assert_eq!(serde_json::from_value::<Op>(json_op)?, op);
Ok(())
}
#[test]
fn user_input_text_serializes_empty_text_elements() -> Result<()> {
let input = UserInput::Text {
+1
View File
@@ -638,6 +638,7 @@ mod tests {
.await
.expect("insert legacy log row");
pool.close().await;
drop(pool);
let runtime = StateRuntime::init(codex_home.clone(), "test-provider".to_string())
.await
+2
View File
@@ -421,6 +421,7 @@ impl AppServerSession {
params: TurnStartParams {
thread_id: thread_id.to_string(),
input: items.into_iter().map(Into::into).collect(),
responsesapi_client_metadata: None,
cwd: Some(cwd),
approval_policy: Some(approval_policy.into()),
approvals_reviewer: Some(approvals_reviewer.into()),
@@ -471,6 +472,7 @@ impl AppServerSession {
params: TurnSteerParams {
thread_id: thread_id.to_string(),
input: items.into_iter().map(Into::into).collect(),
responsesapi_client_metadata: None,
expected_turn_id: turn_id,
},
})