返回 CodeWhale
catalog_tests.rs
根目录 / crates / tui / src / client / catalog_tests.rs
1 //! Local HTTP regressions for complete, bounded provider catalog observations.
2
3 use super::tests::{
4 custom_mock_client_for_identity, mount_models_json, opencode_go_client_for,
5 openrouter_client_for,
6 };
7 use super::*;
8 use crate::config::{ProviderConfig, ProvidersConfig};
9 use tokio::io::{AsyncReadExt, AsyncWriteExt};
10 use tokio::net::{TcpListener, TcpStream};
11 use wiremock::matchers::{method, path};
12 use wiremock::{Mock, MockServer, Request, ResponseTemplate};
13
14 const KEY: &str = "catalog-key-canary-7f092";
15 const CURSOR: &str = "cursor/second +?&=雪-canary";
16
17 #[tokio::test]
18 async fn chatgpt_models_http_uses_visible_roster_and_rejects_secret_labels() {
19 let server = MockServer::start().await;
20 let client = CodewhaleClient::new(&Config {
21 provider: Some("openai-codex".into()),
22 providers: Some(ProvidersConfig {
23 openai_codex: ProviderConfig {
24 api_key: Some(KEY.into()),
25 base_url: Some(format!("{}/v1", server.uri())),
26 ..ProviderConfig::default()
27 },
28 ..ProvidersConfig::default()
29 }),
30 ..Config::default()
31 })
32 .expect("explicit local ChatGPT protocol fixture");
33 mount_models_json(
34 &server,
35 200,
36 json!({"models":[
37 {"slug":"gpt-z","display_name":"GPT Z","visibility":"list"},
38 {"slug":"gpt-hidden","display_name":"Hidden","visibility":"hidden"},
39 {"slug":"gpt-a","display_name":"GPT A","visibility":"list"}
40 ]}),
41 )
42 .await;
43 let models = client.list_models().await.expect("account roster");
44 assert_eq!(
45 models.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
46 ["gpt-z", "gpt-a"]
47 );
48 assert_eq!(models[0].display_name.as_deref(), Some("GPT Z"));
49 let delta = client
50 .fetch_catalog_delta()
51 .await
52 .expect("same parsed roster");
53 assert_eq!(
54 delta
55 .offerings
56 .iter()
57 .map(|row| row.wire_model_id.as_str())
58 .collect::<Vec<_>>(),
59 ["gpt-z", "gpt-a"]
60 );
61 for request in server.received_requests().await.expect("requests") {
62 assert_eq!(request.url.path(), "/v1/models");
63 assert_eq!(
64 request.headers.get("authorization").unwrap(),
65 &format!("Bearer {KEY}")
66 );
67 }
68 server.reset().await;
69 mount_models_json(
70 &server,
71 200,
72 json!({"models":[
73 {"slug":"gpt-a","display_name":KEY,"visibility":"list"}
74 ]}),
75 )
76 .await;
77 let error = client
78 .list_models()
79 .await
80 .expect_err("credential echo is never cached");
81 assert!(!error.to_string().contains(KEY));
82 server.reset().await;
83 Mock::given(method("POST"))
84 .and(path("/v1/responses"))
85 .respond_with(ResponseTemplate::new(307).insert_header("Location", "/other"))
86 .mount(&server)
87 .await;
88 let response = client
89 .http_client
90 .post(format!("{}/v1/responses", server.uri()))
91 .send()
92 .await
93 .expect("local response");
94 assert_eq!(response.status().as_u16(), 307);
95 assert_eq!(
96 server
97 .received_requests()
98 .await
99 .expect("request count")
100 .len(),
101 1,
102 "a plan grant must not follow a redirected inference endpoint"
103 );
104 }
105
106 fn anthropic_client(base_url: &str) -> CodewhaleClient {
107 let mut client = CodewhaleClient::new(&Config {
108 provider: Some("anthropic".into()),
109 providers: Some(ProvidersConfig {
110 anthropic: ProviderConfig {
111 api_key: Some(KEY.into()),
112 base_url: Some(base_url.into()),
113 http_headers: Some(HashMap::from([(
114 "x-private-fixture".into(),
115 "custom-header-canary".into(),
116 )])),
117 ..ProviderConfig::default()
118 },
119 ..ProvidersConfig::default()
120 }),
121 ..Config::default()
122 })
123 .expect("explicit local Anthropic fixture client");
124 client.retry.enabled = false;
125 client.retry.max_retries = 0;
126 client
127 }
128
129 async fn mount_page(server: &MockServer, cursor: Option<&str>, response: ResponseTemplate) {
130 let cursor = cursor.map(str::to_owned);
131 Mock::given(method("GET"))
132 .and(path("/v1/models"))
133 .and(move |request: &Request| {
134 request
135 .url
136 .query_pairs()
137 .find(|(key, _)| key == "after_id")
138 .map(|(_, value)| value.into_owned())
139 == cursor
140 })
141 .respond_with(response)
142 .mount(server)
143 .await;
144 }
145
146 fn page(data: Value, next: Option<&str>) -> ResponseTemplate {
147 ResponseTemplate::new(200).set_body_json(json!({
148 "data": data, "has_more": next.is_some(), "last_id": next,
149 }))
150 }
151
152 // This is an explicit fixture continuation contract, not a claim that any
153 // provider other than Anthropic supports after_id in production.
154 async fn collect_fixture(
155 server: &MockServer,
156 limits: ModelsFetchLimits,
157 ) -> Result<String, ModelsFetchError> {
158 let http = crate::tls::reqwest_client_builder()
159 .redirect(reqwest::redirect::Policy::none())
160 .build()
161 .unwrap();
162 collect_models_document(
163 reqwest::Url::parse(&format!("{}/v1/models", server.uri())).unwrap(),
164 Some("after_id"),
165 limits,
166 |url| {
167 let http = http.clone();
168 async move {
169 http.get(url)
170 .send()
171 .await
172 .map_err(|_| CatalogRefreshError::Network.into())
173 }
174 },
175 )
176 .await
177 .map(|(body, _)| body)
178 }
179
180 fn assert_invalid(result: Result<String, ModelsFetchError>) {
181 assert_eq!(
182 result
183 .expect_err("incomplete or invalid observation must fail")
184 .into_catalog(),
185 CatalogRefreshError::InvalidResponse
186 );
187 }
188
189 #[tokio::test]
190 async fn anthropic_after_id_completes_both_public_consumers_with_frozen_headers() {
191 let server = MockServer::start().await;
192 mount_page(
193 &server,
194 None,
195 page(
196 json!([
197 {"id":"z-model", "owned_by":"first-owner", "created":7},
198 {"id":"a-model"}
199 ]),
200 Some(CURSOR),
201 ),
202 )
203 .await;
204 mount_page(
205 &server,
206 Some(CURSOR),
207 page(
208 json!([
209 {"id":"middle-model", "owned_by":"second-owner", "created":9},
210 {"id":"z-model", "owned_by":"later-owner", "created":10}
211 ]),
212 None,
213 ),
214 )
215 .await;
216 let client = anthropic_client(&server.uri());
217
218 let listed = client.list_models().await.unwrap();
219 assert_eq!(
220 listed.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
221 ["a-model", "middle-model", "z-model"]
222 );
223 assert_eq!(listed[1].owned_by.as_deref(), Some("second-owner"));
224 assert_eq!(listed[1].created, Some(9));
225 assert_eq!(listed[2].owned_by.as_deref(), Some("first-owner"));
226 let delta = client.fetch_catalog_delta().await.unwrap();
227 assert_eq!(delta.provider, "anthropic");
228 assert_eq!(
229 delta.base_url_fingerprint,
230 base_url_fingerprint(&server.uri())
231 );
232 assert_eq!(
233 delta
234 .offerings
235 .iter()
236 .map(|row| row.wire_model_id.as_str())
237 .collect::<Vec<_>>(),
238 ["a-model", "middle-model", "z-model"]
239 );
240 for row in &delta.offerings {
241 assert!(
242 matches!(&row.source, CatalogSource::Live { base_url_fingerprint, fetched_at }
243 if base_url_fingerprint == &delta.base_url_fingerprint && *fetched_at == delta.fetched_at)
244 );
245 }
246 let requests = server.received_requests().await.unwrap();
247 assert_eq!(requests.len(), 4);
248 for (index, request) in requests.iter().enumerate() {
249 assert_eq!(request.headers.get("x-api-key").unwrap(), KEY);
250 assert_eq!(
251 request.headers.get("anthropic-version").unwrap(),
252 "2023-06-01"
253 );
254 assert_eq!(
255 request.headers.get("x-private-fixture").unwrap(),
256 "custom-header-canary"
257 );
258 assert!(request.headers.get("authorization").is_none());
259 let pairs = request.url.query_pairs().collect::<Vec<_>>();
260 if index % 2 == 0 {
261 assert!(pairs.is_empty());
262 } else {
263 assert_eq!(pairs.len(), 1);
264 assert_eq!(pairs[0].0, "after_id");
265 assert_eq!(pairs[0].1, CURSOR);
266 }
267 }
268 }
269
270 #[tokio::test]
271 async fn unpaginated_rosters_stay_single_request_and_unknown_continuation_refuses() {
272 for go in [false, true] {
273 let server = MockServer::start().await;
274 let id = crate::config::opencode_go_models()[0];
275 mount_models_json(&server, 200, json!({"data":[{"id":id}]})).await;
276 let mut client = if go {
277 opencode_go_client_for(&server)
278 } else {
279 openrouter_client_for(&server)
280 };
281 client.retry.enabled = false;
282 client.retry.max_retries = 0;
283 assert_eq!(client.list_models().await.unwrap().len(), 1);
284 assert_eq!(
285 client.fetch_catalog_delta().await.unwrap().offerings.len(),
286 1
287 );
288 assert_eq!(server.received_requests().await.unwrap().len(), 2);
289 server.reset().await;
290 mount_models_json(
291 &server,
292 200,
293 json!({"data":[{"id":id}], "has_more":true, "last_id":"unsupported-next"}),
294 )
295 .await;
296 assert!(client.list_models().await.is_err());
297 assert_eq!(
298 client.fetch_catalog_delta().await.unwrap_err(),
299 CatalogRefreshError::InvalidResponse
300 );
301 let requests = server.received_requests().await.unwrap();
302 assert_eq!(
303 requests.len(),
304 2,
305 "unsupported contract must not speculate a second request"
306 );
307 assert!(requests.iter().all(|request| request.url.query().is_none()));
308 }
309 }
310
311 #[tokio::test]
312 async fn opencode_go_published_unpaginated_roster_keeps_documented_protocols() {
313 let server = MockServer::start().await;
314 let base_url = format!("{}/zen/go/v1", server.uri());
315 // Literal additions from the pinned Go documentation plus retained routes:
316 // do not generate this fixture from the production allowlist it verifies.
317 let positives = [
318 "glm-5.3-flash",
319 "glm-5.3",
320 "longcat-2.0",
321 "deepseek-v4-flash-vision-exp",
322 "hy4-preview",
323 "hy3",
324 "omen-alpha",
325 "deepseek-v4-pro",
326 "grok-4.5",
327 "qwen3.8-max",
328 "qwen3.8-flash",
329 "minimax-m3",
330 "grok-4.6",
331 "gpt-5.6-luna",
332 "muse-spark-1.3-contributor",
333 "muse-spark-1.2-contributor",
334 ];
335 let negatives = ["gpt-unlisted", "claude-unproven"];
336 let rows: Vec<_> = positives
337 .iter()
338 .chain(negatives.iter())
339 .map(
340 |id| json!({"id":id, "object":"model", "created":1_700_000_000, "owned_by":"opencode"}),
341 )
342 .collect();
343 Mock::given(method("GET"))
344 .and(path("/zen/go/v1/models"))
345 .respond_with(
346 ResponseTemplate::new(200).set_body_json(json!({"object":"list", "data":rows})),
347 )
348 .mount(&server)
349 .await;
350 let mut client = CodewhaleClient::new(&Config {
351 provider: Some("opencode-go".into()),
352 providers: Some(ProvidersConfig {
353 opencode_go: ProviderConfig {
354 api_key: Some(KEY.into()),
355 base_url: Some(base_url.clone()),
356 ..ProviderConfig::default()
357 },
358 ..ProvidersConfig::default()
359 }),
360 ..Config::default()
361 })
362 .expect("explicit local Go route");
363 client.retry.enabled = false;
364 client.retry.max_retries = 0;
365
366 let expected: std::collections::BTreeSet<_> = positives.into_iter().collect();
367 let listed = client.list_models().await.unwrap();
368 assert_eq!(
369 listed
370 .iter()
371 .map(|row| row.id.as_str())
372 .collect::<std::collections::BTreeSet<_>>(),
373 expected
374 );
375 assert_eq!(listed.len(), expected.len());
376 assert_eq!(server.received_requests().await.unwrap().len(), 1);
377 let delta = client.fetch_catalog_delta().await.unwrap();
378 assert_eq!(delta.provider, "opencode-go");
379 assert_eq!(delta.base_url_fingerprint, base_url_fingerprint(&base_url));
380 assert_eq!(
381 delta
382 .offerings
383 .iter()
384 .map(|row| row.wire_model_id.as_str())
385 .collect::<std::collections::BTreeSet<_>>(),
386 expected
387 );
388 assert_eq!(delta.offerings.len(), expected.len());
389 for row in &delta.offerings {
390 assert_eq!(row.provider, "opencode-go");
391 assert_eq!(
392 Some(row.endpoint_key.as_str()),
393 codewhale_config::opencode_go_endpoint_key(&row.wire_model_id)
394 );
395 assert_eq!(row.canonical_model, None);
396 assert_eq!(row.family, None);
397 assert_eq!(row.limit, None);
398 assert_eq!(row.cost, None);
399 assert_eq!(row.cost_source, None);
400 assert_eq!(row.modalities, None);
401 assert_eq!(row.attachment, None);
402 assert_eq!(row.reasoning, None);
403 assert_eq!(row.tool_call, None);
404 assert_eq!(row.structured_output, None);
405 assert!(row.reasoning_options.is_empty());
406 assert!(
407 matches!(&row.source, CatalogSource::Live { base_url_fingerprint, fetched_at }
408 if base_url_fingerprint == &delta.base_url_fingerprint && *fetched_at == delta.fetched_at)
409 );
410 }
411 let requests = server.received_requests().await.unwrap();
412 assert_eq!(
413 requests.len(),
414 2,
415 "one unpaginated request per public consumer"
416 );
417 for request in requests {
418 assert_eq!(request.url.path(), "/zen/go/v1/models");
419 assert!(request.url.query().is_none());
420 assert_eq!(
421 request.headers.get("authorization").unwrap(),
422 format!("Bearer {KEY}").as_str()
423 );
424 }
425 }
426
427 #[tokio::test]
428 async fn models_redirects_never_reach_another_server_from_any_consumer() {
429 let destination = MockServer::start().await;
430 for status in [302, 307] {
431 let server = MockServer::start().await;
432 Mock::given(method("GET"))
433 .and(path("/v1/models"))
434 .respond_with(
435 ResponseTemplate::new(status)
436 .insert_header("location", format!("{}/capture", destination.uri())),
437 )
438 .mount(&server)
439 .await;
440 assert!(anthropic_client(&server.uri()).list_models().await.is_err());
441 assert_eq!(
442 anthropic_client(&server.uri())
443 .fetch_catalog_delta()
444 .await
445 .unwrap_err(),
446 CatalogRefreshError::Network
447 );
448 assert!(
449 !anthropic_client(&server.uri())
450 .health_check()
451 .await
452 .unwrap()
453 );
454 let recovery = anthropic_client(&server.uri());
455 {
456 let mut health = recovery.connection_health.lock().await;
457 apply_request_failure(&mut health, Instant::now());
458 apply_request_failure(&mut health, Instant::now());
459 }
460 recovery.maybe_probe_recovery().await;
461 assert!(recovery.connection_health.lock().await.last_probe.is_some());
462 assert!(
463 verify_provider_api_key(ProviderKind::Anthropic, KEY, &server.uri())
464 .await
465 .is_err()
466 );
467 assert_eq!(server.received_requests().await.unwrap().len(), 5);
468 }
469 // Redirects after a valid page are refused by both traversal consumers too.
470 let server = MockServer::start().await;
471 mount_page(
472 &server,
473 None,
474 page(json!([{"id":"first-model"}]), Some(CURSOR)),
475 )
476 .await;
477 mount_page(
478 &server,
479 Some(CURSOR),
480 ResponseTemplate::new(308)
481 .insert_header("location", format!("{}/capture", destination.uri())),
482 )
483 .await;
484 assert!(anthropic_client(&server.uri()).list_models().await.is_err());
485 assert_eq!(
486 anthropic_client(&server.uri())
487 .fetch_catalog_delta()
488 .await
489 .unwrap_err(),
490 CatalogRefreshError::Network
491 );
492 assert_eq!(server.received_requests().await.unwrap().len(), 4);
493 assert!(
494 destination.received_requests().await.unwrap().is_empty(),
495 "no auth/header/cursor may reach a redirect target"
496 );
497 }
498
499 #[tokio::test]
500 async fn malformed_continuations_and_cursor_cycles_fail_without_partial_success() {
501 for body in [
502 json!({"data":[], "has_more":true}),
503 json!({"data":[], "has_more":true, "last_id":""}),
504 json!({"data":[], "has_more":true, "last_id":42}),
505 json!({"data":[], "has_more":"true", "last_id":"x"}),
506 json!({"data":{}, "has_more":false}),
507 json!({"data":[], "has_more":true, "last_id":"12345"}),
508 ] {
509 let server = MockServer::start().await;
510 mount_models_json(&server, 200, body).await;
511 assert_invalid(
512 collect_fixture(
513 &server,
514 ModelsFetchLimits {
515 cursor_bytes: 4,
516 ..MODELS_FETCH_LIMITS
517 },
518 )
519 .await,
520 );
521 assert_eq!(server.received_requests().await.unwrap().len(), 1);
522 }
523 for cycle in [false, true] {
524 let server = MockServer::start().await;
525 mount_page(&server, None, page(json!([{"id":"first"}]), Some("a"))).await;
526 mount_page(
527 &server,
528 Some("a"),
529 page(json!([]), Some(if cycle { "b" } else { "a" })),
530 )
531 .await;
532 if cycle {
533 mount_page(&server, Some("b"), page(json!([]), Some("a"))).await;
534 }
535 assert_invalid(collect_fixture(&server, MODELS_FETCH_LIMITS).await);
536 assert_eq!(
537 server.received_requests().await.unwrap().len(),
538 if cycle { 3 } else { 2 }
539 );
540 }
541 }
542
543 #[tokio::test]
544 async fn cumulative_raw_bytes_rows_and_pages_are_limits_not_truncation() {
545 let first = r#"{"data":[{"id":"same"},{"id":"same"}],"has_more":true,"last_id":"next"}"#;
546 let second = r#"{"data":[{"id":"same"}],"has_more":false}"#;
547 for limits in [
548 ModelsFetchLimits {
549 bytes: first.len() + second.len() - 1,
550 ..MODELS_FETCH_LIMITS
551 },
552 ModelsFetchLimits {
553 rows: 2,
554 ..MODELS_FETCH_LIMITS
555 },
556 ModelsFetchLimits {
557 pages: 1,
558 ..MODELS_FETCH_LIMITS
559 },
560 ] {
561 let server = MockServer::start().await;
562 mount_page(
563 &server,
564 None,
565 ResponseTemplate::new(200).set_body_raw(first, "application/json"),
566 )
567 .await;
568 mount_page(
569 &server,
570 Some("next"),
571 ResponseTemplate::new(200).set_body_raw(second, "application/json"),
572 )
573 .await;
574 assert_invalid(collect_fixture(&server, limits).await);
575 assert_eq!(
576 server.received_requests().await.unwrap().len(),
577 if limits.pages == 1 { 1 } else { 2 }
578 );
579 }
580 let server = MockServer::start().await;
581 mount_page(
582 &server,
583 None,
584 ResponseTemplate::new(200).set_body_raw(first, "application/json"),
585 )
586 .await;
587 mount_page(
588 &server,
589 Some("next"),
590 ResponseTemplate::new(200).set_body_raw(second, "application/json"),
591 )
592 .await;
593 let body = collect_fixture(
594 &server,
595 ModelsFetchLimits {
596 bytes: first.len() + second.len(),
597 rows: 3,
598 pages: 2,
599 ..MODELS_FETCH_LIMITS
600 },
601 )
602 .await
603 .unwrap();
604 assert_eq!(
605 parse_models_response(&body).unwrap().len(),
606 1,
607 "raw duplicates count toward limits before final deduplication"
608 );
609 }
610
611 async fn read_head(stream: &mut TcpStream) -> String {
612 let mut bytes = Vec::new();
613 let mut chunk = [0; 1024];
614 while !bytes.windows(4).any(|window| window == b"\r\n\r\n") {
615 let count = stream.read(&mut chunk).await.unwrap();
616 assert!(
617 count > 0 && bytes.len() + count <= 16_384,
618 "bounded fixture request header"
619 );
620 bytes.extend_from_slice(&chunk[..count]);
621 }
622 String::from_utf8(bytes).unwrap()
623 }
624
625 #[tokio::test]
626 async fn chunked_catalog_body_is_bounded_without_content_length() {
627 let body = r#"{"data":[{"id":"chunked-model"}]}"#;
628 for limit in [20, body.len()] {
629 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
630 let endpoint = format!("http://{}/v1/models", listener.local_addr().unwrap());
631 let (first, second) = body.split_at(12);
632 let response = format!(
633 "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n{:x}\r\n{first}\r\n{:x}\r\n{second}\r\n0\r\n\r\n",
634 first.len(),
635 second.len(),
636 );
637 let server = tokio::spawn(async move {
638 let (mut stream, _) = listener.accept().await.unwrap();
639 read_head(&mut stream).await;
640 stream.write_all(response.as_bytes()).await.unwrap();
641 });
642 let http = crate::tls::reqwest_client_builder().build().unwrap();
643 let result = collect_models_document(
644 reqwest::Url::parse(&endpoint).unwrap(),
645 None,
646 ModelsFetchLimits {
647 bytes: limit,
648 ..MODELS_FETCH_LIMITS
649 },
650 |url| {
651 let http = http.clone();
652 async move {
653 http.get(url)
654 .send()
655 .await
656 .map_err(|_| CatalogRefreshError::Network.into())
657 }
658 },
659 )
660 .await;
661 if limit < body.len() {
662 assert_eq!(
663 result.unwrap_err().into_catalog(),
664 CatalogRefreshError::InvalidResponse
665 );
666 } else {
667 let (collected, _) = result.expect("same valid chunked body fits exact byte budget");
668 assert_eq!(
669 parse_models_response(&collected).unwrap()[0].id,
670 "chunked-model"
671 );
672 }
673 server.await.unwrap();
674 }
675 }
676
677 #[tokio::test]
678 async fn traversal_deadline_bounds_pending_and_completed_later_pages() {
679 let server = MockServer::start().await;
680 mount_models_json(
681 &server,
682 200,
683 json!({"data":[{"id":"first"}], "has_more":true, "last_id":"next"}),
684 )
685 .await;
686 // Fetch the first fixture response before starting the short clock so host
687 // scheduling/connection setup cannot make this accidentally a page-one test.
688 let mut first = Some(
689 crate::tls::reqwest_client_builder()
690 .build()
691 .unwrap()
692 .get(format!("{}/v1/models", server.uri()))
693 .send()
694 .await
695 .unwrap(),
696 );
697 let mut calls = 0;
698 let result = collect_models_document(
699 reqwest::Url::parse(&format!("{}/v1/models", server.uri())).unwrap(),
700 Some("after_id"),
701 ModelsFetchLimits {
702 timeout: Duration::from_millis(100),
703 ..MODELS_FETCH_LIMITS
704 },
705 |_| {
706 calls += 1;
707 let response = first.take();
708 async move {
709 match response {
710 Some(response) => Ok(response),
711 None => std::future::pending().await,
712 }
713 }
714 },
715 )
716 .await;
717 assert_eq!(
718 result.unwrap_err().into_catalog(),
719 CatalogRefreshError::Network
720 );
721 assert_eq!(
722 calls, 2,
723 "the pending later page must be covered by the collector timeout"
724 );
725
726 // Immediate response bodies avoid a scheduler-dependent network margin.
727 // Each completed fetch takes less than the short budget; together they
728 // exceed it. The generous control proves the same pages are otherwise valid.
729 for timeout in [Duration::from_millis(400), Duration::from_secs(5)] {
730 let mut responses = [
731 r#"{"data":[{"id":"first"}],"has_more":true,"last_id":"next"}"#,
732 r#"{"data":[{"id":"second"}],"has_more":false}"#,
733 ]
734 .into_iter()
735 .map(|body| {
736 reqwest::Response::from(
737 axum::http::Response::builder()
738 .status(200)
739 .body(body)
740 .unwrap(),
741 )
742 });
743 let mut calls = 0;
744 let result = collect_models_document(
745 reqwest::Url::parse("http://127.0.0.1/v1/models").unwrap(),
746 Some("after_id"),
747 ModelsFetchLimits {
748 timeout,
749 ..MODELS_FETCH_LIMITS
750 },
751 |_| {
752 calls += 1;
753 let response = responses.next().expect("exactly two fixture pages");
754 async move {
755 // Deliberately ready after bounded synchronous work, like
756 // parsing, so correctness cannot rely only on timer polling.
757 std::thread::sleep(Duration::from_millis(250));
758 Ok(response)
759 }
760 },
761 )
762 .await;
763 assert_eq!(calls, 2);
764 if timeout == Duration::from_millis(400) {
765 assert_eq!(
766 result.unwrap_err().into_catalog(),
767 CatalogRefreshError::Network
768 );
769 } else {
770 let (body, _) = result.expect("complete pages fit the larger shared budget");
771 assert_eq!(parse_models_response(&body).unwrap().len(), 2);
772 }
773 }
774 }
775
776 #[tokio::test]
777 async fn raw_rows_keep_duplicate_known_fields_invalid_in_existing_parsers() {
778 for malformed in [
779 r#"{"id":"second","id":"replacement"}"#,
780 r#"{"id":"second","pricing":{"prompt":"0.000001","prompt":"0.000002"}}"#,
781 ] {
782 let server = MockServer::start().await;
783 mount_page(&server, None, page(json!([{"id":"first"}]), Some("next"))).await;
784 let later = format!(r#"{{"data":[{malformed}],"has_more":false}}"#);
785 mount_page(
786 &server,
787 Some("next"),
788 ResponseTemplate::new(200).set_body_raw(later, "application/json"),
789 )
790 .await;
791 let body = collect_fixture(&server, MODELS_FETCH_LIMITS).await.unwrap();
792 assert!(
793 body.contains(malformed),
794 "collector must retain original row fields"
795 );
796 // OpenRouter decodes per row (#6690): the ambiguous row is skipped as
797 // malformed, never accepted with a last-wins value, and the valid row
798 // survives.
799 let openrouter = parse_openrouter_models_response(&body).unwrap();
800 let ids: Vec<&str> = openrouter.iter().map(|item| item.id.as_str()).collect();
801 assert_eq!(ids, ["first"]);
802 assert_eq!(
803 parse_baseten_models_response(&body).unwrap_err(),
804 CatalogRefreshError::InvalidResponse
805 );
806 }
807 }
808
809 #[tokio::test]
810 async fn existing_provider_parsers_own_cross_page_duplicates_and_full_metadata() {
811 let server = MockServer::start().await;
812 mount_page(
813 &server,
814 None,
815 page(
816 json!([{
817 "id":"same/model", "context_length":32000,
818 "pricing":{"prompt":"0.000001", "completion":"0.000002"},
819 "codewhale":{"protocol":"anthropic-messages", "default":true}
820 }]),
821 Some("next"),
822 ),
823 )
824 .await;
825 mount_page(&server, Some("next"), page(json!([
826 {"id":"same/model", "context_length":99999, "pricing":{"prompt":"0.000009", "completion":"0.000009"}},
827 {"id":"later/model", "context_length":64000, "max_completion_tokens":8000,
828 "pricing":{"prompt":"0.000003", "completion":"0.000004", "input_cache_read":"0.0000003"},
829 "supported_features":["vision"], "supported_parameters":["tools", "reasoning"],
830 "architecture":{"input_modalities":["text","image"],"output_modalities":["text"]},
831 "codewhale":{"protocol":"chat-completions"}}
832 ]), None)).await;
833 let body = collect_fixture(&server, MODELS_FETCH_LIMITS).await.unwrap();
834 let openrouter = parse_openrouter_models_response(&body).unwrap();
835 assert_eq!(openrouter.len(), 2);
836 let first =
837 openrouter_to_catalog_offering(&openrouter[0], "openrouter", "fixture-fp", 42).unwrap();
838 assert_eq!(first.limit.unwrap().context, Some(32000));
839 assert_eq!(first.cost.unwrap().input, Some(1.0));
840 let later =
841 openrouter_to_catalog_offering(&openrouter[1], "openrouter", "fixture-fp", 42).unwrap();
842 assert_eq!(later.limit.unwrap().context, Some(64000));
843 assert_eq!(later.cost.unwrap().cache_read, Some(0.3));
844 assert_eq!(later.reasoning, Some(true));
845 assert_eq!(later.tool_call, Some(true));
846 assert_eq!(later.modalities.unwrap().input, ["text", "image"]);
847 assert_eq!(
848 parse_baseten_models_response(&body).unwrap_err(),
849 CatalogRefreshError::InvalidResponse
850 );
851 let codewhale =
852 codewhale_catalog_offerings_from_body(&body, "codewhale", "fixture-fp", 42).unwrap();
853 assert_eq!(codewhale.len(), 2);
854 assert_eq!(codewhale[0].endpoint_key, "messages");
855 assert!(codewhale[0].default_for_provider);
856 assert_eq!(codewhale[1].endpoint_key, "chat");
857
858 server.reset().await;
859 mount_page(
860 &server,
861 None,
862 page(json!([{"id":" same/model "}]), Some("next")),
863 )
864 .await;
865 mount_page(
866 &server,
867 Some("next"),
868 page(json!([{"id":"same/model"}]), None),
869 )
870 .await;
871 let body = collect_fixture(&server, MODELS_FETCH_LIMITS).await.unwrap();
872 assert_eq!(
873 parse_baseten_models_response(&body).unwrap_err(),
874 CatalogRefreshError::InvalidResponse,
875 "Baseten trims before duplicate rejection across pages"
876 );
877
878 server.reset().await;
879 mount_page(
880 &server,
881 None,
882 page(json!([{"id":"first/model"}]), Some("next")),
883 )
884 .await;
885 mount_page(&server, Some("next"), page(json!([{"id":"later/model", "context_length":"64000", "max_completion_tokens":8000,
886 "pricing":{"prompt":"0.000003", "completion":"0.000004"}, "supported_features":["vision"]}]), None)).await;
887 let body = collect_fixture(&server, MODELS_FETCH_LIMITS).await.unwrap();
888 let rows = parse_baseten_models_response(&body).unwrap();
889 let later = baseten_to_catalog_offering(&rows[1], "base-ten", "fixture-fp", 42).unwrap();
890 assert_eq!(later.provider, "base-ten");
891 assert_eq!(later.limit.unwrap().output, Some(8000));
892 assert_eq!(later.cost.unwrap().output, Some(4.0));
893 assert_eq!(later.attachment, Some(true));
894 }
895
896 #[tokio::test]
897 async fn later_page_failure_preserves_complete_same_scope_cache_and_observation_time() {
898 let server = MockServer::start().await;
899 mount_models_json(
900 &server,
901 200,
902 json!({"data":[{"id":"old-first"},{"id":"old-second"}]}),
903 )
904 .await;
905 let client = anthropic_client(&server.uri());
906 let mut cache = ProviderCatalogCache::new();
907 let mut original = client.fetch_catalog_delta().await.unwrap();
908 original.fetched_at = 17;
909 cache.record_success(original, 3600);
910 let fingerprint = base_url_fingerprint(&server.uri());
911 let before = cache.get("anthropic", &fingerprint).unwrap().clone();
912 for (response, expected) in [
913 (
914 ResponseTemplate::new(401),
915 CatalogRefreshError::Unauthorized,
916 ),
917 (ResponseTemplate::new(429), CatalogRefreshError::RateLimited),
918 (ResponseTemplate::new(500), CatalogRefreshError::Network),
919 (
920 ResponseTemplate::new(200).set_body_string("broken-json"),
921 CatalogRefreshError::InvalidResponse,
922 ),
923 ] {
924 server.reset().await;
925 mount_page(
926 &server,
927 None,
928 page(json!([{"id":"new-partial-only"}]), Some("next")),
929 )
930 .await;
931 mount_page(&server, Some("next"), response).await;
932 assert_eq!(
933 client.refresh_catalog_cache(&mut cache, 3600).await,
934 CatalogStatus::Failed { reason: expected }
935 );
936 let retained = cache.get("anthropic", &fingerprint).unwrap();
937 assert_eq!(retained.offerings, before.offerings);
938 assert_eq!(retained.fetched_at, 17);
939 assert_eq!(retained.ttl_secs, before.ttl_secs);
940 assert_eq!(server.received_requests().await.unwrap().len(), 2);
941 }
942 }
943
944 fn assert_no_canaries(error: &anyhow::Error) {
945 let surfaced = format!("{error:#} {error:?} {:?}", crate::retry_status::snapshot());
946 for canary in [KEY, CURSOR, "cursor%2Fsecond", "custom-header-canary"] {
947 assert!(
948 !surfaced.contains(canary),
949 "model error or retry state exposed a secret/cursor"
950 );
951 }
952 }
953
954 /// #6173: a geo-blocked key produced `Invalid request (400): ` — the colon
955 /// that introduces the provider's reason, with nothing after it, because the
956 /// catalog path discarded the body wholesale. A geo-block, a bad key and a
957 /// wrong endpoint were then indistinguishable, and the reporter had to change
958 /// VPN exits to find out which one it was. The reason is the provider's own
959 /// words; only this client's secrets have to go.
960 #[tokio::test]
961 async fn catalog_errors_surface_the_provider_reason_without_client_secrets() {
962 const REASON: &str = "User location is not supported for the API use.";
963
964 let server = MockServer::start().await;
965 mount_page(
966 &server,
967 None,
968 ResponseTemplate::new(400).set_body_json(json!({
969 "error": {"code": 400, "message": REASON, "status": "FAILED_PRECONDITION"}
970 })),
971 )
972 .await;
973 let client = anthropic_client(&server.uri());
974 let error = client.list_models().await.unwrap_err();
975 assert!(
976 format!("{error:#}").contains(REASON),
977 "the provider's reason must reach the user: {error:#}"
978 );
979 assert_no_canaries(&error);
980
981 // The same reason, from an endpoint that also echoes back things only
982 // this client could have sent it. The reason survives; they do not.
983 let echoing = MockServer::start().await;
984 mount_page(
985 &echoing,
986 None,
987 ResponseTemplate::new(400).set_body_json(json!({
988 "error": {"message": format!("{REASON} key={KEY} header=custom-header-canary")}
989 })),
990 )
991 .await;
992 let client = anthropic_client(&echoing.uri());
993 crate::retry_status::clear();
994 let error = client.list_models().await.unwrap_err();
995 assert!(format!("{error:#}").contains(REASON), "{error:#}");
996 assert_no_canaries(&error);
997 crate::retry_status::clear();
998 }
999
1000 #[tokio::test]
1001 async fn later_page_http_and_transport_errors_do_not_expose_cursor_or_key() {
1002 for isolated in [false, true] {
1003 let server = MockServer::start().await;
1004 mount_page(
1005 &server,
1006 None,
1007 page(json!([{"id":"first-model"}]), Some(CURSOR)),
1008 )
1009 .await;
1010 mount_page(
1011 &server,
1012 Some(CURSOR),
1013 ResponseTemplate::new(500).set_body_json(json!({
1014 "error":{"message":format!("bad cursor {CURSOR}; key {KEY}; custom-header-canary")}
1015 })),
1016 )
1017 .await;
1018 let mut client = anthropic_client(&server.uri());
1019 client.isolated_request_state = isolated;
1020 client.retry.enabled = true;
1021 client.retry.max_retries = 1;
1022 client.retry.initial_delay = 0.001;
1023 client.retry.max_delay = 0.001;
1024 crate::retry_status::clear();
1025 let error = client.list_models().await.unwrap_err();
1026 assert_no_canaries(&error);
1027 if !isolated {
1028 assert!(crate::retry_status::snapshot().is_failed());
1029 }
1030 assert_eq!(server.received_requests().await.unwrap().len(), 3);
1031 crate::retry_status::clear();
1032
1033 server.reset().await;
1034 mount_page(
1035 &server,
1036 None,
1037 page(json!([{"id":"first-model"}]), Some(CURSOR)),
1038 )
1039 .await;
1040 mount_page(
1041 &server,
1042 Some(CURSOR),
1043 page(
1044 json!([{
1045 "id":"second-model", "created":format!("{CURSOR} {KEY} custom-header-canary")
1046 }]),
1047 None,
1048 ),
1049 )
1050 .await;
1051 let error = client.list_models().await.unwrap_err();
1052 assert_no_canaries(&error);
1053 assert!(error.to_string().contains("InvalidResponse"));
1054 assert_eq!(
1055 client.fetch_catalog_delta().await.unwrap_err(),
1056 CatalogRefreshError::InvalidResponse
1057 );
1058 assert_eq!(server.received_requests().await.unwrap().len(), 4);
1059 }
1060
1061 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1062 let base_url = format!("http://{}", listener.local_addr().unwrap());
1063 let heads = Arc::new(StdMutex::new(Vec::new()));
1064 let recorded = heads.clone();
1065 let server = tokio::spawn(async move {
1066 loop {
1067 let (mut stream, _) = listener.accept().await.unwrap();
1068 let head = read_head(&mut stream).await;
1069 let is_first = !head.lines().next().unwrap().contains("after_id=");
1070 recorded.lock().unwrap().push(head);
1071 if is_first {
1072 let body =
1073 json!({"data":[{"id":"first-model"}], "has_more":true, "last_id":CURSOR})
1074 .to_string();
1075 stream.write_all(format!("HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len()).as_bytes()).await.unwrap();
1076 }
1077 // Page two closes before response headers, producing a transport
1078 // error whose reqwest URL used to contain the opaque cursor.
1079 }
1080 });
1081 let client = anthropic_client(&base_url);
1082 let error = client.list_models().await.unwrap_err();
1083 assert_no_canaries(&error);
1084 assert!(
1085 heads
1086 .lock()
1087 .unwrap()
1088 .iter()
1089 .any(|head| head.contains("after_id="))
1090 );
1091 server.abort();
1092 let _ = server.await;
1093 crate::retry_status::clear();
1094 }
1095
1096 #[tokio::test]
1097 async fn unpaginated_custom_identities_keep_exact_identity_and_endpoint_ownership() {
1098 let first = MockServer::start().await;
1099 let second = MockServer::start().await;
1100 mount_models_json(&first, 200, json!({"data":[{"id":"first/model"}]})).await;
1101 mount_models_json(&second, 200, json!({"data":[{"id":"second/model"}]})).await;
1102 let mut cache = ProviderCatalogCache::new();
1103 for (identity, server, expected) in [
1104 ("base-ten", &first, "first/model"),
1105 ("Base-Ten", &first, "first/model"),
1106 ("base-ten", &second, "second/model"),
1107 ] {
1108 let client = custom_mock_client_for_identity(server, identity);
1109 let delta = client.fetch_catalog_delta().await.unwrap();
1110 assert_eq!(delta.provider, identity);
1111 assert_eq!(delta.offerings[0].provider, identity);
1112 assert_eq!(delta.offerings[0].wire_model_id, expected);
1113 assert_eq!(
1114 delta.base_url_fingerprint,
1115 base_url_fingerprint(&format!("{}/v1", server.uri()))
1116 );
1117 cache.record_success(delta, 3600);
1118 }
1119 assert_eq!(cache.entries.len(), 3);
1120 }
1121
1121 lines RUST