返回 CodeWhale
provider_catalog_live.rs
根目录 / crates / tui / src / provider_catalog_live.rs
1 //! Durable, secret-free per-provider `/models` catalog cache.
2 //!
3 //! This is the persistence owner for [`codewhale_config::catalog::ProviderCatalogCache`].
4 //! It replaces the process-only client refresh and provider_lake CLI cache writers: successful
5 //! refreshes replace one exact `(provider kind, identity, base URL fingerprint)`
6 //! partition, failures retain that partition's prior rows, and startup loads
7 //! only the active route's exact partition. Credentials authorize the fetch in
8 //! `client`; they never enter this module or its disk envelope. Baseten and
9 //! Codewhale account rosters are memory-only and cleared before each refresh;
10 //! a named custom route at either official endpoint follows the same rule.
11 //!
12 //! The former unscoped OpenRouter cache reader is retired; installed old files
13 //! are preserved and never imported as provider authority.
14 //! `catalog/provider-catalogs.json` is the sole writer-owned live roster store.
15 //! Older per-endpoint `provider-*.json` files also lack the built-in/custom
16 //! kind boundary and account-roster exclusion, so they are left untouched and
17 //! replaced only by a newly authenticated refresh into this store.
18
19 use std::collections::BTreeMap;
20 use std::fs::{self, OpenOptions};
21 use std::io::Read as _;
22 use std::path::{Path, PathBuf};
23 use std::sync::atomic::{AtomicBool, Ordering};
24 use std::sync::{LazyLock, RwLock};
25
26 use anyhow::{Context, Result};
27 use codewhale_config::catalog::now_unix;
28 use codewhale_config::catalog::{
29 CatalogRefreshError, CatalogSnapshot, CatalogStatus, ProviderCatalogCache,
30 ProviderCatalogDelta, base_url_fingerprint,
31 };
32 use codewhale_config::persistence::atomic_write_json;
33 use codewhale_config::pricing::{Currency, OfferingPricing, PricingProvenance};
34 use serde::{Deserialize, Serialize};
35
36 use crate::config::{Config, ProviderKind};
37
38 const CACHE_SCHEMA_VERSION: u32 = 2;
39 const CACHE_FILE: &str = "provider-catalogs.json";
40 const MAX_CACHE_BYTES: u64 = 32 * 1024 * 1024;
41 const MAX_CACHE_SCOPES: usize = 64;
42 const MAX_CACHE_ROWS: usize = 50_000;
43
44 #[derive(Debug, Clone, Copy)]
45 struct CachePersistenceLimits {
46 max_bytes: u64,
47 max_scopes: usize,
48 max_rows: usize,
49 }
50
51 const CACHE_PERSISTENCE_LIMITS: CachePersistenceLimits = CachePersistenceLimits {
52 max_bytes: MAX_CACHE_BYTES,
53 max_scopes: MAX_CACHE_SCOPES,
54 max_rows: MAX_CACHE_ROWS,
55 };
56
57 /// Provider-owned catalogs are refreshed daily. Past-TTL rows remain visible
58 /// with an explicit stale receipt until a successful replacement arrives.
59 pub const DEFAULT_PROVIDER_CATALOG_TTL_SECS: u64 = 24 * 60 * 60;
60
61 static DISK_LOADED: AtomicBool = AtomicBool::new(false);
62
63 static CACHE: LazyLock<RwLock<ProviderCatalogCache>> =
64 LazyLock::new(|| RwLock::new(ProviderCatalogCache::new()));
65 static REFRESH_GENERATIONS: LazyLock<RwLock<BTreeMap<String, u64>>> =
66 LazyLock::new(|| RwLock::new(BTreeMap::new()));
67
68 #[derive(Debug, Clone)]
69 pub struct ProviderCatalogRefreshTicket {
70 provider: String,
71 provider_kind: ProviderKind,
72 fingerprint: Option<String>,
73 generation: u64,
74 }
75
76 /// Immutable, secret-free catalog rate evidence captured at dispatch.
77 ///
78 /// Rates are stored as canonical decimal strings rather than `f64` so route
79 /// receipts retain exact equality and stable JSON. `catalog_revision` binds
80 /// every identity, scope, timestamp, currency, provenance, and rate field; it
81 /// therefore changes even when two refreshes land in the same Unix second.
82 #[derive(Debug, Clone, PartialEq, Eq)]
83 pub struct ProviderLivePricingQuote {
84 pub(crate) provider: ProviderKind,
85 pub(crate) provider_identity: String,
86 pub(crate) wire_model: String,
87 pub(crate) endpoint_fingerprint: String,
88 pub(crate) catalog_fetched_at: u64,
89 pub(crate) catalog_revision: String,
90 pub(crate) currency: Currency,
91 pub(crate) provenance: PricingProvenance,
92 pub(crate) cloud_facts: Option<CloudFactsPricingSource>,
93 pub(crate) input_per_million: Option<String>,
94 pub(crate) output_per_million: Option<String>,
95 pub(crate) cache_read_per_million: Option<String>,
96 pub(crate) cache_write_per_million: Option<String>,
97 }
98
99 /// Signed source identity and validity, bound into a frozen rate receipt.
100 /// The base URL is admitted only through the canonical official-endpoint
101 /// contract; custom URLs or credential-bearing URLs cannot enter this field.
102 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
103 pub(crate) struct CloudFactsPricingSource {
104 pub(crate) facts_version: u64,
105 pub(crate) key_id: String,
106 pub(crate) valid_until: Option<u64>,
107 pub(crate) base_url: String,
108 }
109
110 #[derive(Serialize, Deserialize)]
111 struct ProviderLivePricingQuoteWire {
112 provider: String,
113 provider_identity: String,
114 wire_model: String,
115 endpoint_fingerprint: String,
116 catalog_fetched_at: u64,
117 catalog_revision: String,
118 currency: Currency,
119 provenance: PricingProvenance,
120 #[serde(default, skip_serializing_if = "Option::is_none")]
121 cloud_facts: Option<CloudFactsPricingSource>,
122 #[serde(default, skip_serializing_if = "Option::is_none")]
123 input_per_million: Option<String>,
124 #[serde(default, skip_serializing_if = "Option::is_none")]
125 output_per_million: Option<String>,
126 #[serde(default, skip_serializing_if = "Option::is_none")]
127 cache_read_per_million: Option<String>,
128 #[serde(default, skip_serializing_if = "Option::is_none")]
129 cache_write_per_million: Option<String>,
130 }
131
132 impl TryFrom<&ProviderLivePricingQuote> for ProviderLivePricingQuoteWire {
133 type Error = &'static str;
134 fn try_from(quote: &ProviderLivePricingQuote) -> Result<Self, Self::Error> {
135 let provider = codewhale_config::descriptors::tui_wire_tag_for_route(
136 quote.provider,
137 &quote.provider_identity,
138 )
139 .ok_or("contradictory pricing provider identity")?;
140 Ok(Self {
141 provider: provider.into(),
142 provider_identity: quote.provider_identity.clone(),
143 wire_model: quote.wire_model.clone(),
144 endpoint_fingerprint: quote.endpoint_fingerprint.clone(),
145 catalog_fetched_at: quote.catalog_fetched_at,
146 catalog_revision: quote.catalog_revision.clone(),
147 currency: quote.currency.clone(),
148 provenance: quote.provenance.clone(),
149 cloud_facts: quote.cloud_facts.clone(),
150 input_per_million: quote.input_per_million.clone(),
151 output_per_million: quote.output_per_million.clone(),
152 cache_read_per_million: quote.cache_read_per_million.clone(),
153 cache_write_per_million: quote.cache_write_per_million.clone(),
154 })
155 }
156 }
157
158 impl Serialize for ProviderLivePricingQuote {
159 fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
160 where
161 S: serde::Serializer,
162 {
163 if !self.is_structurally_valid() {
164 return serializer.serialize_none();
165 }
166 ProviderLivePricingQuoteWire::try_from(self)
167 .map_err(serde::ser::Error::custom)?
168 .serialize(serializer)
169 }
170 }
171
172 impl<'de> Deserialize<'de> for ProviderLivePricingQuote {
173 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
174 where
175 D: serde::Deserializer<'de>,
176 {
177 let wire = ProviderLivePricingQuoteWire::deserialize(deserializer)?;
178 let provider = codewhale_config::descriptors::kind_from_tui_wire_tag(
179 &wire.provider,
180 &wire.provider_identity,
181 )
182 .ok_or_else(|| serde::de::Error::custom("contradictory pricing provider identity"))?;
183 let quote = Self {
184 provider,
185 provider_identity: wire.provider_identity,
186 wire_model: wire.wire_model,
187 endpoint_fingerprint: wire.endpoint_fingerprint,
188 catalog_fetched_at: wire.catalog_fetched_at,
189 catalog_revision: wire.catalog_revision,
190 currency: wire.currency,
191 provenance: wire.provenance,
192 cloud_facts: wire.cloud_facts,
193 input_per_million: wire.input_per_million,
194 output_per_million: wire.output_per_million,
195 cache_read_per_million: wire.cache_read_per_million,
196 cache_write_per_million: wire.cache_write_per_million,
197 };
198 quote
199 .is_structurally_valid()
200 .then_some(quote)
201 .ok_or_else(|| serde::de::Error::custom("invalid provider-live pricing quote"))
202 }
203 }
204
205 pub(crate) fn deserialize_optional_provider_live_pricing<'de, D>(
206 deserializer: D,
207 ) -> std::result::Result<Option<ProviderLivePricingQuote>, D::Error>
208 where
209 D: serde::Deserializer<'de>,
210 {
211 let value = Option::<serde_json::Value>::deserialize(deserializer)?;
212 Ok(value.and_then(|value| serde_json::from_value(value).ok()))
213 }
214
215 impl ProviderLivePricingQuote {
216 /// Whether the quote names at least one rate. A `[[custom_models]]` row
217 /// declared only to add a model to a roster freezes a rate-less quote.
218 #[must_use]
219 pub(crate) fn carries_rates(&self) -> bool {
220 [
221 &self.input_per_million,
222 &self.output_per_million,
223 &self.cache_read_per_million,
224 &self.cache_write_per_million,
225 ]
226 .into_iter()
227 .any(Option::is_some)
228 }
229
230 fn is_structurally_valid(&self) -> bool {
231 self.pricing_for_route(
232 self.provider,
233 &self.provider_identity,
234 &self.wire_model,
235 &self.endpoint_fingerprint,
236 self.catalog_fetched_at,
237 )
238 .is_some()
239 }
240 fn canonical_rate(rate: Option<f64>) -> Option<String> {
241 rate.map(|rate| rate.to_string())
242 }
243
244 fn revision_for(
245 provider: ProviderKind,
246 provider_identity: &str,
247 wire_model: &str,
248 endpoint_fingerprint: &str,
249 catalog_fetched_at: u64,
250 currency: &Currency,
251 provenance: &PricingProvenance,
252 input_per_million: &Option<String>,
253 output_per_million: &Option<String>,
254 cache_read_per_million: &Option<String>,
255 cache_write_per_million: &Option<String>,
256 cloud_facts: Option<&CloudFactsPricingSource>,
257 ) -> Option<String> {
258 let payload = serde_json::to_vec(&(
259 "codewhale-provider-live-pricing-quote-v1",
260 codewhale_config::descriptors::tui_wire_tag_for_route(provider, provider_identity)?,
261 provider_identity,
262 wire_model,
263 endpoint_fingerprint,
264 catalog_fetched_at,
265 currency,
266 provenance,
267 input_per_million,
268 output_per_million,
269 cache_read_per_million,
270 cache_write_per_million,
271 ))
272 .ok()?;
273 // Preserve the existing provider-live wire revision. The additional
274 // cloud source is independently domain-separated and hashes the whole
275 // original binding as well as the signed version/key/expiry.
276 let payload = match cloud_facts {
277 Some(source) => {
278 serde_json::to_vec(&("codewhale-cloud-facts-pricing-quote-v1", payload, source))
279 .ok()?
280 }
281 None => payload,
282 };
283 Some(format!("sha256:{}", crate::hashing::sha256_hex(payload)))
284 }
285
286 fn from_pricing(
287 provider: ProviderKind,
288 provider_identity: &str,
289 wire_model: &str,
290 endpoint_fingerprint: &str,
291 catalog_fetched_at: u64,
292 pricing: &OfferingPricing,
293 ) -> Option<Self> {
294 let provider_identity = provider_identity.trim();
295 let wire_model = wire_model.trim();
296 if crate::cost_status::sanitize_persisted_route_label(provider_identity)
297 != provider_identity
298 || crate::cost_status::sanitize_persisted_route_label(wire_model) != wire_model
299 {
300 return None;
301 }
302 let input_per_million = Self::canonical_rate(pricing.input_per_million);
303 let output_per_million = Self::canonical_rate(pricing.output_per_million);
304 let cache_read_per_million = Self::canonical_rate(pricing.cache_read_per_million);
305 let cache_write_per_million = Self::canonical_rate(pricing.cache_write_per_million);
306 let catalog_revision = Self::revision_for(
307 provider,
308 provider_identity,
309 wire_model,
310 endpoint_fingerprint,
311 catalog_fetched_at,
312 &pricing.currency,
313 &pricing.provenance,
314 &input_per_million,
315 &output_per_million,
316 &cache_read_per_million,
317 &cache_write_per_million,
318 None,
319 )?;
320 Some(Self {
321 provider,
322 provider_identity: provider_identity.to_string(),
323 wire_model: wire_model.to_string(),
324 endpoint_fingerprint: endpoint_fingerprint.to_string(),
325 catalog_fetched_at,
326 catalog_revision,
327 currency: pricing.currency.clone(),
328 provenance: pricing.provenance.clone(),
329 cloud_facts: None,
330 input_per_million,
331 output_per_million,
332 cache_read_per_million,
333 cache_write_per_million,
334 })
335 }
336
337 fn parse_rate(rate: &Option<String>) -> Option<Option<f64>> {
338 let Some(rate) = rate else {
339 return Some(None);
340 };
341 let parsed = rate.parse::<f64>().ok()?;
342 (parsed.is_finite() && parsed >= 0.0 && parsed.to_string() == *rate).then_some(Some(parsed))
343 }
344
345 /// Rehydrate the frozen row only when every receipt binding is intact.
346 /// This is deliberately cache-free: a refresh after dispatch cannot alter
347 /// an earlier turn, while malformed or legacy receipts fail closed.
348 pub(crate) fn pricing_for_route(
349 &self,
350 provider: ProviderKind,
351 provider_identity: &str,
352 wire_model: &str,
353 endpoint_fingerprint: &str,
354 dispatched_at_unix: u64,
355 ) -> Option<OfferingPricing> {
356 let provider_identity = provider_identity.trim();
357 let wire_model = wire_model.trim();
358 if self.provider != provider
359 || crate::cost_status::sanitize_persisted_route_label(&self.provider_identity)
360 != self.provider_identity
361 || crate::cost_status::sanitize_persisted_route_label(&self.wire_model)
362 != self.wire_model
363 || self.endpoint_fingerprint.len() != 64
364 || !self
365 .endpoint_fingerprint
366 .bytes()
367 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
368 || self.provider_identity != provider_identity
369 || self.wire_model != wire_model
370 || self.endpoint_fingerprint != endpoint_fingerprint
371 || self.catalog_fetched_at > dispatched_at_unix
372 || self.currency != Currency::Usd
373 {
374 return None;
375 }
376 match (&self.provenance, &self.cloud_facts) {
377 (PricingProvenance::UserOverride, None) => {}
378 (PricingProvenance::ProviderLive, None)
379 if dispatched_at_unix.saturating_sub(self.catalog_fetched_at)
380 < DEFAULT_PROVIDER_CATALOG_TTL_SECS
381 && reviewed_provider_live_scope(
382 provider,
383 provider_identity,
384 endpoint_fingerprint,
385 ) => {}
386 (PricingProvenance::CloudFacts, Some(source))
387 if source.facts_version > 0
388 && !source.key_id.is_empty()
389 && source.key_id.len() <= 128
390 && source
391 .key_id
392 .bytes()
393 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_'))
394 && source
395 .valid_until
396 .is_none_or(|expires| dispatched_at_unix <= expires)
397 && cloud_pricing_scope(provider, provider_identity, &source.base_url)
398 && base_url_fingerprint(&source.base_url) == endpoint_fingerprint => {}
399 _ => return None,
400 }
401 let input_per_million = Self::parse_rate(&self.input_per_million)?;
402 let output_per_million = Self::parse_rate(&self.output_per_million)?;
403 let cache_read_per_million = Self::parse_rate(&self.cache_read_per_million)?;
404 let cache_write_per_million = Self::parse_rate(&self.cache_write_per_million)?;
405 let cost = codewhale_config::models_dev::ModelsDevCost {
406 input: input_per_million,
407 output: output_per_million,
408 cache_read: cache_read_per_million,
409 cache_write: cache_write_per_million,
410 };
411 if !codewhale_config::pricing::catalog_cost_is_valid(&cost) {
412 return None;
413 }
414 // A reviewed per-token route needs both ordinary request classes. Cache
415 // classes remain optional and fail closed later if a turn used them.
416 if self.provenance == PricingProvenance::ProviderLive
417 && (cost.input.is_none() || cost.output.is_none())
418 {
419 return None;
420 }
421 if cost.input.is_none()
422 && cost.output.is_none()
423 && cost.cache_read.is_none()
424 && cost.cache_write.is_none()
425 && self.provenance != PricingProvenance::UserOverride
426 {
427 return None;
428 }
429 let expected_revision = Self::revision_for(
430 self.provider,
431 &self.provider_identity,
432 &self.wire_model,
433 &self.endpoint_fingerprint,
434 self.catalog_fetched_at,
435 &self.currency,
436 &self.provenance,
437 &self.input_per_million,
438 &self.output_per_million,
439 &self.cache_read_per_million,
440 &self.cache_write_per_million,
441 self.cloud_facts.as_ref(),
442 )?;
443 if self.catalog_revision != expected_revision {
444 return None;
445 }
446 Some(OfferingPricing {
447 provider: self.provider_identity.clone(),
448 wire_model_id: self.wire_model.clone(),
449 canonical_model: None,
450 currency: self.currency.clone(),
451 input_per_million: cost.input,
452 output_per_million: cost.output,
453 cache_read_per_million: cost.cache_read,
454 cache_write_per_million: cost.cache_write,
455 provenance: self.provenance.clone(),
456 effective_at: Some(self.catalog_fetched_at),
457 endpoint_fingerprint: Some(self.endpoint_fingerprint.clone()),
458 })
459 }
460 }
461
462 #[derive(Debug, Clone, Serialize, Deserialize)]
463 struct PersistedProviderCatalogs {
464 schema_version: u32,
465 cache: ProviderCatalogCache,
466 }
467
468 #[derive(Serialize)]
469 struct PersistedProviderCatalogsRef<'a> {
470 schema_version: u32,
471 cache: &'a ProviderCatalogCache,
472 }
473
474 /// Resolve the cache under Codewhale's catalog state directory.
475 ///
476 /// Unguarded tests are confined to the TUI test root, matching the Models.dev
477 /// cache contract, so they never inspect a developer's real provider catalog.
478 #[must_use]
479 pub fn cache_path() -> Option<PathBuf> {
480 #[cfg(test)]
481 {
482 if !crate::test_support::guarded_environment_provides_state_paths() {
483 return Some(
484 crate::test_support::unsealed_test_state_root()
485 .join("catalog")
486 .join(CACHE_FILE),
487 );
488 }
489 }
490 codewhale_config::resolve_state_dir("catalog")
491 .ok()
492 .map(|dir| dir.join(CACHE_FILE))
493 }
494
495 fn canonical_provider_scope(provider: &str) -> String {
496 // Despite the historical name, this is the exact configured ownership
497 // scope. Never collapse a custom table that happens to resemble a built-in
498 // or setup-template alias.
499 crate::provider_lake::catalog_partition_key(provider)
500 }
501
502 #[cfg(test)]
503 fn inferred_provider_kind(identity: &str) -> ProviderKind {
504 // No recognized built-in spelling resolves to a compatible-template id,
505 // so the parse fallback below already answers Custom for every named
506 // custom table (#6289).
507 ProviderKind::parse(identity).unwrap_or(ProviderKind::Custom)
508 }
509
510 fn storage_provider(kind: ProviderKind, identity: &str) -> String {
511 format!("{}:{}", kind.as_str(), identity.trim())
512 }
513
514 /// Providers whose model list is owned by their own `/v1/models` roster
515 /// rather than the cross-provider Models.dev snapshot: the named live
516 /// gateways, plus custom hosts whose private roster no snapshot can serve
517 /// (#6289 widened). The active-provider refresh and the picker's freshness
518 /// receipt both gate on this one predicate, so they cannot drift apart.
519 pub(crate) fn provider_owns_live_catalog(provider: ProviderKind) -> bool {
520 matches!(
521 provider,
522 ProviderKind::Openrouter
523 | ProviderKind::Telecomjs
524 | ProviderKind::Edenai
525 | ProviderKind::Zenmux
526 | ProviderKind::Concentrate
527 | ProviderKind::Codewhale
528 | ProviderKind::Ollama
529 ) || provider == ProviderKind::Custom
530 }
531
532 /// Whether a catalog scope holds an account-scoped roster that must never be
533 /// shared across credentials (#6289).
534 ///
535 /// Baseten's `/models` answers per workspace, so its rows are fenced by
536 /// endpoint fingerprint — never by table name. The Codewhale API's own rows
537 /// are fenced the same way.
538 fn is_account_scoped_scope(provider: &str, fingerprint: &str) -> bool {
539 provider.starts_with("codewhale:")
540 || fingerprint == base_url_fingerprint(codewhale_config::catalog::BASETEN_BASE_URL)
541 || fingerprint
542 == base_url_fingerprint(ProviderKind::Codewhale.provider().default_base_url())
543 }
544
545 fn cache_lock_path(path: &Path) -> PathBuf {
546 let mut name = path
547 .file_name()
548 .map(|name| name.to_os_string())
549 .unwrap_or_else(|| CACHE_FILE.into());
550 name.push(".lock");
551 path.with_file_name(name)
552 }
553
554 fn open_cache_lock(path: &Path) -> Result<fs::File> {
555 let parent = path
556 .parent()
557 .context("provider catalog lock path has no parent")?;
558 fs::create_dir_all(parent)
559 .with_context(|| format!("create provider catalog directory {}", parent.display()))?;
560 let mut options = OpenOptions::new();
561 options.read(true).write(true).create(true).truncate(false);
562 #[cfg(unix)]
563 {
564 use std::os::unix::fs::OpenOptionsExt as _;
565 options
566 .mode(0o600)
567 .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK);
568 }
569 #[cfg(windows)]
570 {
571 use std::os::windows::fs::OpenOptionsExt as _;
572 options.custom_flags(0x0020_0000); // FILE_FLAG_OPEN_REPARSE_POINT
573 }
574 let file = options
575 .open(path)
576 .with_context(|| format!("open provider catalog lock {}", path.display()))?;
577 let metadata = file
578 .metadata()
579 .with_context(|| format!("inspect provider catalog lock {}", path.display()))?;
580 anyhow::ensure!(
581 metadata.is_file(),
582 "provider catalog lock {} must be a regular file",
583 path.display()
584 );
585 #[cfg(unix)]
586 {
587 use std::os::unix::fs::MetadataExt as _;
588 anyhow::ensure!(
589 metadata.nlink() == 1,
590 "provider catalog lock {} must not be hard linked",
591 path.display()
592 );
593 }
594 #[cfg(windows)]
595 {
596 use std::os::windows::fs::MetadataExt as _;
597 const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x0000_0400;
598 anyhow::ensure!(
599 metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT == 0,
600 "provider catalog lock {} must not be a reparse point",
601 path.display()
602 );
603 }
604 Ok(file)
605 }
606
607 fn load_from_disk_unlocked_with_limit(path: &Path, max_bytes: u64) -> Option<ProviderCatalogCache> {
608 let mut options = OpenOptions::new();
609 options.read(true);
610 #[cfg(unix)]
611 {
612 use std::os::unix::fs::OpenOptionsExt as _;
613 options.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK);
614 }
615 #[cfg(windows)]
616 {
617 use std::os::windows::fs::OpenOptionsExt as _;
618 options.custom_flags(0x0020_0000);
619 }
620 let file = options.open(path).ok()?;
621 let metadata = file.metadata().ok()?;
622 if !metadata.is_file() {
623 return None;
624 }
625 #[cfg(unix)]
626 {
627 use std::os::unix::fs::MetadataExt as _;
628 if metadata.nlink() != 1 {
629 return None;
630 }
631 }
632 #[cfg(windows)]
633 {
634 use std::os::windows::fs::MetadataExt as _;
635 if metadata.file_attributes() & 0x0000_0400 != 0 {
636 return None;
637 }
638 }
639 if metadata.len() > max_bytes {
640 tracing::debug!(
641 target: "provider_catalog",
642 path = %path.display(),
643 max_bytes,
644 "provider catalog cache exceeds read limit"
645 );
646 return None;
647 }
648 // Re-check through `take`: the file can grow after metadata is sampled.
649 let mut body = Vec::new();
650 file.take(max_bytes.saturating_add(1))
651 .read_to_end(&mut body)
652 .ok()?;
653 if body.len() as u64 > max_bytes {
654 return None;
655 }
656 let persisted: PersistedProviderCatalogs = serde_json::from_slice(&body).ok()?;
657 if persisted.schema_version != CACHE_SCHEMA_VERSION {
658 return None;
659 }
660 let mut cache = persisted.cache;
661 if cache.entries.len() > MAX_CACHE_SCOPES || cached_row_count(&cache) > MAX_CACHE_ROWS {
662 return None;
663 }
664 if !cache.entries.iter().all(|(key, entry)| {
665 let Some((kind, identity)) = entry.provider.split_once(':') else { return false; };
666 ProviderKind::parse(kind).is_some_and(|parsed| parsed.as_str() == kind)
667 && !identity.is_empty()
668 && key == &ProviderCatalogCache::cache_key(&entry.provider, &entry.base_url_fingerprint)
669 && entry.offerings.iter().all(|row| {
670 row.provider == identity
671 && crate::provider_lake::valid_catalog_model_id(&row.wire_model_id)
672 && provider_cost_source_allowed(row)
673 && matches!(&row.source, codewhale_config::catalog::CatalogSource::Live {
674 base_url_fingerprint, fetched_at
675 } if base_url_fingerprint == &entry.base_url_fingerprint && *fetched_at == entry.fetched_at)
676 })
677 }) { return None; }
678 // Older builds could durably cache account-scoped rosters. Scrub
679 // those entries on every load so upgrading cannot attach one workspace's
680 // catalog to a different credential.
681 cache
682 .entries
683 .retain(|_, entry| !is_account_scoped_scope(&entry.provider, &entry.base_url_fingerprint));
684 Some(cache)
685 }
686
687 fn provider_cost_source_allowed(row: &codewhale_config::catalog::CatalogOffering) -> bool {
688 use codewhale_config::catalog::CatalogSource;
689 matches!(
690 row.cost_source,
691 None | Some(
692 CatalogSource::Bundled
693 | CatalogSource::CodewhaleBundled { .. }
694 | CatalogSource::ModelsDevLive { .. }
695 )
696 )
697 }
698
699 fn load_from_disk_unlocked(path: &Path) -> Option<ProviderCatalogCache> {
700 load_from_disk_unlocked_with_limit(path, MAX_CACHE_BYTES)
701 }
702
703 fn load_from_disk() -> Option<ProviderCatalogCache> {
704 let path = cache_path()?;
705 if !path.is_file() {
706 return None;
707 }
708 let lock_file = open_cache_lock(&cache_lock_path(&path)).ok()?;
709 let lock = fd_lock::RwLock::new(lock_file);
710 let _guard = lock.read().ok()?;
711 load_from_disk_unlocked(&path)
712 }
713
714 fn ensure_cache_loaded() -> Result<()> {
715 if DISK_LOADED.load(Ordering::Acquire) {
716 return Ok(());
717 }
718 let mut cache = CACHE
719 .write()
720 .map_err(|_| anyhow::anyhow!("catalog cache unavailable"))?;
721 if DISK_LOADED.load(Ordering::Acquire) {
722 return Ok(());
723 }
724 if let Some(path) = cache_path() {
725 match fs::symlink_metadata(&path) {
726 Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
727 Err(err) => return Err(err.into()),
728 Ok(_) => {
729 let loaded = load_from_disk().context("invalid provider catalog cache")?;
730 for (key, entry) in loaded.entries {
731 if cache
732 .entries
733 .get(&key)
734 .is_none_or(|local| entry.fetched_at >= local.fetched_at)
735 {
736 cache.entries.insert(key, entry);
737 }
738 }
739 }
740 }
741 }
742 DISK_LOADED.store(true, Ordering::Release);
743 Ok(())
744 }
745
746 /// Read one exact cached route without creating a client or fetching credentials.
747 pub(crate) fn cached_entry_for_route(
748 kind: ProviderKind,
749 identity: &str,
750 base_url: &str,
751 ) -> Result<Option<codewhale_config::catalog::CachedProviderCatalog>> {
752 ensure_cache_loaded()?;
753 let cache = CACHE
754 .read()
755 .map_err(|_| anyhow::anyhow!("catalog cache unavailable"))?;
756 Ok(cache
757 .get(
758 &storage_provider(kind, identity),
759 &base_url_fingerprint(base_url),
760 )
761 .cloned())
762 }
763
764 /// Whether a saved `(provider, model)` pin is absent from that exact route's
765 /// FRESH live roster (#6035). `None` when no fresh roster exists: a stale,
766 /// failed, or absent roster cannot prove drift, and bundled catalog rows say
767 /// nothing about what the account serves today. Absence is a warning, never a
768 /// reason to rewrite the pin: the id may still answer (soft deprecation) and
769 /// other hosts may serve it on their own routes.
770 pub(crate) fn pin_missing_from_fresh_roster(
771 config: &Config,
772 provider: &str,
773 model: &str,
774 ) -> Option<bool> {
775 let captured = config.resolve_provider_pin_identity(provider).ok()?;
776 let kind = captured.provider;
777 let identity = captured.key.as_str();
778 let base_url = config.base_url_for_route(&captured);
779 // `status_for_route` reads memory only. A fresh process (doctor, a
780 // just-started TUI) must see the roster an earlier process persisted.
781 ensure_cache_loaded().ok()?;
782 if status_for_route(kind, identity, &base_url) != CatalogStatus::Fresh {
783 return None;
784 }
785 let listed = cached_entry_for_route(kind, identity, &base_url)
786 .ok()
787 .flatten()
788 .is_some_and(|entry| {
789 entry.offerings.iter().any(|offering| {
790 offering.wire_model_id == model
791 || offering.canonical_model.as_deref() == Some(model)
792 })
793 });
794 Some(!listed)
795 }
796
797 fn merge_durable_scope(
798 mut durable_cache: ProviderCatalogCache,
799 process_cache: &ProviderCatalogCache,
800 provider: &str,
801 fingerprint: &str,
802 ) -> ProviderCatalogCache {
803 durable_cache
804 .entries
805 .retain(|_, entry| !is_account_scoped_scope(&entry.provider, &entry.base_url_fingerprint));
806 if !is_account_scoped_scope(provider, fingerprint)
807 && let Some(entry) = process_cache.get(provider, fingerprint).cloned()
808 {
809 durable_cache.entries.insert(
810 ProviderCatalogCache::cache_key(provider, fingerprint),
811 entry,
812 );
813 }
814 durable_cache
815 }
816
817 fn persisted_envelope_len(cache: &ProviderCatalogCache) -> Result<u64> {
818 struct Counter(u64);
819 impl std::io::Write for Counter {
820 fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
821 self.0 = self.0.saturating_add(bytes.len() as u64);
822 Ok(bytes.len())
823 }
824 fn flush(&mut self) -> std::io::Result<()> {
825 Ok(())
826 }
827 }
828 let envelope = PersistedProviderCatalogsRef {
829 schema_version: CACHE_SCHEMA_VERSION,
830 cache,
831 };
832 let mut counter = Counter(0);
833 serde_json::to_writer_pretty(&mut counter, &envelope)
834 .context("measure provider catalog cache for bounded persistence")?;
835 Ok(counter.0.saturating_add(1))
836 }
837
838 fn cached_row_count(cache: &ProviderCatalogCache) -> usize {
839 cache.entries.values().fold(0usize, |total, entry| {
840 total.saturating_add(entry.offerings.len())
841 })
842 }
843
844 /// Compact a durable cache without ever truncating one provider roster.
845 ///
846 /// The exact scope being written is protected: if that scope alone fits, older
847 /// failed/stale scopes are evicted whole until the envelope is bounded. If the
848 /// protected scope alone does not fit, persistence is refused and the prior
849 /// atomic file remains intact. This avoids both self-bricking the 32 MiB read
850 /// limit and turning a partial provider roster into false authoritative truth.
851 fn bounded_cache_for_persistence(
852 mut cache: ProviderCatalogCache,
853 protected_scope: Option<(&str, &str)>,
854 now: u64,
855 limits: CachePersistenceLimits,
856 ) -> Result<ProviderCatalogCache> {
857 cache
858 .entries
859 .retain(|_, entry| !is_account_scoped_scope(&entry.provider, &entry.base_url_fingerprint));
860
861 let protected_key = protected_scope
862 .filter(|(provider, fingerprint)| !is_account_scoped_scope(provider, fingerprint))
863 .map(|(provider, fingerprint)| ProviderCatalogCache::cache_key(provider, fingerprint));
864
865 if let Some(key) = protected_key.as_deref()
866 && let Some(entry) = cache.entries.get(key).cloned()
867 {
868 let mut protected_only = ProviderCatalogCache::new();
869 protected_only.entries.insert(key.to_string(), entry);
870 anyhow::ensure!(
871 protected_only.entries.len() <= limits.max_scopes.min(MAX_CACHE_SCOPES)
872 && cached_row_count(&protected_only) <= limits.max_rows
873 && persisted_envelope_len(&protected_only)? <= limits.max_bytes,
874 "provider catalog scope {key:?} exceeds bounded persistence limits"
875 );
876 }
877
878 // Rank once while the cache/file locks are held. An older implementation
879 // reserialized and rescanned the entire envelope for every eviction, which
880 // made a valid sub-32-MiB file with many tiny scopes quadratic to compact.
881 let mut eviction_keys = cache
882 .entries
883 .iter()
884 .filter(|(key, _)| protected_key.as_deref() != Some(key.as_str()))
885 .map(|(key, entry)| {
886 let health_rank = if matches!(entry.status, CatalogStatus::Failed { .. }) {
887 0u8
888 } else if entry.is_stale(now) || matches!(entry.status, CatalogStatus::Stale { .. }) {
889 1u8
890 } else {
891 2u8
892 };
893 (health_rank, entry.fetched_at, key.clone())
894 })
895 .collect::<Vec<_>>();
896 eviction_keys.sort();
897 let eviction_keys = eviction_keys
898 .into_iter()
899 .map(|(_, _, key)| key)
900 .collect::<Vec<_>>();
901 let mut eviction_index = 0usize;
902 let mut rows = cached_row_count(&cache);
903 let max_scopes = limits.max_scopes.min(MAX_CACHE_SCOPES);
904
905 let mut evict_next = |cache: &mut ProviderCatalogCache| -> Result<usize> {
906 let key = eviction_keys
907 .get(eviction_index)
908 .context("provider catalog envelope cannot fit even after whole-scope compaction")?;
909 eviction_index = eviction_index.saturating_add(1);
910 let entry = cache
911 .entries
912 .remove(key)
913 .context("provider catalog eviction candidate disappeared")?;
914 Ok(entry.offerings.len())
915 };
916
917 // First enforce the cheap cardinality limits in bulk. Only after at most 64
918 // scopes remain do we serialize to enforce the exact on-disk byte limit.
919 while cache.entries.len() > max_scopes || rows > limits.max_rows {
920 rows = rows.saturating_sub(evict_next(&mut cache)?);
921 }
922 while persisted_envelope_len(&cache)? > limits.max_bytes {
923 let _removed_rows = evict_next(&mut cache)?;
924 }
925
926 Ok(cache)
927 }
928
929 fn write_bounded_cache(
930 path: &Path,
931 cache: ProviderCatalogCache,
932 protected_scope: Option<(&str, &str)>,
933 limits: CachePersistenceLimits,
934 ) -> Result<()> {
935 let cache = bounded_cache_for_persistence(cache, protected_scope, now_unix(), limits)?;
936 let envelope = PersistedProviderCatalogs {
937 schema_version: CACHE_SCHEMA_VERSION,
938 cache,
939 };
940 anyhow::ensure!(
941 persisted_envelope_len(&envelope.cache)? <= limits.max_bytes,
942 "bounded provider catalog cache exceeds its write limit"
943 );
944 atomic_write_json(path, &envelope)
945 }
946
947 fn persist_scope(cache: &ProviderCatalogCache, provider: &str, fingerprint: &str) -> bool {
948 let Some(path) = cache_path() else {
949 return false;
950 };
951 let provider = canonical_provider_scope(provider);
952 let result = (|| -> Result<()> {
953 let lock_file = open_cache_lock(&cache_lock_path(&path))?;
954 let mut lock = fd_lock::RwLock::new(lock_file);
955 let _guard = lock
956 .write()
957 .with_context(|| format!("write-lock provider catalog cache {}", path.display()))?;
958 // Merge only the exact scope this process just changed into the latest
959 // disk snapshot. A stale long-running TUI therefore cannot erase a
960 // different scope written by the Runtime API (or vice versa).
961 let durable_cache = merge_durable_scope(
962 if path.exists() {
963 load_from_disk_unlocked(&path).context("invalid prior catalog cache")?
964 } else {
965 ProviderCatalogCache::new()
966 },
967 cache,
968 &provider,
969 fingerprint,
970 );
971 write_bounded_cache(
972 &path,
973 durable_cache,
974 Some((&provider, fingerprint)),
975 CACHE_PERSISTENCE_LIMITS,
976 )
977 .with_context(|| format!("atomically write provider catalog {}", path.display()))
978 })();
979 if let Err(error) = result {
980 tracing::debug!(
981 target: "provider_catalog",
982 error = %error,
983 "provider catalog cache write failed"
984 );
985 return false;
986 }
987 true
988 }
989
990 /// Persist a failure without letting a stale process replace newer rows from
991 /// another Codewhale process for the same exact scope.
992 ///
993 /// The ordinary scoped merge is sufficient for successes because the response
994 /// being committed is the new roster. A failure is different: its process may
995 /// have started with an older last-known-good entry. Re-read the durable exact
996 /// scope while holding the cross-process write lock, prefer it when it is at
997 /// least as recent, then change only the status before writing. Thus a failed
998 /// refresh can preserve the newest roster without resurrecting its own stale
999 /// snapshot over another process's success.
1000 fn persist_failure_scope(
1001 cache: &mut ProviderCatalogCache,
1002 provider: &str,
1003 fingerprint: &str,
1004 reason: CatalogRefreshError,
1005 ) {
1006 let Some(path) = cache_path() else {
1007 return;
1008 };
1009 let provider = canonical_provider_scope(provider);
1010 let result = (|| -> Result<()> {
1011 let lock_file = open_cache_lock(&cache_lock_path(&path))?;
1012 let mut lock = fd_lock::RwLock::new(lock_file);
1013 let _guard = lock
1014 .write()
1015 .with_context(|| format!("write-lock provider catalog cache {}", path.display()))?;
1016 let durable_cache = if path.exists() {
1017 load_from_disk_unlocked(&path).context("invalid prior catalog cache")?
1018 } else {
1019 ProviderCatalogCache::new()
1020 };
1021
1022 if !is_account_scoped_scope(&provider, fingerprint)
1023 && let Some(durable_entry) = durable_cache.get(&provider, fingerprint).cloned()
1024 {
1025 let durable_is_newer = cache
1026 .get(&provider, fingerprint)
1027 .is_none_or(|local| durable_entry.fetched_at >= local.fetched_at);
1028 if durable_is_newer {
1029 cache.entries.insert(
1030 ProviderCatalogCache::cache_key(&provider, fingerprint),
1031 durable_entry,
1032 );
1033 cache.record_failure(&provider, fingerprint, reason);
1034 }
1035 }
1036
1037 let durable_cache = merge_durable_scope(durable_cache, cache, &provider, fingerprint);
1038 write_bounded_cache(
1039 &path,
1040 durable_cache,
1041 Some((&provider, fingerprint)),
1042 CACHE_PERSISTENCE_LIMITS,
1043 )
1044 .with_context(|| format!("atomically write provider catalog {}", path.display()))
1045 })();
1046 if let Err(error) = result {
1047 tracing::debug!(
1048 target: "provider_catalog",
1049 error = %error,
1050 "provider catalog failure receipt write failed"
1051 );
1052 }
1053 }
1054
1055 fn publish_exact_scope_for_identity(
1056 cache: &ProviderCatalogCache,
1057 provider_kind: ProviderKind,
1058 provider_identity: &str,
1059 fingerprint: &str,
1060 ) -> usize {
1061 let provider = canonical_provider_scope(provider_identity);
1062 let offerings = cache
1063 .get(&storage_provider(provider_kind, &provider), fingerprint)
1064 .map(|entry| entry.offerings.clone())
1065 .unwrap_or_default();
1066 let count = offerings.len();
1067 crate::provider_lake::replace_provider_live_snapshot_for_identity(
1068 provider_kind,
1069 &provider,
1070 CatalogSnapshot { offerings },
1071 );
1072 count
1073 }
1074
1075 /// Load and publish only the active route's exact provider/base-URL scope.
1076 ///
1077 /// A cache created for another custom endpoint or for an old endpoint override
1078 /// is retained on disk but cannot leak into the active picker.
1079 pub fn maybe_load_persisted_cache_for_config(config: &Config) -> usize {
1080 let Ok(identity) = config.active_provider_identity() else {
1081 return 0;
1082 };
1083 let provider = identity.provider;
1084 let provider_identity = canonical_provider_scope(identity.key.as_str());
1085 let fingerprint = base_url_fingerprint(&config.base_url_for_route(&identity));
1086 if is_account_scoped_scope(
1087 &storage_provider(provider, &provider_identity),
1088 &fingerprint,
1089 ) {
1090 forget_account_scoped_provider(provider, &provider_identity);
1091 return 0;
1092 }
1093 if let Ok(mut guard) = CACHE.write()
1094 && let Some(loaded) = load_from_disk()
1095 {
1096 // Keep session-only scopes that cannot exist on disk, while allowing a
1097 // newer durable scope from another Codewhale process to refresh this
1098 // process. Every in-process writer takes CACHE before the file lock, so
1099 // this read/merge cannot overwrite a concurrent local refresh.
1100 for (key, entry) in loaded.entries {
1101 let should_replace = guard
1102 .entries
1103 .get(&key)
1104 .is_none_or(|current| entry.fetched_at >= current.fetched_at);
1105 if should_replace {
1106 guard.entries.insert(key, entry);
1107 }
1108 }
1109 }
1110 CACHE
1111 .read()
1112 .map(|guard| {
1113 publish_exact_scope_for_identity(&guard, provider, &provider_identity, &fingerprint)
1114 })
1115 .unwrap_or(0)
1116 }
1117
1118 fn forget_account_scoped_provider(provider_kind: ProviderKind, provider: &str) {
1119 let provider = canonical_provider_scope(provider);
1120 if let Ok(mut cache) = CACHE.write() {
1121 cache
1122 .entries
1123 .retain(|_, entry| entry.provider != storage_provider(provider_kind, &provider));
1124 }
1125 crate::provider_lake::replace_provider_live_snapshot_for_identity(
1126 provider_kind,
1127 &provider,
1128 CatalogSnapshot::default(),
1129 );
1130 }
1131
1132 /// Begin a provider refresh and invalidate older in-flight results.
1133 ///
1134 /// Account-scoped Baseten and Codewhale routes additionally drop their prior
1135 /// in-memory rosters: the same URL can expose different models after a credential
1136 /// change, and no safe account identifier is available for cache reuse.
1137 #[cfg(test)]
1138 pub fn begin_refresh(provider: &str) -> ProviderCatalogRefreshTicket {
1139 begin_refresh_inner(inferred_provider_kind(provider), provider, None)
1140 }
1141
1142 pub fn begin_refresh_for_identity(
1143 provider_kind: ProviderKind,
1144 provider: &str,
1145 base_url: &str,
1146 ) -> ProviderCatalogRefreshTicket {
1147 begin_refresh_inner(
1148 provider_kind,
1149 provider,
1150 Some(base_url_fingerprint(base_url)),
1151 )
1152 }
1153
1154 fn begin_refresh_inner(
1155 provider_kind: ProviderKind,
1156 provider: &str,
1157 fingerprint: Option<String>,
1158 ) -> ProviderCatalogRefreshTicket {
1159 let provider = canonical_provider_scope(provider);
1160 let scope = storage_provider(provider_kind, &provider);
1161 // Hold the generation gate through account-roster invalidation, so an older
1162 // refresh can never publish between the new ticket and the clear.
1163 let generation = if let Ok(mut generations) = REFRESH_GENERATIONS.write() {
1164 let generation = generations.entry(scope.clone()).or_default();
1165 *generation = generation.saturating_add(1);
1166 if fingerprint
1167 .as_deref()
1168 .is_some_and(|fp| is_account_scoped_scope(&scope, fp))
1169 {
1170 forget_account_scoped_provider(provider_kind, &provider);
1171 }
1172 *generation
1173 } else {
1174 0
1175 };
1176 ProviderCatalogRefreshTicket {
1177 provider,
1178 provider_kind,
1179 fingerprint,
1180 generation,
1181 }
1182 }
1183
1184 fn with_current_ticket<T>(
1185 ticket: &ProviderCatalogRefreshTicket,
1186 provider: &str,
1187 operation: impl FnOnce() -> T,
1188 ) -> Option<T> {
1189 let provider = canonical_provider_scope(provider);
1190 if ticket.provider != provider {
1191 return None;
1192 }
1193 let generations = REFRESH_GENERATIONS.read().ok()?;
1194 if generations
1195 .get(&storage_provider(ticket.provider_kind, &ticket.provider))
1196 .copied()
1197 != Some(ticket.generation)
1198 {
1199 return None;
1200 }
1201 // Keep the generation read guard alive through publication. A newer
1202 // `begin_refresh` needs the write lock, so it cannot slip between the
1203 // current-ticket check and this operation's cache/lake update.
1204 let result = operation();
1205 drop(generations);
1206 Some(result)
1207 }
1208
1209 /// Record a successful refresh only if no newer refresh superseded it.
1210 pub fn record_success_if_current(
1211 ticket: &ProviderCatalogRefreshTicket,
1212 delta: ProviderCatalogDelta,
1213 ) -> Option<CatalogStatus> {
1214 let provider = canonical_provider_scope(&delta.provider);
1215 if ticket
1216 .fingerprint
1217 .as_ref()
1218 .is_some_and(|fp| fp != &delta.base_url_fingerprint)
1219 {
1220 return None;
1221 }
1222 with_current_ticket(ticket, &provider, || {
1223 record_success_for_identity(ticket.provider_kind, delta)
1224 })
1225 }
1226
1227 /// Record a failed refresh only if no newer refresh superseded it.
1228 pub fn record_failure_if_current(
1229 ticket: &ProviderCatalogRefreshTicket,
1230 provider: &str,
1231 fingerprint: &str,
1232 reason: CatalogRefreshError,
1233 ) -> Option<CatalogStatus> {
1234 let provider = canonical_provider_scope(provider);
1235 if ticket
1236 .fingerprint
1237 .as_deref()
1238 .is_some_and(|fp| fp != fingerprint)
1239 {
1240 return None;
1241 }
1242 with_current_ticket(ticket, &provider, || {
1243 record_failure_for_identity(ticket.provider_kind, &provider, fingerprint, reason)
1244 })
1245 }
1246
1247 /// Current freshness receipt for one exact provider/base-URL scope.
1248 ///
1249 /// Runtime route resolution uses this independently from picker visibility:
1250 /// stale or failed rows may remain selectable as an explicit fallback, but
1251 /// their limits, capabilities, and prices are not treated as current endpoint
1252 /// facts during execution.
1253 #[cfg(test)]
1254 pub fn status_for_scope(provider: &str, base_url: &str) -> CatalogStatus {
1255 let fingerprint = base_url_fingerprint(base_url);
1256 status_for_fingerprint(provider, &fingerprint)
1257 }
1258
1259 /// Current freshness receipt when the caller already owns the endpoint
1260 /// fingerprint (for example, an immutable usage-pricing receipt).
1261 #[cfg(test)]
1262 pub(crate) fn status_for_fingerprint(provider: &str, fingerprint: &str) -> CatalogStatus {
1263 status_for_route_fingerprint(inferred_provider_kind(provider), provider, fingerprint)
1264 }
1265
1266 pub(crate) fn status_for_route(
1267 provider: ProviderKind,
1268 identity: &str,
1269 base_url: &str,
1270 ) -> CatalogStatus {
1271 status_for_route_fingerprint(provider, identity, &base_url_fingerprint(base_url))
1272 }
1273
1274 fn status_for_route_fingerprint(
1275 kind: ProviderKind,
1276 provider: &str,
1277 fingerprint: &str,
1278 ) -> CatalogStatus {
1279 let provider = storage_provider(kind, provider);
1280 CACHE
1281 .read()
1282 .map(|cache| cache.status(&provider, fingerprint, now_unix()))
1283 .unwrap_or(CatalogStatus::Unknown)
1284 }
1285
1286 /// Freeze the exact reviewed provider-live rate row fresh at CodeWhale's
1287 /// pre-permit application-dispatch boundary.
1288 ///
1289 /// Status, scope, model, source, and rates are all read beneath one `CACHE`
1290 /// read guard. The returned value owns every fact needed by later auditing, so
1291 /// completion-time code never re-opens mutable catalog or provider-lake state.
1292 fn reviewed_provider_live_scope(
1293 provider: ProviderKind,
1294 provider_identity: &str,
1295 endpoint_fingerprint: &str,
1296 ) -> bool {
1297 match provider {
1298 ProviderKind::Openrouter => {
1299 provider_identity == ProviderKind::Openrouter.as_str()
1300 && endpoint_fingerprint
1301 == base_url_fingerprint(crate::config::DEFAULT_OPENROUTER_BASE_URL)
1302 }
1303 ProviderKind::Custom => {
1304 endpoint_fingerprint
1305 == base_url_fingerprint(codewhale_config::catalog::BASETEN_BASE_URL)
1306 }
1307 _ => false,
1308 }
1309 }
1310
1311 #[must_use]
1312 pub(crate) fn fresh_provider_live_pricing_quote_at(
1313 provider: ProviderKind,
1314 provider_identity: &str,
1315 wire_model: &str,
1316 endpoint_fingerprint: &str,
1317 dispatched_at_unix: u64,
1318 ) -> Option<ProviderLivePricingQuote> {
1319 let provider_identity = canonical_provider_scope(provider_identity);
1320 let wire_model = wire_model.trim();
1321 let endpoint_fingerprint = endpoint_fingerprint.trim();
1322 if provider_identity.is_empty()
1323 || wire_model.is_empty()
1324 || !reviewed_provider_live_scope(provider, &provider_identity, endpoint_fingerprint)
1325 {
1326 return None;
1327 }
1328
1329 let cache = CACHE.read().ok()?;
1330 let storage_scope = storage_provider(provider, &provider_identity);
1331 if cache.status(&storage_scope, endpoint_fingerprint, dispatched_at_unix)
1332 != CatalogStatus::Fresh
1333 {
1334 return None;
1335 }
1336 let entry = cache.get(&storage_scope, endpoint_fingerprint)?;
1337 if entry.provider != storage_scope
1338 || entry.base_url_fingerprint.trim() != endpoint_fingerprint
1339 || entry.fetched_at > dispatched_at_unix
1340 {
1341 return None;
1342 }
1343 let offering = entry.offerings.iter().find(|offering| {
1344 offering.provider.trim() == provider_identity && offering.wire_model_id.trim() == wire_model
1345 })?;
1346 let pricing = OfferingPricing::from_catalog_offering(offering)?;
1347 if pricing.provider.trim() != provider_identity
1348 || pricing.wire_model_id.trim() != wire_model
1349 || pricing.currency != Currency::Usd
1350 || pricing.provenance != PricingProvenance::ProviderLive
1351 || pricing.effective_at != Some(entry.fetched_at)
1352 || pricing.endpoint_fingerprint.as_deref() != Some(endpoint_fingerprint)
1353 || pricing.input_per_million.is_none()
1354 || pricing.output_per_million.is_none()
1355 {
1356 return None;
1357 }
1358 ProviderLivePricingQuote::from_pricing(
1359 provider,
1360 &provider_identity,
1361 wire_model,
1362 endpoint_fingerprint,
1363 entry.fetched_at,
1364 &pricing,
1365 )
1366 }
1367
1368 /// A price is a fact, so it is in scope exactly where a fact is.
1369 ///
1370 /// This defers to [`crate::provider_lake::cloud_facts_apply_to_route`] rather
1371 /// than repeating the scope table. The copy it replaces admitted the dual-wire
1372 /// and regional routes (`deepseek-anthropic`, `siliconflow-CN`) that the
1373 /// catalog gate refuses; no price actually escaped through it, because the
1374 /// offering lookup below independently returns a non-`CloudFacts` row on those
1375 /// routes and the quote then fails — but that is one authority masking another,
1376 /// not agreement, and it would become a real leak the moment either moved. The
1377 /// copy also carried its own `!= OpenaiCodex` test, which `cloud_facts::scope`
1378 /// has always enforced for every consumer.
1379 ///
1380 /// The one condition that is this file's own: the *configured* identity must be
1381 /// the canonical provider. A differently-named provider table pointing at the
1382 /// official host is a separate credential and billing relationship.
1383 fn cloud_pricing_scope(provider: ProviderKind, identity: &str, base_url: &str) -> bool {
1384 crate::provider_lake::cloud_facts_apply_to_route(provider, identity, base_url)
1385 }
1386
1387 /// Capture the effective mutable price authority once. Provider-owned live
1388 /// prices retain priority; signed cloud prices are admitted only on the exact
1389 /// canonical official route. The historical wire field name remains stable.
1390 pub(crate) fn configured_dispatch_pricing_quote_at(
1391 models: &[codewhale_config::catalog::configured::ConfiguredModel],
1392 provider: ProviderKind,
1393 identity: &str,
1394 model: &str,
1395 base_url: &str,
1396 dispatched_at: u64,
1397 ) -> Option<ProviderLivePricingQuote> {
1398 if provider == ProviderKind::OpenaiCodex {
1399 return None;
1400 }
1401 codewhale_config::catalog::configured::validate_configured_models(models).ok()?;
1402 let declared = models
1403 .iter()
1404 .find(|row| row.id == model && row.matches_route(identity, base_url))?;
1405 let cost = declared.cost.clone().unwrap_or_default();
1406 let pricing = OfferingPricing {
1407 provider: identity.to_string(),
1408 wire_model_id: model.to_string(),
1409 canonical_model: None,
1410 currency: Currency::Usd,
1411 input_per_million: cost.input,
1412 output_per_million: cost.output,
1413 cache_read_per_million: cost.cache_read,
1414 cache_write_per_million: cost.cache_write,
1415 provenance: PricingProvenance::UserOverride,
1416 effective_at: None,
1417 endpoint_fingerprint: Some(base_url_fingerprint(base_url)),
1418 };
1419 // Freeze even an unpriced declaration: missing rates must not fall through
1420 // to a same-named bundled or subsequently refreshed price.
1421 ProviderLivePricingQuote::from_pricing(
1422 provider,
1423 identity,
1424 model,
1425 &base_url_fingerprint(base_url),
1426 dispatched_at,
1427 &pricing,
1428 )
1429 }
1430
1431 /// Pick the dispatch quote from an operator declaration and the endpoint's
1432 /// catalog. A declared rate wins. A declaration with no rates yields to the
1433 /// catalog's price for the same exact endpoint (#6690), and is kept (frozen,
1434 /// rate-less) only when the catalog has none, so a same-named bundled price
1435 /// still cannot fill the gap.
1436 pub(crate) fn declared_or_catalog_quote(
1437 declared: Option<ProviderLivePricingQuote>,
1438 catalog: impl FnOnce() -> Option<ProviderLivePricingQuote>,
1439 ) -> Option<ProviderLivePricingQuote> {
1440 match declared {
1441 Some(quote) if quote.carries_rates() => Some(quote),
1442 declared => catalog().or(declared),
1443 }
1444 }
1445
1446 pub(crate) fn fresh_dispatch_pricing_quote_at(
1447 provider: ProviderKind,
1448 provider_identity: &str,
1449 wire_model: &str,
1450 base_url: &str,
1451 dispatched_at_unix: u64,
1452 ) -> Option<ProviderLivePricingQuote> {
1453 let endpoint_fingerprint = base_url_fingerprint(base_url);
1454 if let Some(quote) = fresh_provider_live_pricing_quote_at(
1455 provider,
1456 provider_identity,
1457 wire_model,
1458 &endpoint_fingerprint,
1459 dispatched_at_unix,
1460 ) {
1461 return Some(quote);
1462 }
1463 if !cloud_pricing_scope(provider, provider_identity, base_url) {
1464 return None;
1465 }
1466 let snapshot = codewhale_config::cloud_facts::overlay::snapshot();
1467 let facts = snapshot.facts.as_ref()?;
1468 let offering = crate::provider_lake::catalog_offering_for_route(
1469 provider,
1470 provider_identity,
1471 base_url,
1472 wire_model,
1473 )?;
1474 let codewhale_config::catalog::CatalogSource::CloudFacts {
1475 facts_version,
1476 key_id,
1477 fetched_at,
1478 valid_until,
1479 } = offering.pricing_source()
1480 else {
1481 return None;
1482 };
1483 // Provider cache files cannot authenticate a cloud price by copying a
1484 // source stamp. This row must be a projection of the current verified
1485 // overlay, with matching independent price authority and exact wire ID.
1486 if offering.wire_model_id != wire_model
1487 || !matches!(
1488 offering.source,
1489 codewhale_config::catalog::CatalogSource::CloudFacts { .. }
1490 )
1491 || *facts_version != facts.facts_version
1492 || key_id != &facts.key_id
1493 || *valid_until != facts.valid_until
1494 {
1495 return None;
1496 }
1497 let pricing = OfferingPricing::from_catalog_offering_at(&offering, dispatched_at_unix)?;
1498 let mut quote = ProviderLivePricingQuote::from_pricing(
1499 provider,
1500 provider_identity,
1501 wire_model,
1502 &endpoint_fingerprint,
1503 *fetched_at,
1504 &pricing,
1505 )?;
1506 quote.cloud_facts = Some(CloudFactsPricingSource {
1507 facts_version: *facts_version,
1508 key_id: key_id.clone(),
1509 valid_until: *valid_until,
1510 base_url: base_url.to_string(),
1511 });
1512 quote.catalog_revision = ProviderLivePricingQuote::revision_for(
1513 quote.provider,
1514 &quote.provider_identity,
1515 &quote.wire_model,
1516 &quote.endpoint_fingerprint,
1517 quote.catalog_fetched_at,
1518 &quote.currency,
1519 &quote.provenance,
1520 &quote.input_per_million,
1521 &quote.output_per_million,
1522 &quote.cache_read_per_million,
1523 &quote.cache_write_per_million,
1524 quote.cloud_facts.as_ref(),
1525 )?;
1526 quote.pricing_for_route(
1527 provider,
1528 provider_identity,
1529 wire_model,
1530 &endpoint_fingerprint,
1531 dispatched_at_unix,
1532 )?;
1533 (snapshot.generation == codewhale_config::cloud_facts::overlay::snapshot().generation)
1534 .then_some(quote)
1535 }
1536
1537 /// Record and atomically persist a successful provider refresh.
1538 ///
1539 /// `ProviderCatalogCache::record_success` replaces the exact scope, so models
1540 /// removed upstream disappear instead of accumulating forever.
1541 #[cfg(test)]
1542 pub fn record_success(delta: ProviderCatalogDelta) -> CatalogStatus {
1543 record_success_for_identity(inferred_provider_kind(&delta.provider), delta)
1544 }
1545
1546 fn record_success_for_identity(
1547 kind: ProviderKind,
1548 mut delta: ProviderCatalogDelta,
1549 ) -> CatalogStatus {
1550 let provider = canonical_provider_scope(&delta.provider);
1551 delta.provider = storage_provider(kind, &provider);
1552 if delta.offerings.iter().any(|row| {
1553 row.provider != provider
1554 || !crate::provider_lake::valid_catalog_model_id(&row.wire_model_id)
1555 || !provider_cost_source_allowed(row)
1556 }) {
1557 return record_failure_for_identity(
1558 kind,
1559 &provider,
1560 &delta.base_url_fingerprint,
1561 CatalogRefreshError::InvalidResponse,
1562 );
1563 }
1564 let fingerprint = delta.base_url_fingerprint.clone();
1565 let Ok(mut guard) = CACHE.write() else {
1566 return CatalogStatus::Unknown;
1567 };
1568 guard.record_success(delta, DEFAULT_PROVIDER_CATALOG_TTL_SECS);
1569 let persisted = persist_scope(&guard, &storage_provider(kind, &provider), &fingerprint);
1570 publish_exact_scope_for_identity(&guard, kind, &provider, &fingerprint);
1571 if persisted {
1572 CatalogStatus::Fresh
1573 } else {
1574 CatalogStatus::Unknown
1575 }
1576 }
1577
1578 /// Record a typed failure while preserving and republishing prior rows for the
1579 /// exact route scope.
1580 #[cfg(test)]
1581 pub fn record_failure(
1582 provider: &str,
1583 fingerprint: &str,
1584 reason: CatalogRefreshError,
1585 ) -> CatalogStatus {
1586 record_failure_for_identity(
1587 inferred_provider_kind(provider),
1588 provider,
1589 fingerprint,
1590 reason,
1591 )
1592 }
1593
1594 fn record_failure_for_identity(
1595 kind: ProviderKind,
1596 provider: &str,
1597 fingerprint: &str,
1598 reason: CatalogRefreshError,
1599 ) -> CatalogStatus {
1600 let provider = canonical_provider_scope(provider);
1601 let scope = storage_provider(kind, &provider);
1602 let Ok(mut guard) = CACHE.write() else {
1603 return CatalogStatus::Failed { reason };
1604 };
1605 guard.record_failure(&scope, fingerprint, reason);
1606 persist_failure_scope(&mut guard, &scope, fingerprint, reason);
1607 publish_exact_scope_for_identity(&guard, kind, &provider, fingerprint);
1608 CatalogStatus::Failed { reason }
1609 }
1610
1611 #[cfg(test)]
1612 pub(crate) fn reset_cache_for_test() {
1613 DISK_LOADED.store(false, Ordering::Release);
1614 if let Ok(mut cache) = CACHE.write() {
1615 *cache = ProviderCatalogCache::new();
1616 }
1617 }
1618
1619 #[cfg(test)]
1620 mod tests {
1621 use super::*;
1622 use crate::config::{ProviderConfig, ProviderKind, ProvidersConfig};
1623 use crate::test_support::{EnvVarGuard, lock_test_env};
1624 use codewhale_config::catalog::{CatalogOffering, CatalogSource};
1625
1626 fn delta(provider: &str, fingerprint: &str, ids: &[&str]) -> ProviderCatalogDelta {
1627 delta_at(provider, fingerprint, ids, now_unix())
1628 }
1629
1630 fn delta_at(
1631 provider: &str,
1632 fingerprint: &str,
1633 ids: &[&str],
1634 fetched_at: u64,
1635 ) -> ProviderCatalogDelta {
1636 ProviderCatalogDelta {
1637 provider: provider.to_string(),
1638 base_url_fingerprint: fingerprint.to_string(),
1639 fetched_at,
1640 offerings: ids
1641 .iter()
1642 .map(|id| CatalogOffering {
1643 provider: provider.to_string(),
1644 wire_model_id: (*id).to_string(),
1645 endpoint_key: "chat".to_string(),
1646 source: CatalogSource::Live {
1647 base_url_fingerprint: fingerprint.to_string(),
1648 fetched_at,
1649 },
1650 ..CatalogOffering::default()
1651 })
1652 .collect(),
1653 }
1654 }
1655
1656 fn scope(identity: &str) -> String {
1657 storage_provider(inferred_provider_kind(identity), identity)
1658 }
1659
1660 fn stored_delta(provider: &str, fingerprint: &str, ids: &[&str]) -> ProviderCatalogDelta {
1661 stored_delta_at(provider, fingerprint, ids, now_unix())
1662 }
1663
1664 fn stored_delta_at(
1665 provider: &str,
1666 fingerprint: &str,
1667 ids: &[&str],
1668 fetched_at: u64,
1669 ) -> ProviderCatalogDelta {
1670 let mut delta = delta_at(provider, fingerprint, ids, fetched_at);
1671 delta.provider = scope(provider);
1672 delta
1673 }
1674
1675 #[test]
1676 fn success_replaces_scope_and_failure_preserves_last_rows_on_disk() {
1677 let _env = lock_test_env();
1678 let _live = crate::provider_lake::lock_live_snapshot();
1679 let home = tempfile::tempdir().expect("home");
1680 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1681 if let Ok(mut cache) = CACHE.write() {
1682 *cache = ProviderCatalogCache::new();
1683 }
1684
1685 assert_eq!(
1686 record_success(delta("openrouter", "fp", &["old"])),
1687 CatalogStatus::Fresh
1688 );
1689 assert_eq!(
1690 record_success(delta("openrouter", "fp", &["new"])),
1691 CatalogStatus::Fresh
1692 );
1693 assert!(matches!(
1694 record_failure("openrouter", "fp", CatalogRefreshError::RateLimited),
1695 CatalogStatus::Failed {
1696 reason: CatalogRefreshError::RateLimited
1697 }
1698 ));
1699
1700 let loaded = load_from_disk().expect("persisted cache");
1701 let entry = loaded
1702 .get(&scope("openrouter"), "fp")
1703 .expect("OpenRouter scope");
1704 assert_eq!(entry.offerings.len(), 1);
1705 assert_eq!(entry.offerings[0].wire_model_id, "new");
1706 assert!(matches!(entry.status, CatalogStatus::Failed { .. }));
1707 }
1708
1709 #[test]
1710 fn baseten_workspace_roster_is_session_only_and_clears_before_reauthentication() {
1711 let _env = lock_test_env();
1712 let _live = crate::provider_lake::lock_live_snapshot();
1713 let home = tempfile::tempdir().expect("home");
1714 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1715 reset_cache_for_test();
1716 crate::provider_lake::clear_live_snapshot();
1717
1718 let base_url = codewhale_config::catalog::BASETEN_BASE_URL;
1719 let fingerprint = base_url_fingerprint(base_url);
1720 record_success(delta(
1721 codewhale_config::catalog::BASETEN_PROVIDER_ID,
1722 &fingerprint,
1723 &["workspace-a-only-model"],
1724 ));
1725 assert!(
1726 crate::provider_lake::all_catalog_models_for_provider_identity(
1727 ProviderKind::Custom,
1728 Some(codewhale_config::catalog::BASETEN_PROVIDER_ID),
1729 )
1730 .contains(&"workspace-a-only-model".to_string())
1731 );
1732 assert!(
1733 load_from_disk().is_none_or(|cache| cache
1734 .get(
1735 &scope(codewhale_config::catalog::BASETEN_PROVIDER_ID),
1736 &fingerprint
1737 )
1738 .is_none()),
1739 "an account-scoped Baseten roster must never be durable without a safe account id"
1740 );
1741
1742 let mut custom = std::collections::HashMap::new();
1743 custom.insert(
1744 codewhale_config::catalog::BASETEN_PROVIDER_ID.to_string(),
1745 ProviderConfig {
1746 kind: Some("openai-compatible".to_string()),
1747 base_url: Some(base_url.to_string()),
1748 model: Some(codewhale_config::catalog::BASETEN_DEFAULT_MODEL.to_string()),
1749 ..ProviderConfig::default()
1750 },
1751 );
1752 let config = Config {
1753 provider: Some(codewhale_config::catalog::BASETEN_PROVIDER_ID.to_string()),
1754 providers: Some(ProvidersConfig {
1755 custom,
1756 ..ProvidersConfig::default()
1757 }),
1758 ..Config::default()
1759 };
1760 assert_eq!(maybe_load_persisted_cache_for_config(&config), 0);
1761 assert!(matches!(
1762 status_for_scope(codewhale_config::catalog::BASETEN_PROVIDER_ID, base_url),
1763 CatalogStatus::Unknown
1764 ));
1765 assert!(
1766 !crate::provider_lake::all_catalog_models_for_provider_identity(
1767 ProviderKind::Custom,
1768 Some(codewhale_config::catalog::BASETEN_PROVIDER_ID),
1769 )
1770 .contains(&"workspace-a-only-model".to_string()),
1771 "a new credential attempt must not see the previous workspace roster"
1772 );
1773
1774 reset_cache_for_test();
1775 crate::provider_lake::clear_live_snapshot();
1776 }
1777
1778 #[test]
1779 fn superseded_refresh_ticket_cannot_publish_a_late_response() {
1780 let _env = lock_test_env();
1781 let _live = crate::provider_lake::lock_live_snapshot();
1782 let home = tempfile::tempdir().expect("home");
1783 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1784 reset_cache_for_test();
1785 crate::provider_lake::clear_live_snapshot();
1786
1787 let old = begin_refresh("openrouter");
1788 let current = begin_refresh("openrouter");
1789 assert!(
1790 record_success_if_current(&old, delta("openrouter", "fp", &["late-old-model"]))
1791 .is_none()
1792 );
1793 assert!(
1794 record_success_if_current(&current, delta("openrouter", "fp", &["current-model"]),)
1795 .is_some()
1796 );
1797 assert_eq!(
1798 CACHE
1799 .read()
1800 .expect("cache")
1801 .get(&scope("openrouter"), "fp")
1802 .expect("current scope")
1803 .offerings[0]
1804 .wire_model_id,
1805 "current-model"
1806 );
1807
1808 reset_cache_for_test();
1809 crate::provider_lake::clear_live_snapshot();
1810 }
1811
1812 #[test]
1813 fn current_ticket_holds_generation_gate_through_publication() {
1814 let ticket = begin_refresh("generation-barrier-provider");
1815 let entered = std::sync::Arc::new(std::sync::Barrier::new(2));
1816 let release = std::sync::Arc::new(std::sync::Barrier::new(2));
1817 let publish_entered = std::sync::Arc::clone(&entered);
1818 let publish_release = std::sync::Arc::clone(&release);
1819 let publisher = std::thread::spawn(move || {
1820 with_current_ticket(&ticket, "generation-barrier-provider", || {
1821 publish_entered.wait();
1822 publish_release.wait();
1823 })
1824 });
1825 entered.wait();
1826
1827 let (started_tx, started_rx) = std::sync::mpsc::channel();
1828 let (finished_tx, finished_rx) = std::sync::mpsc::channel();
1829 let newer = std::thread::spawn(move || {
1830 started_tx.send(()).expect("signal refresh start");
1831 let next = begin_refresh("generation-barrier-provider");
1832 finished_tx.send(next).expect("signal refresh finish");
1833 });
1834 started_rx.recv().expect("new refresh thread started");
1835 assert!(
1836 finished_rx
1837 .recv_timeout(std::time::Duration::from_millis(50))
1838 .is_err(),
1839 "a newer generation must wait until the accepted result finishes publication"
1840 );
1841
1842 release.wait();
1843 assert!(publisher.join().expect("publisher thread").is_some());
1844 assert!(
1845 finished_rx
1846 .recv_timeout(std::time::Duration::from_secs(1))
1847 .is_ok()
1848 );
1849 newer.join().expect("newer refresh thread");
1850 }
1851
1852 #[test]
1853 fn stale_process_snapshots_merge_exact_scopes_under_file_lock() {
1854 let _env = lock_test_env();
1855 let home = tempfile::tempdir().expect("home");
1856 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1857
1858 let mut process_a = ProviderCatalogCache::new();
1859 process_a.record_success(stored_delta("CustomA", "fp-a", &["upper-model"]), 60);
1860 persist_scope(&process_a, &scope("CustomA"), "fp-a");
1861
1862 // Simulate another process that started before A wrote and therefore
1863 // has an empty/stale in-memory snapshot. Its scoped write must merge A
1864 // from disk rather than replacing the whole envelope.
1865 let mut process_b = ProviderCatalogCache::new();
1866 process_b.record_success(stored_delta("customa", "fp-b", &["lower-model"]), 60);
1867 persist_scope(&process_b, &scope("customa"), "fp-b");
1868
1869 let loaded = load_from_disk().expect("merged durable cache");
1870 assert_eq!(
1871 loaded
1872 .get(&scope("CustomA"), "fp-a")
1873 .expect("case-sensitive upper scope")
1874 .offerings[0]
1875 .wire_model_id,
1876 "upper-model"
1877 );
1878 assert_eq!(
1879 loaded
1880 .get(&scope("customa"), "fp-b")
1881 .expect("case-sensitive lower scope")
1882 .offerings[0]
1883 .wire_model_id,
1884 "lower-model"
1885 );
1886 }
1887
1888 #[test]
1889 fn stale_process_failure_preserves_newer_durable_rows_for_the_same_scope() {
1890 let _env = lock_test_env();
1891 let _live = crate::provider_lake::lock_live_snapshot();
1892 let home = tempfile::tempdir().expect("home");
1893 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1894 reset_cache_for_test();
1895 crate::provider_lake::clear_live_snapshot();
1896
1897 // Process B began with this old roster and still holds it in memory.
1898 let mut stale_process = ProviderCatalogCache::new();
1899 stale_process.record_success(stored_delta_at("openrouter", "fp", &["old-model"], 1), 60);
1900 persist_scope(&stale_process, &scope("openrouter"), "fp");
1901 *CACHE.write().expect("cache") = stale_process;
1902
1903 // Process A completes a newer successful refresh for the same scope.
1904 let mut newer_process = ProviderCatalogCache::new();
1905 newer_process.record_success(stored_delta_at("openrouter", "fp", &["new-model"], 2), 60);
1906 persist_scope(&newer_process, &scope("openrouter"), "fp");
1907
1908 // B then fails. Its failure status is current, but its old rows are
1909 // not: the transaction must retain A's newer durable roster.
1910 assert!(matches!(
1911 record_failure("openrouter", "fp", CatalogRefreshError::Network),
1912 CatalogStatus::Failed {
1913 reason: CatalogRefreshError::Network
1914 }
1915 ));
1916 let in_memory = CACHE.read().expect("cache");
1917 let entry = in_memory
1918 .get(&scope("openrouter"), "fp")
1919 .expect("failed scope");
1920 assert_eq!(entry.offerings[0].wire_model_id, "new-model");
1921 assert!(matches!(entry.status, CatalogStatus::Failed { .. }));
1922 drop(in_memory);
1923
1924 let durable = load_from_disk().expect("durable cache");
1925 let entry = durable
1926 .get(&scope("openrouter"), "fp")
1927 .expect("durable failed scope");
1928 assert_eq!(entry.offerings[0].wire_model_id, "new-model");
1929 assert!(matches!(entry.status, CatalogStatus::Failed { .. }));
1930
1931 reset_cache_for_test();
1932 crate::provider_lake::clear_live_snapshot();
1933 }
1934
1935 #[test]
1936 fn oversized_cache_file_is_rejected_before_allocation() {
1937 let _env = lock_test_env();
1938 let home = tempfile::tempdir().expect("home");
1939 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
1940 let path = cache_path().expect("cache path");
1941 fs::create_dir_all(path.parent().expect("catalog directory")).expect("catalog directory");
1942 fs::File::create(&path)
1943 .and_then(|file| file.set_len(MAX_CACHE_BYTES + 1))
1944 .expect("sparse oversized cache");
1945 assert!(load_from_disk().is_none());
1946 }
1947
1948 #[test]
1949 fn bounded_persistence_evicts_failed_then_stale_scopes_and_keeps_exact_owner() {
1950 let mut cache = ProviderCatalogCache::new();
1951 cache.record_success(
1952 stored_delta_at("failed", "fp", &["failed-model"], 10),
1953 1_000,
1954 );
1955 cache.record_failure(&scope("failed"), "fp", CatalogRefreshError::Network);
1956 cache.record_success(stored_delta_at("stale", "fp", &["stale-model"], 20), 1);
1957 cache.record_success(stored_delta_at("fresh", "fp", &["fresh-model"], 30), 1_000);
1958 cache.record_success(
1959 stored_delta_at("protected", "fp", &["protected-model"], 40),
1960 1_000,
1961 );
1962
1963 let compacted = bounded_cache_for_persistence(
1964 cache,
1965 Some((&scope("protected"), "fp")),
1966 100,
1967 CachePersistenceLimits {
1968 max_bytes: u64::MAX,
1969 max_scopes: 2,
1970 max_rows: 100,
1971 },
1972 )
1973 .expect("bounded cache");
1974
1975 assert!(compacted.get(&scope("protected"), "fp").is_some());
1976 assert!(compacted.get(&scope("fresh"), "fp").is_some());
1977 assert!(compacted.get(&scope("failed"), "fp").is_none());
1978 assert!(compacted.get(&scope("stale"), "fp").is_none());
1979 }
1980
1981 #[test]
1982 fn bounded_persistence_evicts_whole_scopes_and_refuses_an_oversized_owner() {
1983 let mut cache = ProviderCatalogCache::new();
1984 cache.record_success(
1985 stored_delta_at("protected", "fp", &["one", "two"], 40),
1986 1_000,
1987 );
1988 cache.record_success(stored_delta_at("other", "fp", &["other"], 30), 1_000);
1989 let limits = CachePersistenceLimits {
1990 max_bytes: u64::MAX,
1991 max_scopes: 10,
1992 max_rows: 2,
1993 };
1994
1995 let compacted = bounded_cache_for_persistence(
1996 cache.clone(),
1997 Some((&scope("protected"), "fp")),
1998 50,
1999 limits,
2000 )
2001 .expect("other scope can be evicted whole");
2002 assert_eq!(
2003 compacted
2004 .get(&scope("protected"), "fp")
2005 .expect("protected roster")
2006 .offerings
2007 .len(),
2008 2
2009 );
2010 assert!(compacted.get(&scope("other"), "fp").is_none());
2011
2012 let mut oversized = cache;
2013 oversized.record_success(
2014 stored_delta_at("protected", "fp", &["one", "two", "three"], 50),
2015 1_000,
2016 );
2017 assert!(
2018 bounded_cache_for_persistence(oversized, Some((&scope("protected"), "fp")), 50, limits,)
2019 .is_err(),
2020 "a provider roster must be refused, never partially persisted"
2021 );
2022 }
2023
2024 #[test]
2025 fn bounded_cache_write_matches_read_limit_and_round_trips_after_compaction() {
2026 let directory = tempfile::tempdir().expect("cache directory");
2027 let path = directory.path().join(CACHE_FILE);
2028 let mut protected_only = ProviderCatalogCache::new();
2029 protected_only.record_success(
2030 stored_delta_at("protected", "fp", &["protected-model"], 40),
2031 1_000,
2032 );
2033 let exact_bytes = persisted_envelope_len(&protected_only).expect("encoded length");
2034 let limits = CachePersistenceLimits {
2035 max_bytes: exact_bytes,
2036 max_scopes: 10,
2037 max_rows: 10,
2038 };
2039
2040 let mut combined = protected_only;
2041 combined.record_success(
2042 stored_delta_at(
2043 "evicted",
2044 "fp",
2045 &["this-entire-scope-does-not-fit-the-byte-bound"],
2046 30,
2047 ),
2048 1_000,
2049 );
2050 write_bounded_cache(&path, combined, Some((&scope("protected"), "fp")), limits)
2051 .expect("bounded disk write");
2052
2053 assert!(fs::metadata(&path).expect("cache metadata").len() <= exact_bytes);
2054 let loaded = load_from_disk_unlocked_with_limit(&path, exact_bytes)
2055 .expect("bounded cache must remain readable under the same cap");
2056 assert!(loaded.get(&scope("protected"), "fp").is_some());
2057 assert!(loaded.get(&scope("evicted"), "fp").is_none());
2058 }
2059
2060 #[test]
2061 fn baseten_alias_roster_is_session_only_and_keeps_exact_ownership() {
2062 let _env = lock_test_env();
2063 let _live = crate::provider_lake::lock_live_snapshot();
2064 let home = tempfile::tempdir().expect("home");
2065 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2066 reset_cache_for_test();
2067 crate::provider_lake::clear_live_snapshot();
2068
2069 let alias = "base-ten";
2070 let fingerprint = base_url_fingerprint(codewhale_config::catalog::BASETEN_BASE_URL);
2071 record_success(delta(alias, &fingerprint, &["alias-workspace-model"]));
2072
2073 assert!(
2074 crate::provider_lake::all_catalog_models_for_provider_identity(
2075 ProviderKind::Custom,
2076 Some(alias),
2077 )
2078 .contains(&"alias-workspace-model".to_string())
2079 );
2080 assert!(
2081 !crate::provider_lake::all_catalog_models_for_provider_identity(
2082 ProviderKind::Custom,
2083 Some(codewhale_config::catalog::BASETEN_PROVIDER_ID),
2084 )
2085 .contains(&"alias-workspace-model".to_string()),
2086 "a reviewed schema alias must not collapse distinct exact table ownership"
2087 );
2088 assert!(
2089 load_from_disk().is_none_or(|cache| cache.get(&scope(alias), &fingerprint).is_none()),
2090 "every Baseten schema alias must remain session-only"
2091 );
2092
2093 reset_cache_for_test();
2094 crate::provider_lake::clear_live_snapshot();
2095 }
2096
2097 #[test]
2098 fn different_base_url_fingerprints_do_not_share_rows() {
2099 let mut cache = ProviderCatalogCache::new();
2100 cache.record_success(stored_delta("baseten", "one", &["model-one"]), 60);
2101 cache.record_success(stored_delta("baseten", "two", &["model-two"]), 60);
2102 assert_eq!(
2103 cache.get(&scope("baseten"), "one").unwrap().offerings[0].wire_model_id,
2104 "model-one"
2105 );
2106 assert_eq!(
2107 cache.get(&scope("baseten"), "two").unwrap().offerings[0].wire_model_id,
2108 "model-two"
2109 );
2110 }
2111
2112 #[test]
2113 fn missing_cache_for_changed_base_url_clears_the_previous_provider_partition() {
2114 let _live = crate::provider_lake::lock_live_snapshot();
2115 crate::provider_lake::clear_live_snapshot();
2116 let mut cache = ProviderCatalogCache::new();
2117 cache.record_success(
2118 stored_delta("baseten", "old-fp", &["old-endpoint-model"]),
2119 60,
2120 );
2121
2122 assert_eq!(
2123 publish_exact_scope_for_identity(&cache, ProviderKind::Custom, "baseten", "old-fp"),
2124 1
2125 );
2126 assert_eq!(
2127 crate::provider_lake::all_catalog_models_for_provider_identity(
2128 crate::config::ProviderKind::Custom,
2129 Some("baseten"),
2130 ),
2131 vec!["old-endpoint-model".to_string()]
2132 );
2133
2134 assert_eq!(
2135 publish_exact_scope_for_identity(&cache, ProviderKind::Custom, "baseten", "new-fp"),
2136 0
2137 );
2138 let after_switch = crate::provider_lake::all_catalog_models_for_provider_identity(
2139 crate::config::ProviderKind::Custom,
2140 Some("baseten"),
2141 );
2142 assert!(
2143 after_switch.is_empty(),
2144 "rows from the old Baseten endpoint must not survive a fingerprint change, and no compiled seed replaces them (#6289)"
2145 );
2146 crate::provider_lake::clear_live_snapshot();
2147 }
2148
2149 #[test]
2150 fn disk_reload_rehydrates_and_exposes_six_hundred_openrouter_models() {
2151 let _env = lock_test_env();
2152 let _live = crate::provider_lake::lock_live_snapshot();
2153 let home = tempfile::tempdir().expect("home");
2154 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2155 crate::provider_lake::clear_live_snapshot();
2156 if let Ok(mut cache) = CACHE.write() {
2157 *cache = ProviderCatalogCache::new();
2158 }
2159
2160 let config = Config {
2161 provider: Some("openrouter".to_string()),
2162 providers: Some(ProvidersConfig {
2163 openrouter: ProviderConfig {
2164 base_url: Some("https://synthetic.openrouter.invalid/api/v1".to_string()),
2165 ..ProviderConfig::default()
2166 },
2167 ..ProvidersConfig::default()
2168 }),
2169 ..Config::default()
2170 };
2171 let identity = config.active_provider_identity().unwrap();
2172 let provider = identity.key.as_str();
2173 let fingerprint = base_url_fingerprint(&config.active_route_base_url());
2174 let fetched_at = now_unix();
2175 let ids: Vec<String> = (0..600)
2176 .map(|index| format!("synthetic/openrouter-model-{index:03}"))
2177 .collect();
2178 let status = record_success(ProviderCatalogDelta {
2179 provider: provider.to_string(),
2180 base_url_fingerprint: fingerprint,
2181 fetched_at,
2182 offerings: ids
2183 .iter()
2184 .map(|id| CatalogOffering {
2185 provider: provider.to_string(),
2186 wire_model_id: id.clone(),
2187 endpoint_key: "chat".to_string(),
2188 source: CatalogSource::Live {
2189 base_url_fingerprint: base_url_fingerprint(&config.active_route_base_url()),
2190 fetched_at,
2191 },
2192 ..CatalogOffering::default()
2193 })
2194 .collect(),
2195 });
2196 assert_eq!(status, CatalogStatus::Fresh);
2197 assert!(cache_path().is_some_and(|path| path.is_file()));
2198 assert_eq!(
2199 crate::provider_lake::all_catalog_models_for_provider(ProviderKind::Openrouter),
2200 ids,
2201 "the string compatibility publisher must retain built-in OpenRouter ownership"
2202 );
2203 assert!(
2204 crate::provider_lake::all_catalog_models_for_provider_identity(
2205 ProviderKind::Custom,
2206 Some("openrouter"),
2207 )
2208 .is_empty(),
2209 "built-in OpenRouter rows must not enter the custom namespace"
2210 );
2211
2212 // Simulate a new process: remove both in-memory owners, then republish
2213 // only through the durable startup load path.
2214 if let Ok(mut cache) = CACHE.write() {
2215 *cache = ProviderCatalogCache::new();
2216 }
2217 crate::provider_lake::clear_live_snapshot();
2218
2219 assert_eq!(maybe_load_persisted_cache_for_config(&config), 600);
2220 let visible =
2221 crate::provider_lake::all_catalog_models_for_provider(ProviderKind::Openrouter);
2222 assert_eq!(visible.len(), 600);
2223 assert_eq!(visible.first(), ids.first());
2224 assert_eq!(visible.last(), ids.last());
2225
2226 if let Ok(mut cache) = CACHE.write() {
2227 *cache = ProviderCatalogCache::new();
2228 }
2229 crate::provider_lake::clear_live_snapshot();
2230 }
2231 #[test]
2232 fn typed_refresh_keeps_builtin_and_same_named_custom_scopes_separate_on_disk() {
2233 let _env = lock_test_env();
2234 let _live = crate::provider_lake::lock_live_snapshot();
2235 let home = tempfile::tempdir().unwrap();
2236 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2237 reset_cache_for_test();
2238 crate::provider_lake::clear_live_snapshot();
2239 let endpoint = "https://api.openai.com/v1";
2240 let fingerprint = base_url_fingerprint(endpoint);
2241 let built_in = begin_refresh_for_identity(ProviderKind::Openai, "openai", endpoint);
2242 let custom = begin_refresh_for_identity(ProviderKind::Custom, "openai", endpoint);
2243 assert_eq!(
2244 record_success_if_current(
2245 &built_in,
2246 delta("openai", &fingerprint, &["built-in-model"])
2247 ),
2248 Some(CatalogStatus::Fresh)
2249 );
2250 assert_eq!(
2251 record_success_if_current(&custom, delta("openai", &fingerprint, &["custom-model"])),
2252 Some(CatalogStatus::Fresh)
2253 );
2254 reset_cache_for_test();
2255 crate::provider_lake::clear_live_snapshot();
2256 for (kind, expected) in [
2257 (ProviderKind::Openai, "built-in-model"),
2258 (ProviderKind::Custom, "custom-model"),
2259 ] {
2260 let entry = cached_entry_for_route(kind, "openai", endpoint)
2261 .unwrap()
2262 .unwrap();
2263 assert_eq!(entry.offerings[0].wire_model_id, expected);
2264 assert_eq!(entry.offerings[0].provider, "openai");
2265 }
2266 assert!(
2267 cached_entry_for_route(ProviderKind::Custom, "OpenAI", endpoint)
2268 .unwrap()
2269 .is_none()
2270 );
2271 reset_cache_for_test();
2272 }
2273
2274 #[test]
2275 fn typed_refresh_rejects_wrong_endpoint_and_superseded_endpoint_response() {
2276 let _env = lock_test_env();
2277 let _live = crate::provider_lake::lock_live_snapshot();
2278 let home = tempfile::tempdir().unwrap();
2279 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2280 reset_cache_for_test();
2281 crate::provider_lake::clear_live_snapshot();
2282 let first_url = "https://first.invalid/v1";
2283 let second_url = "https://second.invalid/v1";
2284 let first = begin_refresh_for_identity(ProviderKind::Custom, "ExactRoute", first_url);
2285 assert!(
2286 record_success_if_current(
2287 &first,
2288 delta(
2289 "ExactRoute",
2290 &base_url_fingerprint(second_url),
2291 &["wrong-endpoint"]
2292 )
2293 )
2294 .is_none()
2295 );
2296 let second = begin_refresh_for_identity(ProviderKind::Custom, "ExactRoute", second_url);
2297 assert!(
2298 record_failure_if_current(
2299 &second,
2300 "ExactRoute",
2301 &base_url_fingerprint(first_url),
2302 CatalogRefreshError::Network
2303 )
2304 .is_none()
2305 );
2306 assert_eq!(
2307 record_success_if_current(
2308 &second,
2309 delta(
2310 "ExactRoute",
2311 &base_url_fingerprint(second_url),
2312 &["current-model"]
2313 )
2314 ),
2315 Some(CatalogStatus::Fresh)
2316 );
2317 assert!(
2318 record_success_if_current(
2319 &first,
2320 delta(
2321 "ExactRoute",
2322 &base_url_fingerprint(first_url),
2323 &["late-model"]
2324 )
2325 )
2326 .is_none()
2327 );
2328 assert!(
2329 cached_entry_for_route(ProviderKind::Custom, "ExactRoute", first_url)
2330 .unwrap()
2331 .is_none()
2332 );
2333 assert_eq!(
2334 crate::provider_lake::catalog_models_for_route(
2335 ProviderKind::Custom,
2336 "ExactRoute",
2337 second_url
2338 ),
2339 vec!["current-model"]
2340 );
2341 reset_cache_for_test();
2342 crate::provider_lake::clear_live_snapshot();
2343 }
2344
2345 #[test]
2346 fn differently_named_baseten_endpoint_never_persists_account_roster() {
2347 let _env = lock_test_env();
2348 let _live = crate::provider_lake::lock_live_snapshot();
2349 let home = tempfile::tempdir().unwrap();
2350 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2351 reset_cache_for_test();
2352 crate::provider_lake::clear_live_snapshot();
2353 let endpoint = codewhale_config::catalog::BASETEN_BASE_URL;
2354 let ticket = begin_refresh_for_identity(ProviderKind::Custom, "TeamServing", endpoint);
2355 assert_eq!(
2356 record_success_if_current(
2357 &ticket,
2358 delta(
2359 "TeamServing",
2360 &base_url_fingerprint(endpoint),
2361 &["private-workspace-model"]
2362 )
2363 ),
2364 Some(CatalogStatus::Fresh)
2365 );
2366 assert_eq!(
2367 crate::provider_lake::catalog_models_for_route(
2368 ProviderKind::Custom,
2369 "TeamServing",
2370 endpoint
2371 ),
2372 vec!["private-workspace-model"]
2373 );
2374 assert!(
2375 !fs::read_to_string(cache_path().unwrap())
2376 .unwrap()
2377 .contains("private-workspace-model")
2378 );
2379 let _new_credentials =
2380 begin_refresh_for_identity(ProviderKind::Custom, "TeamServing", endpoint);
2381 assert!(
2382 crate::provider_lake::catalog_models_for_route(
2383 ProviderKind::Custom,
2384 "TeamServing",
2385 endpoint
2386 )
2387 .is_empty()
2388 );
2389 reset_cache_for_test();
2390 assert!(
2391 cached_entry_for_route(ProviderKind::Custom, "TeamServing", endpoint)
2392 .unwrap()
2393 .is_none()
2394 );
2395 crate::provider_lake::clear_live_snapshot();
2396 }
2397
2398 #[test]
2399 fn codewhale_account_rosters_are_memory_only_and_replaced_after_credential_refresh() {
2400 let _env = lock_test_env();
2401 let _live = crate::provider_lake::lock_live_snapshot();
2402 let home = tempfile::tempdir().unwrap();
2403 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2404 reset_cache_for_test();
2405 crate::provider_lake::clear_live_snapshot();
2406 for (kind, identity, endpoint) in [
2407 (
2408 ProviderKind::Codewhale,
2409 "codewhale",
2410 ProviderKind::Codewhale.provider().default_base_url(),
2411 ),
2412 (
2413 ProviderKind::Codewhale,
2414 "codewhale",
2415 "https://codewhale.account.invalid/v1",
2416 ),
2417 (
2418 ProviderKind::Custom,
2419 "PrivateAccount",
2420 ProviderKind::Codewhale.provider().default_base_url(),
2421 ),
2422 ] {
2423 let fingerprint = base_url_fingerprint(endpoint);
2424 let old = begin_refresh_for_identity(kind, identity, endpoint);
2425 assert_eq!(
2426 record_success_if_current(
2427 &old,
2428 delta(identity, &fingerprint, &["private-old-account-model"])
2429 ),
2430 Some(CatalogStatus::Fresh)
2431 );
2432 assert!(
2433 cached_entry_for_route(kind, identity, endpoint)
2434 .unwrap()
2435 .is_some()
2436 );
2437 assert!(
2438 !fs::read_to_string(cache_path().unwrap())
2439 .unwrap()
2440 .contains("private-old-account-model")
2441 );
2442 let current = begin_refresh_for_identity(kind, identity, endpoint);
2443 assert!(
2444 cached_entry_for_route(kind, identity, endpoint)
2445 .unwrap()
2446 .is_none()
2447 );
2448 assert!(
2449 !crate::provider_lake::catalog_models_for_route(kind, identity, endpoint)
2450 .contains(&"private-old-account-model".to_string())
2451 );
2452 assert!(
2453 record_success_if_current(
2454 &old,
2455 delta(identity, &fingerprint, &["private-old-account-model"])
2456 )
2457 .is_none()
2458 );
2459 assert_eq!(
2460 record_success_if_current(
2461 &current,
2462 delta(identity, &fingerprint, &["private-new-account-model"])
2463 ),
2464 Some(CatalogStatus::Fresh)
2465 );
2466 assert_eq!(
2467 crate::provider_lake::catalog_models_for_route(kind, identity, endpoint),
2468 vec!["private-new-account-model"]
2469 );
2470 assert!(
2471 !fs::read_to_string(cache_path().unwrap())
2472 .unwrap()
2473 .contains("private-new-account-model")
2474 );
2475 reset_cache_for_test();
2476 crate::provider_lake::clear_live_snapshot();
2477 assert!(
2478 cached_entry_for_route(kind, identity, endpoint)
2479 .unwrap()
2480 .is_none()
2481 );
2482 }
2483 reset_cache_for_test();
2484 crate::provider_lake::clear_live_snapshot();
2485 }
2486
2487 #[cfg(unix)]
2488 #[test]
2489 fn cache_reader_rejects_links_and_special_files_without_blocking() {
2490 use std::os::unix::ffi::OsStrExt as _;
2491 let dir = tempfile::tempdir().unwrap();
2492 let target = dir.path().join("target.json");
2493 let envelope = PersistedProviderCatalogs {
2494 schema_version: CACHE_SCHEMA_VERSION,
2495 cache: ProviderCatalogCache::new(),
2496 };
2497 fs::write(&target, serde_json::to_vec(&envelope).unwrap()).unwrap();
2498 let symlink = dir.path().join("symlink.json");
2499 std::os::unix::fs::symlink(&target, &symlink).unwrap();
2500 assert!(load_from_disk_unlocked_with_limit(&symlink, MAX_CACHE_BYTES).is_none());
2501 let hardlink = dir.path().join("hardlink.json");
2502 fs::hard_link(&target, &hardlink).unwrap();
2503 assert!(load_from_disk_unlocked_with_limit(&hardlink, MAX_CACHE_BYTES).is_none());
2504 let fifo = dir.path().join("fifo.json");
2505 let fifo_c = std::ffi::CString::new(fifo.as_os_str().as_bytes()).unwrap();
2506 // SAFETY: fifo_c owns a NUL-terminated pathname for the call.
2507 assert_eq!(unsafe { libc::mkfifo(fifo_c.as_ptr(), 0o600) }, 0);
2508 assert!(load_from_disk_unlocked_with_limit(&fifo, MAX_CACHE_BYTES).is_none());
2509 assert!(open_cache_lock(&fifo).is_err());
2510 assert!(load_from_disk_unlocked_with_limit(dir.path(), MAX_CACHE_BYTES).is_none());
2511 }
2512
2513 struct CloudQuoteReset;
2514 impl Drop for CloudQuoteReset {
2515 fn drop(&mut self) {
2516 codewhale_config::cloud_facts::overlay::clear();
2517 crate::provider_lake::clear_live_snapshot();
2518 reset_cache_for_test();
2519 }
2520 }
2521
2522 fn install_cloud_quote_fixture(
2523 channel: &str,
2524 provider: ProviderKind,
2525 version: u64,
2526 input: f64,
2527 now: u64,
2528 ) {
2529 use codewhale_config::cloud_facts::{
2530 CloudFactsState, CloudFactsStatus, FactsOrigin, ModelFact, PricingFact, ScopedFacts,
2531 overlay,
2532 };
2533 let ticket = overlay::configure(true, "cloud-quote-fixture").unwrap();
2534 assert!(overlay::publish(
2535 &ticket,
2536 Some(ScopedFacts {
2537 channel: channel.into(),
2538 facts_version: version,
2539 key_id: "cwf-test-only".into(),
2540 valid_until: Some(now + 60),
2541 models: vec![ModelFact {
2542 provider: provider.as_str().into(),
2543 id: "cloud-quote-fixture".into(),
2544 context_window: Some(16_384),
2545 pricing: Some(PricingFact {
2546 input_per_m: Some(input),
2547 output_per_m: Some(2.0),
2548 cache_read_per_m: None,
2549 }),
2550 ..Default::default()
2551 }],
2552 ..Default::default()
2553 }),
2554 CloudFactsStatus {
2555 state: CloudFactsState::Verified {
2556 channel: channel.into(),
2557 facts_version: version,
2558 key_id: "cwf-test-only".into(),
2559 fetched_at: now,
2560 origin: FactsOrigin::LocalFile,
2561 stale: false,
2562 patches: 1,
2563 defaults: 0,
2564 announcements: 0,
2565 },
2566 ..Default::default()
2567 },
2568 ));
2569 }
2570
2571 #[test]
2572 fn cloud_quote_freezes_actual_dispatch_and_survives_refresh_disable_and_reload() {
2573 let _env = lock_test_env();
2574 let _live = crate::provider_lake::lock_live_snapshot();
2575 let home = tempfile::tempdir().unwrap();
2576 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2577 let _enabled = EnvVarGuard::remove("CODEWHALE_DISABLE_CLOUD_FACTS");
2578 let _reset = CloudQuoteReset;
2579 reset_cache_for_test();
2580 crate::provider_lake::clear_live_snapshot();
2581 let now = chrono::Utc::now();
2582 let at = now.timestamp() as u64;
2583 let base = crate::config::DEFAULT_OPENAI_BASE_URL;
2584 install_cloud_quote_fixture("quote-frozen", ProviderKind::Openai, 91, 1.0, at);
2585 let route = crate::cost_status::EffectiveRouteEnvelope::capture(
2586 None,
2587 ProviderKind::Openai,
2588 "openai",
2589 "cloud-quote-fixture",
2590 Some(base),
2591 now,
2592 );
2593 let quote = route
2594 .provider_live_pricing
2595 .as_ref()
2596 .expect("frozen cloud price");
2597 assert_eq!(quote.provenance, PricingProvenance::CloudFacts);
2598 assert_eq!(quote.cloud_facts.as_ref().unwrap().facts_version, 91);
2599 assert_eq!(quote.input_per_million.as_deref(), Some("1"));
2600 assert!(quote.cache_read_per_million.is_none());
2601 let encoded = serde_json::to_string(&route).unwrap();
2602 install_cloud_quote_fixture("quote-frozen", ProviderKind::Openai, 92, 9.0, at);
2603 let newer = fresh_dispatch_pricing_quote_at(
2604 ProviderKind::Openai,
2605 "openai",
2606 "cloud-quote-fixture",
2607 base,
2608 at,
2609 )
2610 .unwrap();
2611 assert_eq!(newer.input_per_million.as_deref(), Some("9"));
2612 assert_ne!(newer.catalog_revision, quote.catalog_revision);
2613 codewhale_config::cloud_facts::overlay::clear();
2614 assert!(
2615 fresh_dispatch_pricing_quote_at(
2616 ProviderKind::Openai,
2617 "openai",
2618 "cloud-quote-fixture",
2619 base,
2620 at,
2621 )
2622 .is_none()
2623 );
2624 let restored: crate::cost_status::EffectiveRouteEnvelope =
2625 serde_json::from_str(&encoded).unwrap();
2626 assert_eq!(restored, route);
2627 let pricing = restored
2628 .provider_live_pricing
2629 .as_ref()
2630 .unwrap()
2631 .pricing_for_route(
2632 ProviderKind::Openai,
2633 "openai",
2634 "cloud-quote-fixture",
2635 &base_url_fingerprint(base),
2636 at,
2637 )
2638 .unwrap();
2639 assert_eq!(pricing.input_per_million, Some(1.0));
2640 // Expiry is checked at dispatch; loading later does not reprice history.
2641 assert!(
2642 quote
2643 .pricing_for_route(
2644 ProviderKind::Openai,
2645 "openai",
2646 "cloud-quote-fixture",
2647 &base_url_fingerprint(base),
2648 at + 61,
2649 )
2650 .is_none()
2651 );
2652 }
2653
2654 #[test]
2655 fn cloud_quote_rejects_cross_route_and_modified_source_receipts() {
2656 let _env = lock_test_env();
2657 let _live = crate::provider_lake::lock_live_snapshot();
2658 let home = tempfile::tempdir().unwrap();
2659 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2660 let _enabled = EnvVarGuard::remove("CODEWHALE_DISABLE_CLOUD_FACTS");
2661 let _reset = CloudQuoteReset;
2662 reset_cache_for_test();
2663 crate::provider_lake::clear_live_snapshot();
2664 let at = now_unix();
2665 let base = crate::config::DEFAULT_OPENAI_BASE_URL;
2666 install_cloud_quote_fixture("quote-binding", ProviderKind::Openai, 93, 1.0, at);
2667 let quote = fresh_dispatch_pricing_quote_at(
2668 ProviderKind::Openai,
2669 "openai",
2670 "cloud-quote-fixture",
2671 base,
2672 at,
2673 )
2674 .unwrap();
2675 for (provider, identity, endpoint) in [
2676 (ProviderKind::Openai, "openai", "https://proxy.example/v1"),
2677 (ProviderKind::Custom, "openai", base),
2678 (ProviderKind::Openai, "named-openai", base),
2679 (
2680 ProviderKind::Openai,
2681 "openai",
2682 "https://secret@api.openai.com/v1",
2683 ),
2684 ] {
2685 assert!(
2686 fresh_dispatch_pricing_quote_at(
2687 provider,
2688 identity,
2689 "cloud-quote-fixture",
2690 endpoint,
2691 at
2692 )
2693 .is_none()
2694 );
2695 }
2696 let value = serde_json::to_value(&quote).unwrap();
2697 for (field, replacement) in [
2698 ("facts_version", serde_json::json!(94)),
2699 ("key_id", serde_json::json!("cwf-relabelled")),
2700 ("valid_until", serde_json::json!(at + 600)),
2701 ("base_url", serde_json::json!("https://proxy.example/v1")),
2702 ] {
2703 let mut modified = value.clone();
2704 modified["cloud_facts"][field] = replacement;
2705 assert!(
2706 serde_json::from_value::<ProviderLivePricingQuote>(modified).is_err(),
2707 "changed {field} accepted"
2708 );
2709 }
2710 let mut invalid = quote.clone();
2711 invalid.cloud_facts.as_mut().unwrap().base_url = "https://secret@api.openai.com/v1".into();
2712 assert!(serde_json::to_value(invalid).unwrap().is_null());
2713 assert!(
2714 quote
2715 .pricing_for_route(
2716 ProviderKind::Openai,
2717 "openai",
2718 "different-model",
2719 &base_url_fingerprint(base),
2720 at,
2721 )
2722 .is_none()
2723 );
2724 }
2725
2726 #[test]
2727 fn cloud_quote_static_endpoint_survives_operator_endpoint_changes() {
2728 let _env = lock_test_env();
2729 let _live = crate::provider_lake::lock_live_snapshot();
2730 let home = tempfile::tempdir().unwrap();
2731 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2732 let _enabled = EnvVarGuard::remove("CODEWHALE_DISABLE_CLOUD_FACTS");
2733 let _reset = CloudQuoteReset;
2734 reset_cache_for_test();
2735 crate::provider_lake::clear_live_snapshot();
2736 let at = now_unix();
2737 let base = ProviderKind::Codewhale.provider().default_base_url();
2738 let private = "https://private.example/v1?api_key=public-test-marker";
2739 let declared = EnvVarGuard::set("CODEWHALE_API_BASE", private);
2740 // The ordinary runtime credential contract intentionally accepts this
2741 // explicit operator route. It is not a public signed-facts authority.
2742 assert!(codewhale_config::provider_base_url_is_official(
2743 codewhale_config::ProviderKind::Codewhale,
2744 private,
2745 ));
2746 install_cloud_quote_fixture("quote-static", ProviderKind::Codewhale, 94, 1.0, at);
2747 assert!(
2748 fresh_dispatch_pricing_quote_at(
2749 ProviderKind::Codewhale,
2750 "codewhale",
2751 "cloud-quote-fixture",
2752 private,
2753 at,
2754 )
2755 .is_none()
2756 );
2757 let quote = fresh_dispatch_pricing_quote_at(
2758 ProviderKind::Codewhale,
2759 "codewhale",
2760 "cloud-quote-fixture",
2761 base,
2762 at,
2763 )
2764 .unwrap();
2765 let encoded = serde_json::to_string(&quote).unwrap();
2766 assert!(!encoded.contains("private.example"));
2767 assert!(!encoded.contains("public-test-marker"));
2768 drop(declared);
2769 let removed = EnvVarGuard::remove("CODEWHALE_API_BASE");
2770 let reloaded: ProviderLivePricingQuote = serde_json::from_str(&encoded).unwrap();
2771 assert_eq!(reloaded, quote);
2772 drop(removed);
2773 let _changed = EnvVarGuard::set("CODEWHALE_API_BASE", "https://another-private.example/v1");
2774 assert_eq!(
2775 serde_json::from_str::<ProviderLivePricingQuote>(&encoded).unwrap(),
2776 quote
2777 );
2778 assert!(
2779 quote
2780 .pricing_for_route(
2781 ProviderKind::Codewhale,
2782 "codewhale",
2783 "cloud-quote-fixture",
2784 &base_url_fingerprint(base),
2785 at,
2786 )
2787 .is_some()
2788 );
2789 }
2790
2791 #[test]
2792 fn durable_provider_catalog_cannot_claim_cloud_or_override_price_authority() {
2793 let dir = tempfile::tempdir().unwrap();
2794 let path = dir.path().join("catalog.json");
2795 let mut cache = ProviderCatalogCache::new();
2796 cache.record_success(
2797 stored_delta("openai", "fp", &["fixture"]),
2798 DEFAULT_PROVIDER_CATALOG_TTL_SECS,
2799 );
2800 let baseline = PersistedProviderCatalogs {
2801 schema_version: CACHE_SCHEMA_VERSION,
2802 cache,
2803 };
2804 fs::write(&path, serde_json::to_vec(&baseline).unwrap()).unwrap();
2805 assert!(load_from_disk_unlocked(&path).is_some());
2806 for source in [
2807 CatalogSource::CloudFacts {
2808 facts_version: 7,
2809 key_id: "cwf-test-only".into(),
2810 fetched_at: now_unix(),
2811 valid_until: None,
2812 },
2813 CatalogSource::ConfigOverride,
2814 CatalogSource::UserOverride,
2815 CatalogSource::Live {
2816 base_url_fingerprint: "other".into(),
2817 fetched_at: now_unix(),
2818 },
2819 ] {
2820 let mut forged = baseline.clone();
2821 forged.cache.entries.values_mut().next().unwrap().offerings[0].cost_source =
2822 Some(source);
2823 fs::write(&path, serde_json::to_vec(&forged).unwrap()).unwrap();
2824 assert!(load_from_disk_unlocked(&path).is_none());
2825 }
2826 }
2827
2828 #[test]
2829 fn pin_drift_reads_a_fresh_roster_persisted_by_an_earlier_process() {
2830 // #6035: `codewhale doctor` and a just-started TUI have not touched
2831 // the in-process cache yet; the durable fresh roster must still count.
2832 let _env = lock_test_env();
2833 let _live = crate::provider_lake::lock_live_snapshot();
2834 let home = tempfile::tempdir().expect("home");
2835 let _home = EnvVarGuard::set("CODEWHALE_HOME", home.path());
2836 reset_cache_for_test();
2837 let config = Config::default();
2838 let base_url = config.base_url_for_route(
2839 &config
2840 .resolve_provider_selection_identity("deepseek")
2841 .unwrap(),
2842 );
2843 let fingerprint = base_url_fingerprint(&base_url);
2844 assert_eq!(
2845 record_success(delta("deepseek", &fingerprint, &["deepseek-flash"])),
2846 CatalogStatus::Fresh
2847 );
2848 // A new process: nothing loaded in memory, the roster only on disk.
2849 reset_cache_for_test();
2850 assert_eq!(
2851 pin_missing_from_fresh_roster(&config, "deepseek", "deepseek-retired"),
2852 Some(true)
2853 );
2854 assert_eq!(
2855 pin_missing_from_fresh_roster(&config, "deepseek", "deepseek-flash"),
2856 Some(false)
2857 );
2858 reset_cache_for_test();
2859 }
2860 }
2861
2861 lines RUST