返回 CodeWhale
tests.rs
根目录 / crates / tui / src / client / responses / tests.rs
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
1797 lines RUST