| 1 | use super::*; |
| 2 | |
| 3 | use std::sync::Arc; |
| 4 | use std::sync::atomic::{AtomicUsize, Ordering}; |
| 5 | |
| 6 | use futures_util::StreamExt; |
| 7 | |
| 8 | use crate::config::{Config, ProviderConfig, ProvidersConfig, RetryConfig}; |
| 9 | use codewhale_models::Message; |
| 10 | use codewhale_models::Role; |
| 11 | use codewhale_models::SystemPrompt; |
| 12 | use wiremock::matchers::{method, path}; |
| 13 | use wiremock::{Mock, MockServer, Request, Respond, ResponseTemplate}; |
| 14 | |
| 15 | #[derive(Clone)] |
| 16 | struct RetryThenSuccess { |
| 17 | attempts: Arc<AtomicUsize>, |
| 18 | retry_status: u16, |
| 19 | retry_body: &'static str, |
| 20 | } |
| 21 | |
| 22 | impl Respond for RetryThenSuccess { |
| 23 | fn respond(&self, _request: &Request) -> ResponseTemplate { |
| 24 | if self.attempts.fetch_add(1, Ordering::SeqCst) == 0 { |
| 25 | let mut response = |
| 26 | ResponseTemplate::new(self.retry_status).set_body_string(self.retry_body); |
| 27 | if self.retry_status == 429 { |
| 28 | response = response.insert_header("Retry-After", "0"); |
| 29 | } |
| 30 | return response; |
| 31 | } |
| 32 | |
| 33 | ResponseTemplate::new(200) |
| 34 | .insert_header("Content-Type", "text/event-stream") |
| 35 | .set_body_string("data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n") |
| 36 | } |
| 37 | } |
| 38 | |
| 39 | #[derive(Clone)] |
| 40 | struct AlwaysError { |
| 41 | attempts: Arc<AtomicUsize>, |
| 42 | status: u16, |
| 43 | body: &'static str, |
| 44 | } |
| 45 | |
| 46 | impl Respond for AlwaysError { |
| 47 | fn respond(&self, _request: &Request) -> ResponseTemplate { |
| 48 | self.attempts.fetch_add(1, Ordering::SeqCst); |
| 49 | ResponseTemplate::new(self.status).set_body_string(self.body) |
| 50 | } |
| 51 | } |
| 52 | |
| 53 | fn minimal_responses_request() -> MessageRequest { |
| 54 | MessageRequest { |
| 55 | model: "gpt-5.5".to_string(), |
| 56 | messages: vec![Message { |
| 57 | role: Role::User, |
| 58 | content: vec![ContentBlock::Text { |
| 59 | text: "hello".to_string(), |
| 60 | cache_control: None, |
| 61 | }], |
| 62 | }], |
| 63 | max_tokens: 128, |
| 64 | system: None, |
| 65 | tools: None, |
| 66 | tool_choice: None, |
| 67 | metadata: None, |
| 68 | thinking: None, |
| 69 | reasoning_effort: None, |
| 70 | stream: None, |
| 71 | temperature: None, |
| 72 | top_p: None, |
| 73 | } |
| 74 | } |
| 75 | |
| 76 | fn test_codex_config(server: &MockServer) -> Config { |
| 77 | Config { |
| 78 | provider: Some("openai-codex".to_string()), |
| 79 | retry: Some(RetryConfig { |
| 80 | enabled: Some(true), |
| 81 | max_retries: Some(1), |
| 82 | initial_delay: Some(0.0), |
| 83 | max_delay: Some(0.0), |
| 84 | exponential_base: Some(1.0), |
| 85 | jitter: None, |
| 86 | jitter_factor: None, |
| 87 | respect_retry_after: None, |
| 88 | }), |
| 89 | providers: Some(ProvidersConfig { |
| 90 | openai_codex: ProviderConfig { |
| 91 | base_url: Some(format!("{}/v1", server.uri())), |
| 92 | api_key: Some("test-token".to_string()), |
| 93 | ..ProviderConfig::default() |
| 94 | }, |
| 95 | ..ProvidersConfig::default() |
| 96 | }), |
| 97 | ..Config::default() |
| 98 | } |
| 99 | } |
| 100 | |
| 101 | #[tokio::test] |
| 102 | async fn responses_stream_retries_rate_limited_request() { |
| 103 | let server = MockServer::start().await; |
| 104 | let attempts = Arc::new(AtomicUsize::new(0)); |
| 105 | Mock::given(method("POST")) |
| 106 | .and(path("/v1/responses")) |
| 107 | .respond_with(RetryThenSuccess { |
| 108 | attempts: Arc::clone(&attempts), |
| 109 | retry_status: 429, |
| 110 | retry_body: "rate limited", |
| 111 | }) |
| 112 | .mount(&server) |
| 113 | .await; |
| 114 | |
| 115 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 116 | let mut request = minimal_responses_request(); |
| 117 | request.max_tokens = 384_000; |
| 118 | let prepared = client |
| 119 | .prepare_outbound_request(request, true) |
| 120 | .expect("responses request prepares"); |
| 121 | assert_eq!( |
| 122 | prepared.endpoint.url, |
| 123 | format!("{}/v1/responses", server.uri()) |
| 124 | ); |
| 125 | // The official ChatGPT plan preview does not support output-cap fields; |
| 126 | // omit them while retaining the allowance in the resolved envelope. |
| 127 | assert!(prepared.body.get("max_output_tokens").is_none()); |
| 128 | let mut stream = client.handle_responses_stream(&prepared).await.unwrap(); |
| 129 | |
| 130 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 131 | while let Some(event) = stream.next().await { |
| 132 | event.unwrap(); |
| 133 | } |
| 134 | }) |
| 135 | .await |
| 136 | .expect("Responses retry stream should finish after response.completed"); |
| 137 | |
| 138 | assert_eq!(attempts.load(Ordering::SeqCst), 2); |
| 139 | let requests = server |
| 140 | .received_requests() |
| 141 | .await |
| 142 | .expect("recorded retry requests"); |
| 143 | assert_eq!(requests.len(), 2); |
| 144 | for request in requests { |
| 145 | let body: Value = serde_json::from_slice(&request.body).expect("Responses JSON"); |
| 146 | assert!( |
| 147 | body.get("max_output_tokens").is_none(), |
| 148 | "ChatGPT plan body must not name the unsupported output cap: {body}" |
| 149 | ); |
| 150 | } |
| 151 | } |
| 152 | |
| 153 | #[tokio::test] |
| 154 | async fn responses_stream_retries_transient_server_error() { |
| 155 | let server = MockServer::start().await; |
| 156 | let attempts = Arc::new(AtomicUsize::new(0)); |
| 157 | Mock::given(method("POST")) |
| 158 | .and(path("/v1/responses")) |
| 159 | .respond_with(RetryThenSuccess { |
| 160 | attempts: Arc::clone(&attempts), |
| 161 | retry_status: 503, |
| 162 | retry_body: "temporarily unavailable", |
| 163 | }) |
| 164 | .mount(&server) |
| 165 | .await; |
| 166 | |
| 167 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 168 | let mut stream = client |
| 169 | .handle_responses_stream( |
| 170 | &client |
| 171 | .prepare_outbound_request(minimal_responses_request(), true) |
| 172 | .expect("responses request prepares"), |
| 173 | ) |
| 174 | .await |
| 175 | .unwrap(); |
| 176 | |
| 177 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 178 | while let Some(event) = stream.next().await { |
| 179 | event.unwrap(); |
| 180 | } |
| 181 | }) |
| 182 | .await |
| 183 | .expect("Responses retry stream should finish after response.completed"); |
| 184 | |
| 185 | assert_eq!(attempts.load(Ordering::SeqCst), 2); |
| 186 | } |
| 187 | |
| 188 | #[tokio::test] |
| 189 | async fn responses_stream_retries_upstream_499_before_streaming() { |
| 190 | let server = MockServer::start().await; |
| 191 | let attempts = Arc::new(AtomicUsize::new(0)); |
| 192 | Mock::given(method("POST")) |
| 193 | .and(path("/v1/responses")) |
| 194 | .respond_with(RetryThenSuccess { |
| 195 | attempts: Arc::clone(&attempts), |
| 196 | retry_status: 499, |
| 197 | retry_body: "upstream request cancelled", |
| 198 | }) |
| 199 | .mount(&server) |
| 200 | .await; |
| 201 | |
| 202 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 203 | let mut stream = client |
| 204 | .handle_responses_stream( |
| 205 | &client |
| 206 | .prepare_outbound_request(minimal_responses_request(), true) |
| 207 | .expect("responses request prepares"), |
| 208 | ) |
| 209 | .await |
| 210 | .unwrap(); |
| 211 | |
| 212 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 213 | while let Some(event) = stream.next().await { |
| 214 | event.unwrap(); |
| 215 | } |
| 216 | }) |
| 217 | .await |
| 218 | .expect("Responses retry stream should finish after response.completed"); |
| 219 | |
| 220 | assert_eq!(attempts.load(Ordering::SeqCst), 2); |
| 221 | } |
| 222 | |
| 223 | async fn collect_responses_stream(sse_body: &str) -> Vec<Result<StreamEvent>> { |
| 224 | let server = MockServer::start().await; |
| 225 | Mock::given(method("POST")) |
| 226 | .and(path("/v1/responses")) |
| 227 | .respond_with( |
| 228 | ResponseTemplate::new(200) |
| 229 | .insert_header("Content-Type", "text/event-stream") |
| 230 | .set_body_string(sse_body), |
| 231 | ) |
| 232 | .mount(&server) |
| 233 | .await; |
| 234 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 235 | let stream = client |
| 236 | .handle_responses_stream( |
| 237 | &client |
| 238 | .prepare_outbound_request(minimal_responses_request(), true) |
| 239 | .expect("responses request prepares"), |
| 240 | ) |
| 241 | .await |
| 242 | .expect("Responses stream opens"); |
| 243 | tokio::time::timeout(std::time::Duration::from_secs(5), stream.collect()) |
| 244 | .await |
| 245 | .expect("stream ends") |
| 246 | } |
| 247 | |
| 248 | #[tokio::test] |
| 249 | async fn responses_stream_eof_without_a_terminal_event_is_an_error() { |
| 250 | let events = collect_responses_stream(concat!( |
| 251 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"message\",\"id\":\"m\"}}\n\n", |
| 252 | "data: {\"type\":\"response.output_text.delta\",\"delta\":\"partial\"}\n\n", |
| 253 | )) |
| 254 | .await; |
| 255 | assert!( |
| 256 | !events |
| 257 | .iter() |
| 258 | .any(|event| matches!(event, Ok(StreamEvent::MessageStop))), |
| 259 | "a truncated stream must not report MessageStop: {events:?}" |
| 260 | ); |
| 261 | assert!( |
| 262 | events |
| 263 | .last() |
| 264 | .is_some_and(|event| event.as_ref().is_err_and(|error| error |
| 265 | .to_string() |
| 266 | .contains("closed before a successful completion event"))), |
| 267 | "{events:?}" |
| 268 | ); |
| 269 | } |
| 270 | |
| 271 | #[tokio::test] |
| 272 | async fn chatgpt_stream_requires_valid_successful_completion() { |
| 273 | for body in [ |
| 274 | "data: [DONE]\n\n", |
| 275 | "data: {\"type\":\"response.completed\"}\n\n", |
| 276 | "data: {\"type\":\"response.completed\",\"response\":{}}\n\n", |
| 277 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"failed\"}}\n\n", |
| 278 | "data: {\"type\":\"response.incomplete\",\"response\":{\"status\":\"incomplete\",\"incomplete_details\":{\"reason\":\"max_output_tokens\"}}}\n\n", |
| 279 | concat!( |
| 280 | "data: {invalid-json}\n\n", |
| 281 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 282 | ), |
| 283 | ] { |
| 284 | let events = collect_responses_stream(body).await; |
| 285 | assert!( |
| 286 | events.last().is_some_and(Result::is_err), |
| 287 | "{body}: {events:?}" |
| 288 | ); |
| 289 | assert!( |
| 290 | !events |
| 291 | .iter() |
| 292 | .any(|event| matches!(event, Ok(StreamEvent::MessageStop))), |
| 293 | "an invalid or incomplete response must not settle the turn: {body}: {events:?}" |
| 294 | ); |
| 295 | } |
| 296 | } |
| 297 | |
| 298 | #[tokio::test] |
| 299 | async fn chatgpt_stream_validates_returned_function_namespace() { |
| 300 | let events = collect_responses_stream(concat!( |
| 301 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"function_call\",\"namespace\":\"codewhale\",\"call_id\":\"call_1\",\"id\":\"fc_1\",\"name\":\"read\"}}\n\n", |
| 302 | "data: {\"type\":\"response.function_call_arguments.delta\",\"delta\":\"{}\"}\n\n", |
| 303 | "data: {\"type\":\"response.output_item.done\"}\n\n", |
| 304 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 305 | )).await; |
| 306 | assert!( |
| 307 | events.iter().any(|event| matches!( |
| 308 | event, |
| 309 | Ok(StreamEvent::ContentBlockStart { |
| 310 | content_block: ContentBlockStart::ToolUse { id, name, .. }, .. |
| 311 | }) if id == "call_1|fc_1" && name == "read" |
| 312 | )), |
| 313 | "{events:?}" |
| 314 | ); |
| 315 | assert!( |
| 316 | events.iter().any(|event| matches!( |
| 317 | event, Ok(StreamEvent::MessageDelta { delta, .. }) |
| 318 | if delta.stop_reason.as_deref() == Some("tool_use") |
| 319 | )), |
| 320 | "{events:?}" |
| 321 | ); |
| 322 | assert!(matches!(events.last(), Some(Ok(StreamEvent::MessageStop)))); |
| 323 | |
| 324 | for body in [ |
| 325 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"function_call\",\"namespace\":\"untrusted\",\"call_id\":\"call_1\",\"name\":\"read\"}}\n\n", |
| 326 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"function_call\",\"call_id\":\"call_1\",\"name\":\"read\"}}\n\n", |
| 327 | ] { |
| 328 | let events = collect_responses_stream(body).await; |
| 329 | assert!(events.last().is_some_and(Result::is_err), "{events:?}"); |
| 330 | assert!( |
| 331 | !events.iter().any(|event| matches!( |
| 332 | event, |
| 333 | Ok(StreamEvent::ContentBlockStart { |
| 334 | content_block: ContentBlockStart::ToolUse { .. }, |
| 335 | .. |
| 336 | }) |
| 337 | )), |
| 338 | "unrecognized namespaces must not invoke local tools: {events:?}" |
| 339 | ); |
| 340 | } |
| 341 | } |
| 342 | |
| 343 | #[tokio::test] |
| 344 | async fn chatgpt_plan_usage_errors_explain_where_to_check_allowance() { |
| 345 | for body in [ |
| 346 | "data: {\"type\":\"error\",\"code\":\"subscription_sharing_usage_limit_exceeded\",\"message\":\"opaque upstream message\"}\n\n", |
| 347 | "data: {\"type\":\"response.failed\",\"response\":{\"error\":{\"code\":\"subscription_sharing_usage_unavailable\",\"message\":\"opaque upstream message\"}}}\n\n", |
| 348 | ] { |
| 349 | let events = collect_responses_stream(body).await; |
| 350 | let error = events |
| 351 | .last() |
| 352 | .unwrap() |
| 353 | .as_ref() |
| 354 | .expect_err("plan usage error"); |
| 355 | let typed = error |
| 356 | .downcast_ref::<crate::llm_client::LlmError>() |
| 357 | .expect("typed allowance error"); |
| 358 | assert!(matches!( |
| 359 | typed, |
| 360 | crate::llm_client::LlmError::QuotaExhausted(_) |
| 361 | )); |
| 362 | assert!(!typed.is_retryable()); |
| 363 | assert!(error.to_string().contains("ChatGPT plan usage"), "{error}"); |
| 364 | assert!( |
| 365 | error.to_string().contains("ChatGPT Settings > Usage"), |
| 366 | "{error}" |
| 367 | ); |
| 368 | assert!( |
| 369 | !events |
| 370 | .iter() |
| 371 | .any(|event| matches!(event, Ok(StreamEvent::MessageStop))), |
| 372 | "{events:?}" |
| 373 | ); |
| 374 | } |
| 375 | } |
| 376 | |
| 377 | #[tokio::test] |
| 378 | async fn chatgpt_terminal_failures_preserve_usage_without_settling() { |
| 379 | for response in [ |
| 380 | json!({"status":"incomplete","incomplete_details":{"reason":"max_output_tokens"}}), |
| 381 | json!({"status":"failed","error":{"code":"subscription_sharing_usage_unavailable"}}), |
| 382 | ] { |
| 383 | let event_type = if response["status"] == "incomplete" { |
| 384 | "response.incomplete" |
| 385 | } else { |
| 386 | "response.failed" |
| 387 | }; |
| 388 | let mut response = response; |
| 389 | response["usage"] = json!({"input_tokens":13,"output_tokens":5}); |
| 390 | let events = collect_responses_stream(&format!( |
| 391 | "data: {}\n\n", |
| 392 | json!({"type":event_type,"response":response}) |
| 393 | )) |
| 394 | .await; |
| 395 | assert!(events.iter().any(|event| matches!(event, Ok(StreamEvent::MessageDelta { usage: Some(usage), .. }) if usage.input_tokens == 13 && usage.output_tokens == 5))); |
| 396 | let error = events.last().unwrap().as_ref().unwrap_err(); |
| 397 | let typed = error.downcast_ref::<crate::llm_client::LlmError>().unwrap(); |
| 398 | assert!(!typed.is_retryable()); |
| 399 | assert!( |
| 400 | !events |
| 401 | .iter() |
| 402 | .any(|event| matches!(event, Ok(StreamEvent::MessageStop))) |
| 403 | ); |
| 404 | } |
| 405 | } |
| 406 | |
| 407 | #[tokio::test] |
| 408 | async fn chatgpt_http_usage_limit_is_not_retried() { |
| 409 | let server = MockServer::start().await; |
| 410 | let attempts = Arc::new(AtomicUsize::new(0)); |
| 411 | Mock::given(method("POST")) |
| 412 | .and(path("/v1/responses")) |
| 413 | .respond_with(AlwaysError { |
| 414 | attempts: Arc::clone(&attempts), |
| 415 | status: 429, |
| 416 | body: r#"{"error":{"code":"subscription_sharing_usage_limit_exceeded","message":"opaque"}}"#, |
| 417 | }).mount(&server).await; |
| 418 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 419 | let request = client |
| 420 | .prepare_outbound_request(minimal_responses_request(), true) |
| 421 | .unwrap(); |
| 422 | let error = match client.handle_responses_stream(&request).await { |
| 423 | Ok(_) => panic!("allowance error must fail"), |
| 424 | Err(error) => error, |
| 425 | }; |
| 426 | assert_eq!(attempts.load(Ordering::SeqCst), 1); |
| 427 | let typed = error.downcast_ref::<crate::llm_client::LlmError>().unwrap(); |
| 428 | assert!(matches!( |
| 429 | typed, |
| 430 | crate::llm_client::LlmError::QuotaExhausted(_) |
| 431 | )); |
| 432 | assert!(!typed.is_retryable()); |
| 433 | } |
| 434 | |
| 435 | #[tokio::test] |
| 436 | async fn responses_stream_joins_multiline_data_fields_into_one_event() { |
| 437 | let events = collect_responses_stream(concat!( |
| 438 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"message\",\"id\":\"m\"}}\n\n", |
| 439 | "data: {\"type\":\"response.output_text.delta\",\n", |
| 440 | "data: \"delta\":\"joined\"}\n\n", |
| 441 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 442 | )) |
| 443 | .await; |
| 444 | assert!( |
| 445 | events.iter().any(|event| matches!( |
| 446 | event, |
| 447 | Ok(StreamEvent::ContentBlockDelta { |
| 448 | delta: Delta::TextDelta { text }, |
| 449 | .. |
| 450 | }) if text == "joined" |
| 451 | )), |
| 452 | "the split event was lost: {events:?}" |
| 453 | ); |
| 454 | assert!(matches!(events.last(), Some(Ok(StreamEvent::MessageStop)))); |
| 455 | } |
| 456 | |
| 457 | #[tokio::test] |
| 458 | async fn responses_stream_finishes_on_semantic_terminal_event_without_done_marker() { |
| 459 | let server = MockServer::start().await; |
| 460 | let sse_body = concat!( |
| 461 | "data: {\"type\":\"response.created\",\"response\":{\"status\":\"in_progress\"}}\n\n", |
| 462 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\",\"usage\":{\"input_tokens\":3,\"output_tokens\":2}}}\n\n", |
| 463 | ); |
| 464 | Mock::given(method("POST")) |
| 465 | .and(path("/v1/responses")) |
| 466 | .respond_with( |
| 467 | ResponseTemplate::new(200) |
| 468 | .insert_header("Content-Type", "text/event-stream") |
| 469 | .set_body_string(sse_body), |
| 470 | ) |
| 471 | .mount(&server) |
| 472 | .await; |
| 473 | |
| 474 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 475 | let mut stream = client |
| 476 | .handle_responses_stream( |
| 477 | &client |
| 478 | .prepare_outbound_request(minimal_responses_request(), true) |
| 479 | .expect("responses request prepares"), |
| 480 | ) |
| 481 | .await |
| 482 | .expect("semantic Responses stream opens"); |
| 483 | |
| 484 | let mut saw_stop = false; |
| 485 | let mut usage = None; |
| 486 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 487 | while let Some(event) = stream.next().await { |
| 488 | match event.unwrap() { |
| 489 | StreamEvent::MessageStop => saw_stop = true, |
| 490 | StreamEvent::MessageDelta { usage: value, .. } => usage = value, |
| 491 | _ => {} |
| 492 | } |
| 493 | } |
| 494 | }) |
| 495 | .await |
| 496 | .expect("terminal event ends the stream without [DONE]"); |
| 497 | assert!(saw_stop); |
| 498 | let usage = usage.expect("completed response carries usage"); |
| 499 | assert_eq!(usage.input_tokens, 3); |
| 500 | assert_eq!(usage.output_tokens, 2); |
| 501 | } |
| 502 | |
| 503 | #[tokio::test] |
| 504 | async fn responses_stream_surfaces_notice_for_web_search_call_items() { |
| 505 | let server = MockServer::start().await; |
| 506 | let sse_body = concat!( |
| 507 | "data: {\"type\":\"response.created\",\"response\":{\"status\":\"in_progress\"}}\n\n", |
| 508 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"web_search_call\",\"id\":\"ws_1\",\"call_id\":\"call_1\"}}\n\n", |
| 509 | "data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"web_search_call\",\"id\":\"ws_1\"}}\n\n", |
| 510 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\",\"usage\":{\"input_tokens\":3,\"output_tokens\":2}}}\n\n", |
| 511 | ); |
| 512 | Mock::given(method("POST")) |
| 513 | .and(path("/v1/responses")) |
| 514 | .respond_with( |
| 515 | ResponseTemplate::new(200) |
| 516 | .insert_header("Content-Type", "text/event-stream") |
| 517 | .set_body_string(sse_body), |
| 518 | ) |
| 519 | .mount(&server) |
| 520 | .await; |
| 521 | |
| 522 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 523 | let mut stream = client |
| 524 | .handle_responses_stream( |
| 525 | &client |
| 526 | .prepare_outbound_request(minimal_responses_request(), true) |
| 527 | .expect("responses request prepares"), |
| 528 | ) |
| 529 | .await |
| 530 | .expect("semantic Responses stream opens"); |
| 531 | |
| 532 | let mut saw_notice = false; |
| 533 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 534 | while let Some(event) = stream.next().await { |
| 535 | if let Ok(StreamEvent::ContentBlockStart { |
| 536 | content_block: ContentBlockStart::Text { text }, |
| 537 | .. |
| 538 | }) = event |
| 539 | && text.contains("not replayed") |
| 540 | { |
| 541 | saw_notice = true; |
| 542 | } |
| 543 | } |
| 544 | }) |
| 545 | .await |
| 546 | .expect("stream terminates"); |
| 547 | assert!(saw_notice, "web_search_call must surface a visible notice"); |
| 548 | } |
| 549 | |
| 550 | #[tokio::test] |
| 551 | async fn responses_stream_fails_fast_on_non_retryable_provider_error() { |
| 552 | let server = MockServer::start().await; |
| 553 | let attempts = Arc::new(AtomicUsize::new(0)); |
| 554 | Mock::given(method("POST")) |
| 555 | .and(path("/v1/responses")) |
| 556 | .respond_with(AlwaysError { |
| 557 | attempts: Arc::clone(&attempts), |
| 558 | status: 403, |
| 559 | body: "<html><title>Access Denied</title><body>Security alert. Contact support. Ray ID 1234abcd.</body></html>", |
| 560 | }) |
| 561 | .mount(&server) |
| 562 | .await; |
| 563 | |
| 564 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 565 | |
| 566 | let err = match client |
| 567 | .handle_responses_stream( |
| 568 | &client |
| 569 | .prepare_outbound_request(minimal_responses_request(), true) |
| 570 | .expect("responses request prepares"), |
| 571 | ) |
| 572 | .await |
| 573 | { |
| 574 | Ok(_) => panic!("non-retryable Responses errors should fail fast"), |
| 575 | Err(err) => err, |
| 576 | }; |
| 577 | |
| 578 | assert_eq!(attempts.load(Ordering::SeqCst), 1); |
| 579 | let message = format!("{err:#}"); |
| 580 | assert!( |
| 581 | message.contains("Responses API request failed"), |
| 582 | "{message}" |
| 583 | ); |
| 584 | assert!(message.contains("OpenAI Codex"), "{message}"); |
| 585 | assert!(message.contains("Access Denied"), "{message}"); |
| 586 | assert!( |
| 587 | message.contains("blocked before it reached the model"), |
| 588 | "{message}" |
| 589 | ); |
| 590 | // #3884: the structured LlmError must stay downcastable through the |
| 591 | // context layers so sub-agent failure records can classify it. |
| 592 | assert!( |
| 593 | err.downcast_ref::<crate::llm_client::LlmError>().is_some(), |
| 594 | "LlmError should survive the anyhow chain" |
| 595 | ); |
| 596 | } |
| 597 | |
| 598 | #[test] |
| 599 | fn responses_body_serializes_the_child_catalog_without_duplication() { |
| 600 | // Mirror of the Anthropic contract: the real child catalog fixture |
| 601 | // maps 1:1 into Responses function tools with one canonical `read` entry. |
| 602 | // `load_skill` is eager in DEFAULT_ACTIVE_NATIVE_TOOLS and children resolve |
| 603 | // the same catalog authority, so it maps through exactly once too. |
| 604 | let tools = crate::tools::subagent::kimi_general_child_request_tools_fixture(); |
| 605 | let mut request = minimal_responses_request(); |
| 606 | request.tools = Some(tools); |
| 607 | let body = build_responses_body(&request); |
| 608 | assert_eq!(body["parallel_tool_calls"], false); |
| 609 | let generic = build_responses_body_for_provider(&request, ProviderKind::Openai, None); |
| 610 | assert_eq!(generic["parallel_tool_calls"], true); |
| 611 | assert_eq!(body["tools"].as_array().unwrap().len(), 1); |
| 612 | assert_eq!(body["tools"][0]["type"], "namespace"); |
| 613 | assert_eq!(body["tools"][0]["name"], "codewhale"); |
| 614 | let serialized = body["tools"][0]["tools"] |
| 615 | .as_array() |
| 616 | .expect("functions serialize inside the Codewhale namespace"); |
| 617 | let reads: Vec<_> = serialized |
| 618 | .iter() |
| 619 | .filter(|tool| tool["name"] == "read") |
| 620 | .collect(); |
| 621 | assert_eq!( |
| 622 | reads.len(), |
| 623 | 1, |
| 624 | "exactly one canonical read definition reaches the Responses wire" |
| 625 | ); |
| 626 | assert!( |
| 627 | reads[0]["parameters"]["properties"].is_object(), |
| 628 | "read keeps a valid parameters schema: {}", |
| 629 | reads[0] |
| 630 | ); |
| 631 | assert_eq!( |
| 632 | serialized |
| 633 | .iter() |
| 634 | .filter(|tool| tool["name"] == "load_skill") |
| 635 | .count(), |
| 636 | 1, |
| 637 | "exactly one canonical load_skill definition reaches the Responses wire" |
| 638 | ); |
| 639 | } |
| 640 | |
| 641 | #[tokio::test] |
| 642 | async fn responses_stream_open_preserves_wire_headers_through_shared_seam() { |
| 643 | use wiremock::matchers::header; |
| 644 | |
| 645 | let server = MockServer::start().await; |
| 646 | // Public API requests retain bearer/SSE headers and Codewhale identity |
| 647 | // through the shared stream-entry transport; no backend headers survive. |
| 648 | Mock::given(method("POST")) |
| 649 | .and(path("/v1/responses")) |
| 650 | .and(header("Accept", "text/event-stream")) |
| 651 | .and(header("Authorization", "Bearer test-token")) |
| 652 | .respond_with( |
| 653 | ResponseTemplate::new(200) |
| 654 | .insert_header("Content-Type", "text/event-stream") |
| 655 | .set_body_string("data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n"), |
| 656 | ) |
| 657 | .expect(1) |
| 658 | .mount(&server) |
| 659 | .await; |
| 660 | |
| 661 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 662 | let mut stream = client |
| 663 | .handle_responses_stream( |
| 664 | &client |
| 665 | .prepare_outbound_request(minimal_responses_request(), true) |
| 666 | .expect("responses request prepares"), |
| 667 | ) |
| 668 | .await |
| 669 | .expect("stream opens with preserved headers"); |
| 670 | |
| 671 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 672 | while let Some(event) = stream.next().await { |
| 673 | event.unwrap(); |
| 674 | } |
| 675 | }) |
| 676 | .await |
| 677 | .expect("stream should finish after response.completed"); |
| 678 | |
| 679 | let requests = server |
| 680 | .received_requests() |
| 681 | .await |
| 682 | .expect("recorded public API request"); |
| 683 | assert_eq!(requests.len(), 1); |
| 684 | let user_agent = requests[0] |
| 685 | .headers |
| 686 | .get("user-agent") |
| 687 | .unwrap() |
| 688 | .to_str() |
| 689 | .unwrap(); |
| 690 | assert!(user_agent.contains("codewhale/"), "{user_agent}"); |
| 691 | assert!(!user_agent.contains("codex_cli_rs"), "{user_agent}"); |
| 692 | for header in ["openai-beta", "originator", "chatgpt-account-id"] { |
| 693 | assert!( |
| 694 | requests[0].headers.get(header).is_none(), |
| 695 | "{header} must not reach the public API" |
| 696 | ); |
| 697 | } |
| 698 | } |
| 699 | |
| 700 | #[tokio::test] |
| 701 | async fn responses_stream_inserts_boundary_between_reasoning_summary_parts() { |
| 702 | let server = MockServer::start().await; |
| 703 | let sse_body = concat!( |
| 704 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"reasoning\",\"id\":\"rs_1\"}}\n\n", |
| 705 | "data: {\"type\":\"response.reasoning_summary_part.added\",\"item_id\":\"rs_1\",\"summary_index\":0,\"part\":{\"type\":\"summary_text\",\"text\":\"\"}}\n\n", |
| 706 | "data: {\"type\":\"response.reasoning_summary_text.delta\",\"delta\":\"partA\"}\n\n", |
| 707 | "data: {\"type\":\"response.reasoning_summary_part.added\",\"item_id\":\"rs_1\",\"summary_index\":1,\"part\":{\"type\":\"summary_text\",\"text\":\"\"}}\n\n", |
| 708 | "data: {\"type\":\"response.reasoning_summary_text.delta\",\"delta\":\"partB\"}\n\n", |
| 709 | "data: {\"type\":\"response.output_item.done\"}\n\n", |
| 710 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 711 | ); |
| 712 | Mock::given(method("POST")) |
| 713 | .and(path("/v1/responses")) |
| 714 | .respond_with( |
| 715 | ResponseTemplate::new(200) |
| 716 | .insert_header("Content-Type", "text/event-stream") |
| 717 | .set_body_string(sse_body), |
| 718 | ) |
| 719 | .mount(&server) |
| 720 | .await; |
| 721 | |
| 722 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 723 | let mut stream = client |
| 724 | .handle_responses_stream( |
| 725 | &client |
| 726 | .prepare_outbound_request(minimal_responses_request(), true) |
| 727 | .expect("responses request prepares"), |
| 728 | ) |
| 729 | .await |
| 730 | .unwrap(); |
| 731 | |
| 732 | let mut thinking = String::new(); |
| 733 | tokio::time::timeout(std::time::Duration::from_secs(5), async { |
| 734 | while let Some(event) = stream.next().await { |
| 735 | if let StreamEvent::ContentBlockDelta { |
| 736 | delta: Delta::ThinkingDelta { thinking: chunk }, |
| 737 | .. |
| 738 | } = event.unwrap() |
| 739 | { |
| 740 | thinking.push_str(&chunk); |
| 741 | } |
| 742 | } |
| 743 | }) |
| 744 | .await |
| 745 | .expect("Responses reasoning stream should finish after response.completed"); |
| 746 | |
| 747 | // The second summary part must be separated from the first by a |
| 748 | // paragraph break, and no separator may precede the first part. |
| 749 | assert_eq!(thinking, "partA\n\npartB"); |
| 750 | } |
| 751 | |
| 752 | #[test] |
| 753 | fn codex_reasoning_effort_uses_responses_labels() { |
| 754 | assert_eq!(codex_responses_reasoning_effort("max"), Some("max")); |
| 755 | assert_eq!(codex_responses_reasoning_effort("maximum"), Some("max")); |
| 756 | assert_eq!(codex_responses_reasoning_effort("xhigh"), Some("xhigh")); |
| 757 | assert_eq!(codex_responses_reasoning_effort("ultra"), Some("ultra")); |
| 758 | assert_eq!(codex_responses_reasoning_effort("ultracode"), Some("ultra")); |
| 759 | assert_eq!(codex_responses_reasoning_effort("high"), Some("high")); |
| 760 | assert_eq!(codex_responses_reasoning_effort("medium"), Some("medium")); |
| 761 | assert_eq!(codex_responses_reasoning_effort("minimal"), Some("low")); |
| 762 | assert_eq!(codex_responses_reasoning_effort("auto"), Some("medium")); |
| 763 | assert_eq!(codex_responses_reasoning_effort("off"), Some("low")); |
| 764 | } |
| 765 | |
| 766 | #[tokio::test] |
| 767 | async fn codex_selected_effort_reaches_preview_wire_and_restored_receipt_unchanged() { |
| 768 | use crate::reasoning_preference::{EffectiveReasoningEffort, ReasoningEffort}; |
| 769 | use crate::work_graph::WorkActivityEvent; |
| 770 | |
| 771 | let server = MockServer::start().await; |
| 772 | Mock::given(method("POST")) |
| 773 | .and(path("/v1/responses")) |
| 774 | .respond_with( |
| 775 | ResponseTemplate::new(200) |
| 776 | .insert_header("Content-Type", "text/event-stream") |
| 777 | .set_body_string("data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n"), |
| 778 | ) |
| 779 | .expect(6) |
| 780 | .mount(&server) |
| 781 | .await; |
| 782 | let client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 783 | let receipts = tempfile::tempdir().unwrap(); |
| 784 | for effort in ["low", "medium", "high", "xhigh", "max", "ultra"] { |
| 785 | let selected = ReasoningEffort::parse_strict(effort).unwrap(); |
| 786 | let activity = WorkActivityEvent::ReasoningEffortChanged { |
| 787 | requested: selected.into(), |
| 788 | effective: selected.into(), |
| 789 | provider_kind: Some(ProviderKind::OpenaiCodex), |
| 790 | provider: "openai-codex".to_string(), |
| 791 | endpoint_identity: Some(crate::config::DEFAULT_OPENAI_CODEX_BASE_URL.to_string()), |
| 792 | model: Some("gpt-6-astra".to_string()), |
| 793 | ts: 1, |
| 794 | operation: None, |
| 795 | }; |
| 796 | let persisted = serde_json::to_value(activity).unwrap(); |
| 797 | assert_eq!(persisted["requested"], effort); |
| 798 | assert_eq!(persisted["effective"], effort); |
| 799 | let receipt_path = receipts.path().join(format!("{effort}.json")); |
| 800 | std::fs::write(&receipt_path, serde_json::to_vec(&persisted).unwrap()).unwrap(); |
| 801 | let WorkActivityEvent::ReasoningEffortChanged { effective, .. } = |
| 802 | serde_json::from_slice(&std::fs::read(receipt_path).unwrap()).unwrap(); |
| 803 | let restored = EffectiveReasoningEffort::from(effective) |
| 804 | .request_tier_for_replay() |
| 805 | .unwrap(); |
| 806 | assert_eq!(restored, selected); |
| 807 | let mut request = minimal_responses_request(); |
| 808 | request.model = "gpt-6-astra".to_string(); |
| 809 | request.reasoning_effort = restored |
| 810 | .api_value_for_provider(ProviderKind::OpenaiCodex) |
| 811 | .map(str::to_string); |
| 812 | let prepared = client.prepare_outbound_request(request, true).unwrap(); |
| 813 | assert_eq!( |
| 814 | prepared.reasoning.wire_effort(), |
| 815 | Some(("reasoning.effort", effort)) |
| 816 | ); |
| 817 | assert_eq!(prepared.body["reasoning"]["effort"], effort); |
| 818 | let mut stream = client.handle_responses_stream(&prepared).await.unwrap(); |
| 819 | while let Some(event) = stream.next().await { |
| 820 | event.unwrap(); |
| 821 | } |
| 822 | } |
| 823 | let requests = server.received_requests().await.unwrap(); |
| 824 | assert_eq!(requests.len(), 6); |
| 825 | for (request, effort) in requests |
| 826 | .iter() |
| 827 | .zip(["low", "medium", "high", "xhigh", "max", "ultra"]) |
| 828 | { |
| 829 | let body: Value = serde_json::from_slice(&request.body).unwrap(); |
| 830 | assert_eq!(body["model"], "gpt-6-astra"); |
| 831 | assert_eq!(body["reasoning"]["effort"], effort); |
| 832 | } |
| 833 | } |
| 834 | |
| 835 | #[test] |
| 836 | fn codex_tiers_do_not_change_other_responses_provider_dialects() { |
| 837 | let mut request = minimal_responses_request(); |
| 838 | for effort in ["max", "ultra"] { |
| 839 | request.reasoning_effort = Some(effort.to_string()); |
| 840 | assert_eq!( |
| 841 | build_responses_body_for_provider(&request, ProviderKind::Concentrate, None)["reasoning"] |
| 842 | ["effort"], |
| 843 | "xhigh" |
| 844 | ); |
| 845 | assert_eq!( |
| 846 | build_responses_body_for_provider(&request, ProviderKind::Deepseek, None)["reasoning"] |
| 847 | ["effort"], |
| 848 | "max" |
| 849 | ); |
| 850 | } |
| 851 | } |
| 852 | |
| 853 | /// Concentrate's parameter reference documents `model`, `input`, `stream`, |
| 854 | /// `max_output_tokens`, `tools`/`tool_choice`/`parallel_tool_calls`, and |
| 855 | /// `reasoning.effort`; `store`, `include`, `instructions`, and |
| 856 | /// `reasoning.summary` are absent. The body sends only documented fields and |
| 857 | /// carries the system prompt as a leading `system` input item. |
| 858 | /// https://concentrate.ai/docs/api-reference/endpoint/request-parameters |
| 859 | #[test] |
| 860 | fn concentrate_responses_body_sends_only_documented_fields() { |
| 861 | let mut request = minimal_responses_request(); |
| 862 | request.model = "openai/gpt-5.6-sol".to_string(); |
| 863 | request.system = Some(SystemPrompt::Text( |
| 864 | "You are the Codewhale test system prompt.".to_string(), |
| 865 | )); |
| 866 | request.reasoning_effort = Some("high".to_string()); |
| 867 | request.tools = Some(vec![Tool { |
| 868 | tool_type: None, |
| 869 | name: "read".to_string(), |
| 870 | description: "Read a file".to_string(), |
| 871 | input_schema: serde_json::json!({ |
| 872 | "type": "object", |
| 873 | "properties": { "path": { "type": "string" } }, |
| 874 | "required": ["path"] |
| 875 | }), |
| 876 | allowed_callers: None, |
| 877 | defer_loading: None, |
| 878 | input_examples: None, |
| 879 | strict: None, |
| 880 | cache_control: None, |
| 881 | }]); |
| 882 | |
| 883 | let body = build_responses_body_for_provider(&request, ProviderKind::Concentrate, None); |
| 884 | let documented = [ |
| 885 | "model", |
| 886 | "input", |
| 887 | "max_output_tokens", |
| 888 | "temperature", |
| 889 | "top_p", |
| 890 | "stream", |
| 891 | "text", |
| 892 | "reasoning", |
| 893 | "tools", |
| 894 | "tool_choice", |
| 895 | "parallel_tool_calls", |
| 896 | "routing", |
| 897 | "cache_control", |
| 898 | "prompt_cache_options", |
| 899 | ]; |
| 900 | for key in body.as_object().expect("object body").keys() { |
| 901 | assert!( |
| 902 | documented.contains(&key.as_str()), |
| 903 | "undocumented top-level field `{key}` reached the Concentrate wire: {body}" |
| 904 | ); |
| 905 | } |
| 906 | assert_eq!( |
| 907 | body["model"], "openai/gpt-5.6-sol", |
| 908 | "provider/model ids pass through verbatim" |
| 909 | ); |
| 910 | assert_eq!(body["stream"], true); |
| 911 | assert!(body.get("store").is_none(), "{body}"); |
| 912 | assert!(body.get("include").is_none(), "{body}"); |
| 913 | assert!(body.get("instructions").is_none(), "{body}"); |
| 914 | let input = body["input"].as_array().expect("input array"); |
| 915 | assert_eq!(input[0]["type"], "message"); |
| 916 | assert_eq!(input[0]["role"], "system"); |
| 917 | assert_eq!(input[0]["content"][0]["type"], "input_text"); |
| 918 | assert_eq!( |
| 919 | input[0]["content"][0]["text"], |
| 920 | "You are the Codewhale test system prompt." |
| 921 | ); |
| 922 | assert_eq!(input[1]["role"], "user"); |
| 923 | assert_eq!(body["reasoning"], serde_json::json!({ "effort": "high" })); |
| 924 | assert_eq!(body["tools"][0]["type"], "function"); |
| 925 | assert_eq!(body["tools"][0]["name"], "read"); |
| 926 | assert_eq!(body["tools"][0]["strict"], false); |
| 927 | assert_eq!(body["tool_choice"], "auto"); |
| 928 | assert_eq!(body["parallel_tool_calls"], true); |
| 929 | |
| 930 | // The same request on the generic Responses path still carries the |
| 931 | // OpenAI-only fields, so the Concentrate branch is a deliberate subset. |
| 932 | let generic = build_responses_body_for_provider(&request, ProviderKind::Openai, None); |
| 933 | assert!( |
| 934 | generic.get("store").is_some() |
| 935 | && generic.get("include").is_some() |
| 936 | && generic.get("instructions").is_some() |
| 937 | ); |
| 938 | } |
| 939 | |
| 940 | #[test] |
| 941 | fn deepseek_flash_responses_body_uses_stateless_0731_contract() { |
| 942 | let mut request = minimal_responses_request(); |
| 943 | request.model = "deepseek-v4-flash".to_string(); |
| 944 | request.reasoning_effort = Some("xhigh".to_string()); |
| 945 | request.temperature = Some(1.0); |
| 946 | request.top_p = Some(0.95); |
| 947 | request.messages.insert( |
| 948 | 0, |
| 949 | Message { |
| 950 | role: Role::Assistant, |
| 951 | content: vec![ContentBlock::Thinking { |
| 952 | thinking: "preserve this tool-loop reasoning".to_string(), |
| 953 | signature: None, |
| 954 | state: None, |
| 955 | }], |
| 956 | }, |
| 957 | ); |
| 958 | |
| 959 | let body = build_responses_body_for_provider(&request, ProviderKind::Deepseek, None); |
| 960 | |
| 961 | assert_eq!(body["model"], "deepseek-v4-flash"); |
| 962 | assert_eq!(body["max_output_tokens"], 128); |
| 963 | assert_eq!(body["temperature"], 1.0); |
| 964 | assert!( |
| 965 | (body["top_p"].as_f64().expect("top_p number") - 0.95).abs() < 1e-6, |
| 966 | "{}", |
| 967 | body["top_p"] |
| 968 | ); |
| 969 | assert_eq!(body.pointer("/reasoning/effort"), Some(&json!("high"))); |
| 970 | assert!(body.pointer("/reasoning/summary").is_none()); |
| 971 | assert!(body.get("include").is_none()); |
| 972 | assert!(body.get("store").is_none()); |
| 973 | assert_eq!( |
| 974 | body.pointer("/input/0/content/0/type"), |
| 975 | Some(&json!("reasoning_text")) |
| 976 | ); |
| 977 | assert_eq!( |
| 978 | body.pointer("/input/0/content/0/text"), |
| 979 | Some(&json!("preserve this tool-loop reasoning")) |
| 980 | ); |
| 981 | } |
| 982 | |
| 983 | #[test] |
| 984 | fn chatgpt_plan_body_omits_unsupported_output_caps() { |
| 985 | // The official ChatGPT plan preview does not support output-cap fields. |
| 986 | // Other Responses providers keep the central cap on the wire. |
| 987 | let mut request = minimal_responses_request(); |
| 988 | request.max_tokens = 4_096; |
| 989 | |
| 990 | let codex = build_responses_body_for_provider(&request, ProviderKind::OpenaiCodex, None); |
| 991 | assert!( |
| 992 | codex.get("max_output_tokens").is_none(), |
| 993 | "ChatGPT plan body names an unsupported output cap: {codex}" |
| 994 | ); |
| 995 | assert!( |
| 996 | codex.get("max_tokens").is_none() && codex.get("max_completion_tokens").is_none(), |
| 997 | "no alternate output-cap spelling may sneak onto the Codex wire: {codex}" |
| 998 | ); |
| 999 | |
| 1000 | let deepseek = build_responses_body_for_provider(&request, ProviderKind::Deepseek, None); |
| 1001 | assert_eq!(deepseek["max_output_tokens"], json!(4_096)); |
| 1002 | } |
| 1003 | |
| 1004 | #[test] |
| 1005 | fn chatgpt_replays_only_exact_grant_and_model_opaque_reasoning_state() { |
| 1006 | const SENTINEL: &str = "readable private reasoning must not be replayed"; |
| 1007 | const SCOPE: &str = "openai-responses-siwc-v1:test-grant"; |
| 1008 | let state = OpaqueReasoningState { |
| 1009 | provider: ProviderKind::OpenaiCodex.as_str().to_string(), |
| 1010 | api: SCOPE.to_string(), |
| 1011 | model: "gpt-5.5".to_string(), |
| 1012 | id: Some("rs_opaque".to_string()), |
| 1013 | encrypted_content: "enc_opaque_payload".to_string(), |
| 1014 | }; |
| 1015 | let mut request = minimal_responses_request(); |
| 1016 | request.messages.insert( |
| 1017 | 0, |
| 1018 | Message { |
| 1019 | role: Role::Assistant, |
| 1020 | content: vec![ContentBlock::Thinking { |
| 1021 | thinking: SENTINEL.to_string(), |
| 1022 | signature: None, |
| 1023 | state: Some(state), |
| 1024 | }], |
| 1025 | }, |
| 1026 | ); |
| 1027 | |
| 1028 | let exact = build_responses_body_for_provider(&request, ProviderKind::OpenaiCodex, Some(SCOPE)); |
| 1029 | let exact_wire = exact.to_string(); |
| 1030 | assert!(!exact_wire.contains(SENTINEL), "{exact}"); |
| 1031 | assert_eq!(exact.pointer("/input/0/type"), Some(&json!("reasoning"))); |
| 1032 | assert_eq!(exact.pointer("/input/0/id"), Some(&json!("rs_opaque"))); |
| 1033 | assert_eq!(exact.pointer("/input/0/summary"), Some(&json!([]))); |
| 1034 | assert_eq!( |
| 1035 | exact.pointer("/input/0/encrypted_content"), |
| 1036 | Some(&json!("enc_opaque_payload")) |
| 1037 | ); |
| 1038 | |
| 1039 | for other_scope in [None, Some("openai-responses-siwc-v1:another-grant")] { |
| 1040 | let body = |
| 1041 | build_responses_body_for_provider(&request, ProviderKind::OpenaiCodex, other_scope); |
| 1042 | assert!(!body.to_string().contains("enc_opaque_payload"), "{body}"); |
| 1043 | assert!(!body.to_string().contains(SENTINEL), "{body}"); |
| 1044 | } |
| 1045 | let mut legacy = request.clone(); |
| 1046 | if let ContentBlock::Thinking { |
| 1047 | state: Some(state), .. |
| 1048 | } = &mut legacy.messages[0].content[0] |
| 1049 | { |
| 1050 | state.api = "openai-responses".to_string(); |
| 1051 | } |
| 1052 | let legacy_body = |
| 1053 | build_responses_body_for_provider(&legacy, ProviderKind::OpenaiCodex, Some(SCOPE)); |
| 1054 | assert!(!legacy_body.to_string().contains("enc_opaque_payload")); |
| 1055 | |
| 1056 | request.model = "gpt-5.6".to_string(); |
| 1057 | let switched_model = |
| 1058 | build_responses_body_for_provider(&request, ProviderKind::OpenaiCodex, Some(SCOPE)); |
| 1059 | assert!(!switched_model.to_string().contains(SENTINEL)); |
| 1060 | assert!( |
| 1061 | switched_model |
| 1062 | .get("input") |
| 1063 | .and_then(Value::as_array) |
| 1064 | .is_some_and(|items| items.iter().all(|item| item["type"] != "reasoning")), |
| 1065 | "{switched_model}" |
| 1066 | ); |
| 1067 | |
| 1068 | let switched_provider = |
| 1069 | build_responses_body_for_provider(&request, ProviderKind::Deepseek, None); |
| 1070 | let switched_wire = switched_provider.to_string(); |
| 1071 | assert!(!switched_wire.contains(SENTINEL), "{switched_provider}"); |
| 1072 | assert!( |
| 1073 | !switched_wire.contains("enc_opaque_payload"), |
| 1074 | "{switched_provider}" |
| 1075 | ); |
| 1076 | } |
| 1077 | |
| 1078 | #[tokio::test] |
| 1079 | async fn chatgpt_stream_captures_only_scoped_encrypted_reasoning() { |
| 1080 | let server = MockServer::start().await; |
| 1081 | let sse_body = concat!( |
| 1082 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"reasoning\",\"id\":\"rs_1\"}}\n\n", |
| 1083 | "data: {\"type\":\"response.reasoning_summary_text.delta\",\"delta\":\"visible summary\"}\n\n", |
| 1084 | "data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"reasoning\",\"id\":\"rs_1\",\"summary\":[],\"encrypted_content\":\"enc_state\"}}\n\n", |
| 1085 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 1086 | ); |
| 1087 | Mock::given(method("POST")) |
| 1088 | .and(path("/v1/responses")) |
| 1089 | .respond_with( |
| 1090 | ResponseTemplate::new(200) |
| 1091 | .insert_header("Content-Type", "text/event-stream") |
| 1092 | .set_body_string(sse_body), |
| 1093 | ) |
| 1094 | .mount(&server) |
| 1095 | .await; |
| 1096 | |
| 1097 | for scope in [Some("openai-responses-siwc-v1:verified-test-grant"), None] { |
| 1098 | let mut client = CodewhaleClient::new(&test_codex_config(&server)).unwrap(); |
| 1099 | // The local HTTP fixture stands in for the official transport; production |
| 1100 | // obtains this frozen marker only from its selected verified grant. |
| 1101 | client.chatgpt_reasoning_api = scope.map(str::to_string); |
| 1102 | let mut stream = client |
| 1103 | .handle_responses_stream( |
| 1104 | &client |
| 1105 | .prepare_outbound_request(minimal_responses_request(), true) |
| 1106 | .expect("responses request prepares"), |
| 1107 | ) |
| 1108 | .await |
| 1109 | .unwrap(); |
| 1110 | let mut captured = None; |
| 1111 | while let Some(event) = stream.next().await { |
| 1112 | if let StreamEvent::ContentBlockDelta { |
| 1113 | delta: Delta::ReasoningStateDelta { state }, |
| 1114 | .. |
| 1115 | } = event.unwrap() |
| 1116 | { |
| 1117 | captured = Some(state); |
| 1118 | } |
| 1119 | } |
| 1120 | |
| 1121 | let Some(scope) = scope else { |
| 1122 | assert!( |
| 1123 | captured.is_none(), |
| 1124 | "a custom API-key route cannot mint grant state" |
| 1125 | ); |
| 1126 | continue; |
| 1127 | }; |
| 1128 | let state = captured.expect("encrypted reasoning state delta"); |
| 1129 | assert_eq!(state.provider, ProviderKind::OpenaiCodex.as_str()); |
| 1130 | assert_eq!(state.api, scope); |
| 1131 | assert_eq!(state.model, "gpt-5.5"); |
| 1132 | assert_eq!(state.id.as_deref(), Some("rs_1")); |
| 1133 | assert_eq!(state.encrypted_content, "enc_state"); |
| 1134 | } |
| 1135 | } |
| 1136 | |
| 1137 | #[test] |
| 1138 | fn deepseek_responses_reasoning_effort_uses_documented_labels() { |
| 1139 | assert_eq!(responses_reasoning_effort("low", true), Some("low")); |
| 1140 | assert_eq!(responses_reasoning_effort("medium", true), Some("high")); |
| 1141 | assert_eq!(responses_reasoning_effort("high", true), Some("high")); |
| 1142 | assert_eq!(responses_reasoning_effort("xhigh", true), Some("high")); |
| 1143 | assert_eq!(responses_reasoning_effort("max", true), Some("max")); |
| 1144 | // The off tier must disable thinking on the wire, not collapse into |
| 1145 | // low: DeepSeek documents `reasoning.effort: "none"` as the off value. |
| 1146 | assert_eq!(responses_reasoning_effort("off", true), Some("none")); |
| 1147 | assert_eq!(responses_reasoning_effort("disabled", true), Some("none")); |
| 1148 | assert_eq!(responses_reasoning_effort("none", true), Some("none")); |
| 1149 | assert_eq!(responses_reasoning_effort("false", true), Some("none")); |
| 1150 | // minimal stays a low tier for DeepSeek (undocumented label preserved |
| 1151 | // for Codex compatibility). |
| 1152 | assert_eq!(responses_reasoning_effort("minimal", true), Some("low")); |
| 1153 | } |
| 1154 | |
| 1155 | #[tokio::test] |
| 1156 | async fn generic_responses_captures_and_replays_opaque_reasoning() { |
| 1157 | let server = MockServer::start().await; |
| 1158 | Mock::given(method("POST")) |
| 1159 | .and(path("/v1/responses")) |
| 1160 | .respond_with(ResponseTemplate::new(200) |
| 1161 | .insert_header("Content-Type", "text/event-stream") |
| 1162 | .set_body_string(concat!( |
| 1163 | "data: {\"type\":\"response.output_item.added\",\"item\":{\"type\":\"reasoning\",\"id\":\"rs_generic\"}}\n\n", |
| 1164 | "data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"reasoning\",\"id\":\"rs_generic\",\"encrypted_content\":\"enc_generic\"}}\n\n", |
| 1165 | "data: {\"type\":\"response.completed\",\"response\":{\"status\":\"completed\"}}\n\n", |
| 1166 | ))) |
| 1167 | .mount(&server).await; |
| 1168 | let config = Config { |
| 1169 | provider: Some("openai".into()), |
| 1170 | providers: Some(ProvidersConfig { |
| 1171 | openai: ProviderConfig { |
| 1172 | api_key: Some("test-token".into()), |
| 1173 | base_url: Some(format!("{}/v1", server.uri())), |
| 1174 | ..Default::default() |
| 1175 | }, |
| 1176 | ..Default::default() |
| 1177 | }), |
| 1178 | ..Default::default() |
| 1179 | }; |
| 1180 | let client = CodewhaleClient::from_parts( |
| 1181 | format!("{}/v1", server.uri()), |
| 1182 | "gpt-5.5".into(), |
| 1183 | codewhale_config::provider::WireFormat::Responses, |
| 1184 | None, |
| 1185 | &config, |
| 1186 | ) |
| 1187 | .unwrap(); |
| 1188 | assert!(client.chatgpt_reasoning_api.is_none()); |
| 1189 | let prepared = client |
| 1190 | .prepare_outbound_request(minimal_responses_request(), true) |
| 1191 | .unwrap(); |
| 1192 | assert_eq!( |
| 1193 | prepared.endpoint.url, |
| 1194 | format!("{}/v1/responses", server.uri()) |
| 1195 | ); |
| 1196 | let mut stream = client.handle_responses_stream(&prepared).await.unwrap(); |
| 1197 | let mut captured = None; |
| 1198 | while let Some(event) = stream.next().await { |
| 1199 | if let StreamEvent::ContentBlockDelta { |
| 1200 | delta: Delta::ReasoningStateDelta { state }, |
| 1201 | .. |
| 1202 | } = event.unwrap() |
| 1203 | { |
| 1204 | captured = Some(state); |
| 1205 | } |
| 1206 | } |
| 1207 | let state = captured.expect("generic Responses must retain encrypted state"); |
| 1208 | assert_eq!(state.provider, "openai"); |
| 1209 | assert_eq!(state.api, "openai-responses"); |
| 1210 | assert_eq!(state.model, "gpt-5.5"); |
| 1211 | let mut continuation = minimal_responses_request(); |
| 1212 | continuation.messages.push(Message { |
| 1213 | role: Role::Assistant, |
| 1214 | content: vec![ContentBlock::Thinking { |
| 1215 | thinking: "readable-private-summary".into(), |
| 1216 | signature: None, |
| 1217 | state: Some(state), |
| 1218 | }], |
| 1219 | }); |
| 1220 | let replay = client |
| 1221 | .prepare_outbound_request(continuation.clone(), true) |
| 1222 | .unwrap(); |
| 1223 | assert!( |
| 1224 | replay.body["input"] |
| 1225 | .as_array() |
| 1226 | .unwrap() |
| 1227 | .iter() |
| 1228 | .any(|item| item["type"] == "reasoning" && item["encrypted_content"] == "enc_generic") |
| 1229 | ); |
| 1230 | assert!(!replay.body.to_string().contains("readable-private-summary")); |
| 1231 | continuation.model = "another-model".into(); |
| 1232 | let wrong_model = build_responses_body_for_provider(&continuation, ProviderKind::Openai, None); |
| 1233 | assert!(!wrong_model.to_string().contains("enc_generic")); |
| 1234 | let wrong_provider = |
| 1235 | build_responses_body_for_provider(&continuation, ProviderKind::Deepseek, None); |
| 1236 | assert!(!wrong_provider.to_string().contains("enc_generic")); |
| 1237 | } |
| 1238 | |
| 1239 | #[test] |
| 1240 | fn codex_responses_body_uses_responses_reasoning_not_deepseek_thinking() { |
| 1241 | let request = MessageRequest { |
| 1242 | model: "gpt-6-astra".to_string(), |
| 1243 | messages: vec![Message { |
| 1244 | role: Role::User, |
| 1245 | content: vec![ContentBlock::Text { |
| 1246 | text: "hello".to_string(), |
| 1247 | cache_control: None, |
| 1248 | }], |
| 1249 | }], |
| 1250 | max_tokens: 128, |
| 1251 | system: None, |
| 1252 | tools: None, |
| 1253 | tool_choice: None, |
| 1254 | metadata: None, |
| 1255 | thinking: None, |
| 1256 | reasoning_effort: Some("max".to_string()), |
| 1257 | stream: None, |
| 1258 | temperature: None, |
| 1259 | top_p: None, |
| 1260 | }; |
| 1261 | |
| 1262 | let body = build_responses_body(&request); |
| 1263 | |
| 1264 | assert_eq!( |
| 1265 | body.pointer("/reasoning/effort").and_then(Value::as_str), |
| 1266 | Some("max") |
| 1267 | ); |
| 1268 | assert_eq!( |
| 1269 | body.pointer("/reasoning/summary").and_then(Value::as_str), |
| 1270 | Some("auto") |
| 1271 | ); |
| 1272 | assert!(body.get("thinking").is_none()); |
| 1273 | assert!(body.get("reasoning_effort").is_none()); |
| 1274 | } |
| 1275 | |
| 1276 | #[test] |
| 1277 | fn responses_failed_event_reports_nested_error() { |
| 1278 | let event = json!({ |
| 1279 | "type": "response.failed", |
| 1280 | "response": { |
| 1281 | "id": "resp_123", |
| 1282 | "error": { |
| 1283 | "code": "rate_limit_exceeded", |
| 1284 | "message": "Please retry later" |
| 1285 | } |
| 1286 | } |
| 1287 | }); |
| 1288 | |
| 1289 | let (code, message) = responses_event_error_details(&event); |
| 1290 | |
| 1291 | assert_eq!(code, "rate_limit_exceeded"); |
| 1292 | assert_eq!(message, "Please retry later"); |
| 1293 | } |
| 1294 | |
| 1295 | #[test] |
| 1296 | fn responses_incomplete_event_reports_reason() { |
| 1297 | let event = json!({ |
| 1298 | "type": "response.incomplete", |
| 1299 | "response": { |
| 1300 | "id": "resp_123", |
| 1301 | "status": "incomplete", |
| 1302 | "error": null, |
| 1303 | "incomplete_details": { |
| 1304 | "reason": "content_filter" |
| 1305 | } |
| 1306 | } |
| 1307 | }); |
| 1308 | |
| 1309 | let (code, message) = responses_event_error_details(&event); |
| 1310 | |
| 1311 | assert_eq!(code, "content_filter"); |
| 1312 | assert_eq!(message, "response incomplete: content_filter"); |
| 1313 | } |
| 1314 | |
| 1315 | #[test] |
| 1316 | fn responses_incomplete_stop_reason_preserves_provider_reason() { |
| 1317 | assert_eq!( |
| 1318 | responses_stop_reason( |
| 1319 | &json!({ |
| 1320 | "status": "incomplete", |
| 1321 | "incomplete_details": { "reason": "max_output_tokens" } |
| 1322 | }), |
| 1323 | false, |
| 1324 | ), |
| 1325 | "incomplete:max_output_tokens" |
| 1326 | ); |
| 1327 | assert_eq!( |
| 1328 | responses_stop_reason(&json!({"status": "incomplete"}), false), |
| 1329 | "incomplete:max_tokens" |
| 1330 | ); |
| 1331 | } |
| 1332 | |
| 1333 | #[test] |
| 1334 | fn parse_responses_usage_derives_cache_miss_and_reasoning() { |
| 1335 | let usage = json!({ |
| 1336 | "input_tokens": 1000, |
| 1337 | "output_tokens": 200, |
| 1338 | "input_tokens_details": { "cached_tokens": 600 }, |
| 1339 | "output_tokens_details": { "reasoning_tokens": 120 } |
| 1340 | }); |
| 1341 | |
| 1342 | let parsed = parse_responses_usage(&usage); |
| 1343 | |
| 1344 | assert_eq!(parsed.input_tokens, 1000); |
| 1345 | assert_eq!(parsed.output_tokens, 200); |
| 1346 | assert_eq!(parsed.prompt_cache_hit_tokens, Some(600)); |
| 1347 | // Cache-miss is derived as input minus the cached hit when cached > 0. |
| 1348 | assert_eq!(parsed.prompt_cache_miss_tokens, Some(400)); |
| 1349 | // Reasoning surfaces from output_tokens_details (Responses dialect). |
| 1350 | assert_eq!(parsed.reasoning_tokens, Some(120)); |
| 1351 | |
| 1352 | // Without cached/reasoning details, the derived fields stay None. |
| 1353 | let bare = json!({ "input_tokens": 1000, "output_tokens": 200 }); |
| 1354 | let parsed_bare = parse_responses_usage(&bare); |
| 1355 | assert_eq!(parsed_bare.prompt_cache_hit_tokens, None); |
| 1356 | assert_eq!(parsed_bare.prompt_cache_miss_tokens, None); |
| 1357 | assert_eq!(parsed_bare.reasoning_tokens, None); |
| 1358 | } |
| 1359 | |
| 1360 | #[test] |
| 1361 | fn parse_responses_usage_saturates_u64_fields() { |
| 1362 | let parsed = parse_responses_usage(&json!({ |
| 1363 | "input_tokens": u64::MAX, |
| 1364 | "output_tokens": u64::MAX, |
| 1365 | "input_tokens_details": { "cached_tokens": u64::MAX }, |
| 1366 | "output_tokens_details": { "reasoning_tokens": u64::MAX } |
| 1367 | })); |
| 1368 | assert_eq!(parsed.input_tokens, u32::MAX); |
| 1369 | assert_eq!(parsed.output_tokens, u32::MAX); |
| 1370 | assert_eq!(parsed.prompt_cache_hit_tokens, Some(u32::MAX)); |
| 1371 | assert_eq!(parsed.prompt_cache_miss_tokens, Some(0)); |
| 1372 | assert_eq!(parsed.reasoning_tokens, Some(u32::MAX)); |
| 1373 | } |
| 1374 | |
| 1375 | #[test] |
| 1376 | fn parse_responses_usage_reads_deepseek_top_level_cache_fields() { |
| 1377 | // DeepSeek's Responses dialect reports cache telemetry as top-level |
| 1378 | // `prompt_cache_hit_tokens` / `prompt_cache_miss_tokens` with |
| 1379 | // `cache_write_tokens` nested under `input_tokens_details` -- none of |
| 1380 | // which the old parser read (it only looked at |
| 1381 | // `input_tokens_details.cached_tokens`, which DeepSeek leaves unset, |
| 1382 | // so every V4 Flash turn recorded cache_hit = None). |
| 1383 | let usage = json!({ |
| 1384 | "input_tokens": 1_000, |
| 1385 | "output_tokens": 200, |
| 1386 | "prompt_cache_hit_tokens": 600, |
| 1387 | "prompt_cache_miss_tokens": 200, |
| 1388 | "input_tokens_details": { "cached_tokens": 999, "cache_write_tokens": 100 }, |
| 1389 | "output_tokens_details": { "reasoning_tokens": 120 } |
| 1390 | }); |
| 1391 | |
| 1392 | let parsed = parse_responses_usage(&usage); |
| 1393 | |
| 1394 | // Top-level DeepSeek fields win over the nested OpenAI-style shape, |
| 1395 | // and the explicit miss is trusted over the derived fallback. |
| 1396 | assert_eq!(parsed.prompt_cache_hit_tokens, Some(600)); |
| 1397 | assert_eq!(parsed.prompt_cache_miss_tokens, Some(200)); |
| 1398 | assert_eq!(parsed.prompt_cache_write_tokens, Some(100)); |
| 1399 | // `input_tokens` remains the provider-reported total; the pricing |
| 1400 | // layer partitions it into hit / miss / write classes. |
| 1401 | assert_eq!(parsed.input_tokens, 1_000); |
| 1402 | assert_eq!(parsed.output_tokens, 200); |
| 1403 | assert_eq!(parsed.reasoning_tokens, Some(120)); |
| 1404 | |
| 1405 | // The parsed fields must reach the pricing classes unchanged: 600 hit |
| 1406 | // at the cache-read rate, 100 write at the creation rate, and the |
| 1407 | // remaining 300 (200 reported miss + 100 uncategorized) at the miss |
| 1408 | // rate -- instead of the pre-fix all-raw-input miss billing. |
| 1409 | let classes = crate::pricing::token_usage_for_pricing(&parsed); |
| 1410 | assert_eq!(classes.input, 300); |
| 1411 | assert_eq!(classes.cache_read, 600); |
| 1412 | assert_eq!(classes.cache_write, 100); |
| 1413 | } |
| 1414 | |
| 1415 | #[test] |
| 1416 | fn parse_responses_usage_keeps_old_shape_with_cache_write_fallback() { |
| 1417 | // OpenAI-style payloads still parse from `input_tokens_details` alone: |
| 1418 | // hit from `cached_tokens` (fallback), miss derived as input minus |
| 1419 | // hit, and the write class from `cache_write_tokens` when present. |
| 1420 | let usage = json!({ |
| 1421 | "input_tokens": 1_000, |
| 1422 | "output_tokens": 200, |
| 1423 | "input_tokens_details": { "cached_tokens": 600, "cache_write_tokens": 100 } |
| 1424 | }); |
| 1425 | |
| 1426 | let parsed = parse_responses_usage(&usage); |
| 1427 | |
| 1428 | assert_eq!(parsed.input_tokens, 1_000); |
| 1429 | assert_eq!(parsed.prompt_cache_hit_tokens, Some(600)); |
| 1430 | assert_eq!(parsed.prompt_cache_miss_tokens, Some(400)); |
| 1431 | assert_eq!(parsed.prompt_cache_write_tokens, Some(100)); |
| 1432 | assert_eq!(parsed.reasoning_tokens, None); |
| 1433 | } |
| 1434 | |
| 1435 | /// Regression fixture for the reasoning double-billing bug: a real |
| 1436 | /// Responses usage payload has to survive the whole way into the pricing |
| 1437 | /// conversion without reasoning tokens being charged twice. OpenAI's |
| 1438 | /// `output_tokens` is already the *total* billable completion count, with |
| 1439 | /// `output_tokens_details.reasoning_tokens` a subset of it. |
| 1440 | #[test] |
| 1441 | fn responses_usage_reaches_pricing_conversion_without_double_billing_reasoning() { |
| 1442 | use crate::config::ProviderKind; |
| 1443 | use crate::pricing::{calculate_turn_cost_estimate_for_provider, token_usage_for_pricing}; |
| 1444 | |
| 1445 | let usage = parse_responses_usage(&json!({ |
| 1446 | "input_tokens": 10_000, |
| 1447 | "output_tokens": 4_000, |
| 1448 | "total_tokens": 14_000, |
| 1449 | "input_tokens_details": { "cached_tokens": 6_000 }, |
| 1450 | "output_tokens_details": { "reasoning_tokens": 3_500 } |
| 1451 | })); |
| 1452 | |
| 1453 | let classes = token_usage_for_pricing(&usage); |
| 1454 | assert_eq!(classes.output, 4_000, "reasoning must not inflate output"); |
| 1455 | assert_eq!(classes.input, 4_000); |
| 1456 | assert_eq!(classes.cache_read, 6_000); |
| 1457 | assert_eq!(classes.cache_write, 0); |
| 1458 | |
| 1459 | // gpt-5.5: 0.50 cache-read / 5.00 input / 30.00 output per million. |
| 1460 | let cost = calculate_turn_cost_estimate_for_provider(ProviderKind::Openai, "gpt-5.5", &usage) |
| 1461 | .expect("direct OpenAI route is priced"); |
| 1462 | let expected = 0.006 * 0.50 + 0.004 * 5.00 + 0.004 * 30.00; |
| 1463 | assert!( |
| 1464 | (cost.usd - expected).abs() < 1e-12, |
| 1465 | "expected {expected}, got {}", |
| 1466 | cost.usd |
| 1467 | ); |
| 1468 | |
| 1469 | // The bug charged the 3_500 reasoning tokens a second time at the |
| 1470 | // output rate; assert the difference explicitly so a reintroduction is |
| 1471 | // unambiguous rather than a silent number change. |
| 1472 | let double_billed = expected + 0.0035 * 30.00; |
| 1473 | assert!((cost.usd - double_billed).abs() > 1e-6); |
| 1474 | } |
| 1475 | |
| 1476 | #[test] |
| 1477 | fn responses_input_includes_user_role_tool_results() { |
| 1478 | let request = MessageRequest { |
| 1479 | model: "gpt-5.5".to_string(), |
| 1480 | messages: vec![ |
| 1481 | Message { |
| 1482 | role: Role::Assistant, |
| 1483 | content: vec![ContentBlock::ToolUse { |
| 1484 | execution_id: None, |
| 1485 | id: "call_abc|fc_123".to_string(), |
| 1486 | name: "checklist_write".to_string(), |
| 1487 | input: json!({"items": []}), |
| 1488 | caller: None, |
| 1489 | thought_signature: None, |
| 1490 | }], |
| 1491 | }, |
| 1492 | Message { |
| 1493 | role: Role::User, |
| 1494 | content: vec![ContentBlock::ToolResult { |
| 1495 | execution_id: None, |
| 1496 | tool_use_id: "call_abc|fc_123".to_string(), |
| 1497 | content: "<6 items>".to_string(), |
| 1498 | is_error: None, |
| 1499 | content_blocks: None, |
| 1500 | }], |
| 1501 | }, |
| 1502 | ], |
| 1503 | max_tokens: 128, |
| 1504 | system: None, |
| 1505 | tools: None, |
| 1506 | tool_choice: None, |
| 1507 | metadata: None, |
| 1508 | thinking: None, |
| 1509 | reasoning_effort: None, |
| 1510 | stream: None, |
| 1511 | temperature: None, |
| 1512 | top_p: None, |
| 1513 | }; |
| 1514 | |
| 1515 | let input = convert_messages_to_responses_input(&request, ProviderKind::OpenaiCodex, None); |
| 1516 | |
| 1517 | assert_eq!(input[0]["type"], "function_call"); |
| 1518 | assert_eq!(input[0]["call_id"], "call_abc"); |
| 1519 | assert_eq!(input[0]["name"], "checklist_write"); |
| 1520 | assert_eq!(input[0]["namespace"], "codewhale"); |
| 1521 | assert_eq!(input[1]["type"], "function_call_output"); |
| 1522 | assert_eq!(input[1]["call_id"], "call_abc"); |
| 1523 | assert_eq!(input[1]["output"], "<6 items>"); |
| 1524 | } |
| 1525 | |
| 1526 | #[test] |
| 1527 | fn responses_input_encodes_tool_call_names() { |
| 1528 | let request = MessageRequest { |
| 1529 | model: "gpt-5.5".to_string(), |
| 1530 | messages: vec![Message { |
| 1531 | role: Role::Assistant, |
| 1532 | content: vec![ContentBlock::ToolUse { |
| 1533 | execution_id: None, |
| 1534 | id: "call_abc|fc_123".to_string(), |
| 1535 | name: "web.run".to_string(), |
| 1536 | input: json!({}), |
| 1537 | caller: None, |
| 1538 | thought_signature: None, |
| 1539 | }], |
| 1540 | }], |
| 1541 | max_tokens: 128, |
| 1542 | system: None, |
| 1543 | tools: None, |
| 1544 | tool_choice: None, |
| 1545 | metadata: None, |
| 1546 | thinking: None, |
| 1547 | reasoning_effort: None, |
| 1548 | stream: None, |
| 1549 | temperature: None, |
| 1550 | top_p: None, |
| 1551 | }; |
| 1552 | |
| 1553 | let input = convert_messages_to_responses_input(&request, ProviderKind::OpenaiCodex, None); |
| 1554 | |
| 1555 | assert_eq!(input[0]["type"], "function_call"); |
| 1556 | assert_eq!(input[0]["name"], to_api_tool_name("web.run")); |
| 1557 | assert_eq!(input[0]["namespace"], "codewhale"); |
| 1558 | let generic = convert_messages_to_responses_input(&request, ProviderKind::Openai, None); |
| 1559 | assert!(generic[0].get("namespace").is_none()); |
| 1560 | } |
| 1561 | |
| 1562 | #[test] |
| 1563 | fn responses_function_tool_sanitizes_root_composition_schema() { |
| 1564 | let tool = Tool { |
| 1565 | tool_type: None, |
| 1566 | name: "web.run".to_string(), |
| 1567 | description: "Apply patch".to_string(), |
| 1568 | input_schema: json!({ |
| 1569 | "type": "object", |
| 1570 | "properties": { |
| 1571 | "patch": {"type": "string"}, |
| 1572 | "replace": {"type": "array"}, |
| 1573 | "changes": {"type": "array"} |
| 1574 | }, |
| 1575 | "oneOf": [ |
| 1576 | {"required": ["patch"]}, |
| 1577 | {"required": ["replace"]}, |
| 1578 | {"required": ["changes"]} |
| 1579 | ] |
| 1580 | }), |
| 1581 | allowed_callers: None, |
| 1582 | defer_loading: None, |
| 1583 | input_examples: None, |
| 1584 | strict: None, |
| 1585 | cache_control: None, |
| 1586 | }; |
| 1587 | |
| 1588 | let payload = tool_to_responses_function(&tool); |
| 1589 | let parameters = &payload["parameters"]; |
| 1590 | |
| 1591 | assert_eq!(payload["name"], to_api_tool_name("web.run")); |
| 1592 | assert_eq!(parameters["type"], "object"); |
| 1593 | assert!(parameters.get("oneOf").is_none()); |
| 1594 | assert!(parameters.get("anyOf").is_none()); |
| 1595 | assert!(parameters.get("allOf").is_none()); |
| 1596 | assert!(parameters.get("enum").is_none()); |
| 1597 | assert!(parameters.get("not").is_none()); |
| 1598 | assert!(parameters["properties"].get("patch").is_some()); |
| 1599 | assert!(parameters["properties"].get("replace").is_some()); |
| 1600 | assert!(parameters["properties"].get("changes").is_some()); |
| 1601 | assert_eq!( |
| 1602 | payload["description"], |
| 1603 | "Apply patch\n\nExactly one of these parameter groups must be provided: `changes` | `patch` | `replace`." |
| 1604 | ); |
| 1605 | assert!(tool.input_schema.get("oneOf").is_some()); |
| 1606 | } |
| 1607 | |
| 1608 | #[test] |
| 1609 | fn responses_function_tool_trims_description_before_constraint_note() { |
| 1610 | let tool = Tool { |
| 1611 | tool_type: None, |
| 1612 | name: "apply_patch".to_string(), |
| 1613 | description: "Apply patch\n".to_string(), |
| 1614 | input_schema: json!({ |
| 1615 | "type": "object", |
| 1616 | "properties": { |
| 1617 | "patch": {"type": "string"}, |
| 1618 | "replace": {"type": "array"}, |
| 1619 | "changes": {"type": "array"} |
| 1620 | }, |
| 1621 | "oneOf": [ |
| 1622 | {"required": ["patch"]}, |
| 1623 | {"required": ["replace"]}, |
| 1624 | {"required": ["changes"]} |
| 1625 | ] |
| 1626 | }), |
| 1627 | allowed_callers: None, |
| 1628 | defer_loading: None, |
| 1629 | input_examples: None, |
| 1630 | strict: None, |
| 1631 | cache_control: None, |
| 1632 | }; |
| 1633 | |
| 1634 | let payload = tool_to_responses_function(&tool); |
| 1635 | |
| 1636 | assert_eq!( |
| 1637 | payload["description"], |
| 1638 | "Apply patch\n\nExactly one of these parameter groups must be provided: `changes` | `patch` | `replace`." |
| 1639 | ); |
| 1640 | } |
| 1641 | |
| 1642 | #[test] |
| 1643 | fn responses_function_tool_leaves_description_unchanged_without_constraint_note() { |
| 1644 | let tool = Tool { |
| 1645 | tool_type: None, |
| 1646 | name: "lookup".to_string(), |
| 1647 | description: "Lookup".to_string(), |
| 1648 | input_schema: json!({ |
| 1649 | "type": "object", |
| 1650 | "properties": { |
| 1651 | "query": {"type": "string"} |
| 1652 | } |
| 1653 | }), |
| 1654 | allowed_callers: None, |
| 1655 | defer_loading: None, |
| 1656 | input_examples: None, |
| 1657 | strict: None, |
| 1658 | cache_control: None, |
| 1659 | }; |
| 1660 | |
| 1661 | let payload = tool_to_responses_function(&tool); |
| 1662 | |
| 1663 | assert_eq!(payload["description"], "Lookup"); |
| 1664 | } |
| 1665 | |
| 1666 | /// The Responses API projection of [`ContentBlock::ImageUrl`]. |
| 1667 | /// |
| 1668 | /// Responses is the odd one out: the image part carries `image_url` as a bare |
| 1669 | /// string rather than the nested object Chat Completions uses. Getting that |
| 1670 | /// wrong produces a schema error from OpenAI rather than anything that names |
| 1671 | /// the image, so it is worth pinning explicitly. |
| 1672 | #[test] |
| 1673 | fn user_image_becomes_an_input_image_item() { |
| 1674 | const DATA_URL: &str = "data:image/png;base64,QUJD"; |
| 1675 | |
| 1676 | let mut request = minimal_responses_request(); |
| 1677 | request.messages[0].content.push(ContentBlock::ImageUrl { |
| 1678 | image_url: codewhale_models::ImageUrlContent { |
| 1679 | url: DATA_URL.to_string(), |
| 1680 | }, |
| 1681 | }); |
| 1682 | |
| 1683 | let items = convert_messages_to_responses_input(&request, ProviderKind::OpenaiCodex, None); |
| 1684 | |
| 1685 | let user = items |
| 1686 | .iter() |
| 1687 | .find(|item| item["role"] == "user") |
| 1688 | .expect("a user item"); |
| 1689 | let content = user["content"].as_array().expect("content items"); |
| 1690 | |
| 1691 | let image = content |
| 1692 | .iter() |
| 1693 | .find(|part| part["type"] == "input_image") |
| 1694 | .expect("an input_image part"); |
| 1695 | assert_eq!( |
| 1696 | image["image_url"], DATA_URL, |
| 1697 | "Responses takes image_url as a bare string, not a nested object: {image}" |
| 1698 | ); |
| 1699 | |
| 1700 | assert!( |
| 1701 | content.iter().any(|part| part["type"] == "input_text"), |
| 1702 | "the accompanying question must survive: {user}" |
| 1703 | ); |
| 1704 | } |
| 1705 | |
| 1706 | #[test] |
| 1707 | fn tool_result_image_becomes_native_function_output_content() { |
| 1708 | let mut request = minimal_responses_request(); |
| 1709 | request.messages = vec![ |
| 1710 | Message { |
| 1711 | role: Role::Assistant, |
| 1712 | content: vec![ContentBlock::ToolUse { |
| 1713 | execution_id: None, |
| 1714 | id: "call_image_1".to_string(), |
| 1715 | name: "read".to_string(), |
| 1716 | input: serde_json::json!({"path": "shot.png"}), |
| 1717 | caller: None, |
| 1718 | thought_signature: None, |
| 1719 | }], |
| 1720 | }, |
| 1721 | Message { |
| 1722 | role: Role::User, |
| 1723 | content: vec![ContentBlock::ToolResult { |
| 1724 | execution_id: None, |
| 1725 | tool_use_id: "call_image_1".to_string(), |
| 1726 | content: "screenshot captured".to_string(), |
| 1727 | is_error: Some(false), |
| 1728 | content_blocks: Some(vec![serde_json::json!({ |
| 1729 | "type": "image", |
| 1730 | "mime_type": "image/png", |
| 1731 | "data": "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR4nGP4z8DwHwAFAAH/iZk9HQAAAABJRU5ErkJggg==", |
| 1732 | })]), |
| 1733 | }], |
| 1734 | }, |
| 1735 | ]; |
| 1736 | |
| 1737 | let items = convert_messages_to_responses_input(&request, ProviderKind::OpenaiCodex, None); |
| 1738 | let output = items |
| 1739 | .iter() |
| 1740 | .find(|item| item["type"] == "function_call_output") |
| 1741 | .expect("function output"); |
| 1742 | let content = output["output"].as_array().expect("rich output array"); |
| 1743 | |
| 1744 | assert_eq!( |
| 1745 | content[0], |
| 1746 | serde_json::json!({ |
| 1747 | "type": "input_text", |
| 1748 | "text": "screenshot captured", |
| 1749 | }) |
| 1750 | ); |
| 1751 | assert_eq!(content[1]["type"], "input_image"); |
| 1752 | assert_eq!( |
| 1753 | content[1]["image_url"], |
| 1754 | "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR4nGP4z8DwHwAFAAH/iZk9HQAAAABJRU5ErkJggg==" |
| 1755 | ); |
| 1756 | } |
| 1757 | |
| 1758 | /// A `system`-role history message — the shape a compaction summary, a branch |
| 1759 | /// summary, or an imported journal `system` entry takes once it reaches |
| 1760 | /// `MessageRequest::messages` — must survive the Responses conversion. The |
| 1761 | /// Chat Completions adapter already keeps it |
| 1762 | /// (`request_builder_preserves_internal_system_messages`); dropping it here |
| 1763 | /// silently deletes the only record of everything the compaction replaced. |
| 1764 | #[test] |
| 1765 | fn responses_input_preserves_system_history_with_the_provider_role() { |
| 1766 | let mut request = minimal_responses_request(); |
| 1767 | request.messages.insert( |
| 1768 | 0, |
| 1769 | Message { |
| 1770 | role: Role::System, |
| 1771 | content: vec![ContentBlock::Text { |
| 1772 | text: "[compaction summary] the user is porting the parser".to_string(), |
| 1773 | cache_control: None, |
| 1774 | }], |
| 1775 | }, |
| 1776 | ); |
| 1777 | |
| 1778 | for (provider, role) in [ |
| 1779 | (ProviderKind::OpenaiCodex, "developer"), |
| 1780 | (ProviderKind::Openai, "system"), |
| 1781 | ] { |
| 1782 | let items = convert_messages_to_responses_input(&request, provider, None); |
| 1783 | let system = items |
| 1784 | .iter() |
| 1785 | .find(|item| item["role"] == role) |
| 1786 | .expect("system history survives under the provider-supported role"); |
| 1787 | assert_eq!(system["type"], "message"); |
| 1788 | assert_eq!( |
| 1789 | system["content"][0], |
| 1790 | serde_json::json!({ |
| 1791 | "type": "input_text", |
| 1792 | "text": "[compaction summary] the user is porting the parser", |
| 1793 | }) |
| 1794 | ); |
| 1795 | } |
| 1796 | } |
| 1797 |