1use std::collections::{BTreeMap, HashMap, HashSet};
2use std::sync::{Arc, Mutex, OnceLock};
3use std::time::{Duration, Instant};
4
5use anyhow::{Context, Result};
6use futures_util::stream::{FuturesUnordered, StreamExt};
7use reqwest::Client;
8use serde::{Deserialize, Serialize};
9use tracing::{info, warn};
10
11use crate::balance::{
12 BalanceSnapshotStatus, ProviderBalanceSnapshot, ProviderUsageAlert, ProviderUsageAlertKind,
13 ProviderUsageModelStat, ProviderUsageRateSnapshot, ProviderUsageWindow,
14};
15use crate::config::{ProxyConfig, ServiceConfigManager, proxy_home_dir};
16use crate::lb::LbState;
17use crate::policy_actions::{PolicyAction, PolicyActionKind};
18use crate::pricing::UsdAmount;
19use crate::provider_signals::{
20 ProviderSignal, ProviderSignalKind, ProviderSignalSource, ProviderSignalTarget,
21};
22use crate::runtime_identity::ProviderEndpointKey;
23use crate::state::ProxyState;
24
25#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq)]
26#[serde(rename_all = "snake_case")]
27enum ProviderKind {
28 BudgetHttpJson,
30 YescodeProfile,
32 #[serde(
34 rename = "openai_balance_http_json",
35 alias = "open_ai_balance_http_json",
36 alias = "relay_balance_http_json"
37 )]
38 OpenAiBalanceHttpJson,
39 #[serde(rename = "sub2api_usage", alias = "sub2api_usage_http_json")]
41 Sub2ApiUsage,
42 #[serde(rename = "sub2api_auth_me", alias = "sub2api_auth_me_http_json")]
44 Sub2ApiAuthMe,
45 #[serde(
47 rename = "new_api_token_usage",
48 alias = "new_api_token_usage_http_json"
49 )]
50 NewApiTokenUsage,
51 NewApiUserSelf,
53 #[serde(
55 rename = "rightcode_account_summary",
56 alias = "right_code_account_summary",
57 alias = "rightcode"
58 )]
59 RightCodeAccountSummary,
60 #[serde(
62 rename = "openai_organization_costs",
63 alias = "openai_org_costs",
64 alias = "openai_costs"
65 )]
66 OpenAiOrganizationCosts,
67}
68
69impl ProviderKind {
70 fn source_name(&self) -> &'static str {
71 match self {
72 ProviderKind::BudgetHttpJson => "usage_provider:budget_http_json",
73 ProviderKind::YescodeProfile => "usage_provider:yescode_profile",
74 ProviderKind::OpenAiBalanceHttpJson => "usage_provider:openai_balance_http_json",
75 ProviderKind::Sub2ApiUsage => "usage_provider:sub2api_usage",
76 ProviderKind::Sub2ApiAuthMe => "usage_provider:sub2api_auth_me",
77 ProviderKind::NewApiTokenUsage => "usage_provider:new_api_token_usage",
78 ProviderKind::NewApiUserSelf => "usage_provider:new_api_user_self",
79 ProviderKind::RightCodeAccountSummary => "usage_provider:rightcode_account_summary",
80 ProviderKind::OpenAiOrganizationCosts => "usage_provider:openai_organization_costs",
81 }
82 }
83
84 fn default_endpoint(&self) -> Option<&'static str> {
85 match self {
86 ProviderKind::OpenAiBalanceHttpJson => Some("{{base_url}}/user/balance"),
87 ProviderKind::Sub2ApiUsage => Some("{{base_url}}/v1/usage"),
88 ProviderKind::Sub2ApiAuthMe => Some("{{base_url}}/api/v1/auth/me"),
89 ProviderKind::NewApiTokenUsage => Some("{{base_url}}/api/usage/token/"),
90 ProviderKind::NewApiUserSelf => Some("{{base_url}}/api/user/self"),
91 ProviderKind::RightCodeAccountSummary => {
92 Some("https://www.right.codes/account/summary")
93 }
94 ProviderKind::OpenAiOrganizationCosts => {
95 Some("{{base_url}}/v1/organization/costs?start_time={{unix_days_ago:30}}&limit=30")
96 }
97 _ => None,
98 }
99 }
100}
101
102#[derive(Debug, Deserialize, Serialize, Default, Clone)]
103#[serde(default)]
104struct UsageProviderExtractConfig {
105 #[serde(skip_serializing_if = "Vec::is_empty")]
106 remaining_balance_paths: Vec<String>,
107 #[serde(skip_serializing_if = "Vec::is_empty")]
108 subscription_balance_paths: Vec<String>,
109 #[serde(skip_serializing_if = "Vec::is_empty")]
110 paygo_balance_paths: Vec<String>,
111 #[serde(skip_serializing_if = "Vec::is_empty")]
112 monthly_budget_paths: Vec<String>,
113 #[serde(skip_serializing_if = "Vec::is_empty")]
114 monthly_spent_paths: Vec<String>,
115 #[serde(skip_serializing_if = "Vec::is_empty")]
116 exhausted_paths: Vec<String>,
117 #[serde(skip_serializing_if = "Option::is_none")]
118 remaining_divisor: Option<u64>,
119 #[serde(skip_serializing_if = "Option::is_none")]
120 monthly_budget_divisor: Option<u64>,
121 #[serde(skip_serializing_if = "Option::is_none")]
122 monthly_spent_divisor: Option<u64>,
123 #[serde(skip_serializing_if = "bool_is_false")]
124 derive_budget_from_remaining_and_spent: bool,
125 #[serde(skip_serializing_if = "bool_is_false")]
126 derive_remaining_from_budget_and_spent: bool,
127}
128
129impl UsageProviderExtractConfig {
130 fn is_empty(&self) -> bool {
131 self.remaining_balance_paths.is_empty()
132 && self.subscription_balance_paths.is_empty()
133 && self.paygo_balance_paths.is_empty()
134 && self.monthly_budget_paths.is_empty()
135 && self.monthly_spent_paths.is_empty()
136 && self.exhausted_paths.is_empty()
137 && self.remaining_divisor.is_none()
138 && self.monthly_budget_divisor.is_none()
139 && self.monthly_spent_divisor.is_none()
140 && !self.derive_budget_from_remaining_and_spent
141 && !self.derive_remaining_from_budget_and_spent
142 }
143}
144
145#[derive(Debug, Deserialize, Serialize)]
146struct UsageProviderConfig {
147 id: String,
148 kind: ProviderKind,
149 domains: Vec<String>,
150 #[serde(default)]
151 endpoint: String,
152 #[serde(default)]
153 token_env: Option<String>,
154 #[serde(default, skip_serializing_if = "bool_is_false")]
155 require_token_env: bool,
156 #[serde(default)]
157 poll_interval_secs: Option<u64>,
158 #[serde(
159 default = "default_refresh_on_request",
160 skip_serializing_if = "bool_is_true"
161 )]
162 refresh_on_request: bool,
163 #[serde(
164 default = "default_trust_exhaustion_for_routing",
165 skip_serializing_if = "bool_is_true"
166 )]
167 trust_exhaustion_for_routing: bool,
168 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
169 headers: BTreeMap<String, String>,
170 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
171 variables: BTreeMap<String, String>,
172 #[serde(default, skip_serializing_if = "UsageProviderExtractConfig::is_empty")]
173 extract: UsageProviderExtractConfig,
174}
175
176#[derive(Debug, Deserialize, Serialize, Default)]
177struct UsageProvidersFile {
178 #[serde(default)]
179 providers: Vec<UsageProviderConfig>,
180}
181
182#[derive(Debug, Clone)]
183struct UpstreamRef {
184 station_name: String,
185 index: usize,
186 provider_endpoint: Option<ProviderEndpointKey>,
187}
188
189#[derive(Debug, Clone)]
190struct UsageProviderTarget {
191 upstream: UpstreamRef,
192 base_url: String,
193 provider_id: Option<String>,
194}
195
196#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
197struct UsageProviderTargetKey {
198 station_name: String,
199 upstream_index: usize,
200}
201
202#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
203pub struct UsageProviderRefreshSummary {
204 pub providers_configured: usize,
205 pub providers_matched: usize,
206 pub upstreams_matched: usize,
207 pub attempted: usize,
208 pub refreshed: usize,
209 pub failed: usize,
210 pub missing_token: usize,
211 #[serde(skip_serializing_if = "usize_is_zero")]
212 pub auto_attempted: usize,
213 #[serde(skip_serializing_if = "usize_is_zero")]
214 pub auto_refreshed: usize,
215 #[serde(skip_serializing_if = "usize_is_zero")]
216 pub auto_failed: usize,
217 #[serde(skip_serializing_if = "usize_is_zero")]
218 pub deduplicated: usize,
219}
220
221#[derive(Debug, Clone, Copy, PartialEq, Eq)]
222enum UsageProviderRefreshOutcome {
223 Refreshed,
224 Failed,
225 MissingToken,
226}
227
228struct RefreshProviderTargetParams<'a> {
229 client: &'a Client,
230 provider: &'a UsageProviderConfig,
231 target: &'a UsageProviderTarget,
232 cfg: &'a ProxyConfig,
233 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
234 state: &'a Arc<ProxyState>,
235 service_name: &'a str,
236 interval_secs: u64,
237}
238
239static LAST_USAGE_POLL: OnceLock<Mutex<HashMap<String, Instant>>> = OnceLock::new();
241static REQUEST_BALANCE_QUEUE: OnceLock<Mutex<HashMap<RequestBalanceQueueKey, Instant>>> =
242 OnceLock::new();
243
244const DEFAULT_POLL_INTERVAL_SECS: u64 = 60;
245const MIN_POLL_INTERVAL_SECS: u64 = 20;
247pub const REQUEST_BALANCE_REFRESH_DELAY: Duration = Duration::from_secs(8);
248const BALANCE_REFRESH_CONCURRENCY: usize = 6;
249const BALANCE_HTTP_REQUEST_TIMEOUT: Duration = Duration::from_secs(6);
250const LOW_BALANCE_ALERT_THRESHOLD_USD: &str = "10";
251const EXPIRING_SOON_WINDOW_SECS: u64 = 7 * 24 * 60 * 60;
252const AUTO_PROVIDER_ID_PREFIX: &str = "auto:balance:";
253const AUTO_PROBE_KINDS: [ProviderKind; 5] = [
254 ProviderKind::RightCodeAccountSummary,
255 ProviderKind::Sub2ApiUsage,
256 ProviderKind::NewApiTokenUsage,
257 ProviderKind::NewApiUserSelf,
258 ProviderKind::OpenAiBalanceHttpJson,
259];
260
261#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
262enum RequestBalanceQueueKey {
263 ProviderEndpoint(ProviderEndpointKey),
264 LegacyUpstream {
265 service_name: String,
266 station_name: String,
267 upstream_index: usize,
268 },
269}
270
271#[derive(Debug, Clone, Copy, PartialEq, Eq)]
272enum RequestBalanceQueueDue {
273 Due,
274 NotDue(Duration),
275 Missing,
276}
277
278#[derive(Debug, Clone, Copy, PartialEq, Eq)]
279enum RequestBalancePollOutcome {
280 Attempted,
281 Deferred(Duration),
282 Skipped,
283}
284
285fn bool_is_false(value: &bool) -> bool {
286 !*value
287}
288
289fn bool_is_true(value: &bool) -> bool {
290 *value
291}
292
293fn usize_is_zero(value: &usize) -> bool {
294 *value == 0
295}
296
297fn default_refresh_on_request() -> bool {
298 true
299}
300
301fn default_trust_exhaustion_for_routing() -> bool {
302 true
303}
304
305fn unix_now_ms() -> u64 {
306 std::time::SystemTime::now()
307 .duration_since(std::time::UNIX_EPOCH)
308 .map(|d| d.as_millis() as u64)
309 .unwrap_or(0)
310}
311
312fn unix_now_secs() -> u64 {
313 std::time::SystemTime::now()
314 .duration_since(std::time::UNIX_EPOCH)
315 .map(|d| d.as_secs())
316 .unwrap_or(0)
317}
318
319fn stale_after_ms(fetched_at_ms: u64, interval_secs: u64) -> Option<u64> {
320 fetched_at_ms.checked_add(interval_secs.saturating_mul(3).saturating_mul(1_000))
321}
322
323fn snapshot_refresh_interval_secs(provider: &UsageProviderConfig) -> u64 {
324 let interval_secs = provider
325 .poll_interval_secs
326 .unwrap_or(DEFAULT_POLL_INTERVAL_SECS);
327 if interval_secs == 0 {
328 DEFAULT_POLL_INTERVAL_SECS
329 } else {
330 interval_secs.max(MIN_POLL_INTERVAL_SECS)
331 }
332}
333
334fn effective_poll_interval_secs(provider: &UsageProviderConfig) -> Option<u64> {
335 if !provider.refresh_on_request {
336 return None;
337 }
338
339 let interval_secs = provider
340 .poll_interval_secs
341 .unwrap_or(DEFAULT_POLL_INTERVAL_SECS);
342 if interval_secs == 0 {
343 return None;
344 }
345 Some(interval_secs.max(MIN_POLL_INTERVAL_SECS))
346}
347
348fn remaining_poll_cooldown(last: Instant, interval_secs: u64, now: Instant) -> Option<Duration> {
349 let interval = Duration::from_secs(interval_secs);
350 let elapsed = now.saturating_duration_since(last);
351 interval.checked_sub(elapsed).filter(|d| !d.is_zero())
352}
353
354fn usage_providers_path() -> std::path::PathBuf {
355 proxy_home_dir().join("usage_providers.json")
356}
357
358fn service_manager<'a>(cfg: &'a ProxyConfig, service_name: &str) -> &'a ServiceConfigManager {
359 match service_name {
360 "claude" => &cfg.claude,
361 _ => &cfg.codex,
362 }
363}
364
365fn default_provider_config(
366 id: &str,
367 kind: ProviderKind,
368 domains: Vec<&str>,
369 endpoint: &str,
370 extract: UsageProviderExtractConfig,
371) -> UsageProviderConfig {
372 UsageProviderConfig {
373 id: id.to_string(),
374 kind,
375 domains: domains.into_iter().map(str::to_string).collect(),
376 endpoint: endpoint.to_string(),
377 token_env: None,
378 require_token_env: false,
379 poll_interval_secs: Some(60),
380 refresh_on_request: true,
381 trust_exhaustion_for_routing: true,
382 headers: BTreeMap::new(),
383 variables: BTreeMap::new(),
384 extract,
385 }
386}
387
388fn default_rightcode_provider_config(id: &str) -> UsageProviderConfig {
389 let mut provider = default_provider_config(
390 id,
391 ProviderKind::RightCodeAccountSummary,
392 vec!["www.right.codes", "right.codes"],
393 "https://www.right.codes/account/summary",
394 UsageProviderExtractConfig::default(),
395 );
396 provider.trust_exhaustion_for_routing = false;
400 provider
401}
402
403fn host_from_base_url(base_url: &str) -> Option<String> {
404 reqwest::Url::parse(base_url)
405 .ok()
406 .and_then(|url| url.host_str().map(|host| host.to_ascii_lowercase()))
407}
408
409fn is_official_openai_base_url(base_url: &str) -> bool {
410 host_from_base_url(base_url).as_deref() == Some("api.openai.com")
411}
412
413fn is_rightcode_base_url(base_url: &str) -> bool {
414 matches!(
415 host_from_base_url(base_url).as_deref(),
416 Some("www.right.codes" | "right.codes")
417 )
418}
419
420fn provider_id_component(value: &str) -> String {
421 let component = value
422 .chars()
423 .map(|ch| {
424 if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' || ch == '.' {
425 ch
426 } else {
427 '-'
428 }
429 })
430 .collect::<String>()
431 .trim_matches('-')
432 .to_string();
433 if component.is_empty() {
434 "station".to_string()
435 } else {
436 component
437 }
438}
439
440fn auto_provider_id(target: &UsageProviderTarget) -> String {
441 if let Some(provider_id) = target
442 .provider_id
443 .as_deref()
444 .map(str::trim)
445 .filter(|value| !value.is_empty())
446 {
447 return provider_id.to_string();
448 }
449 format!(
450 "{}{}:{}",
451 AUTO_PROVIDER_ID_PREFIX,
452 provider_id_component(&target.upstream.station_name),
453 target.upstream.index
454 )
455}
456
457fn auto_usage_provider(target: &UsageProviderTarget, kind: ProviderKind) -> UsageProviderConfig {
458 let mut provider = UsageProviderConfig {
459 id: auto_provider_id(target),
460 kind,
461 domains: host_from_base_url(&target.base_url)
462 .into_iter()
463 .collect::<Vec<_>>(),
464 endpoint: String::new(),
465 token_env: None,
466 require_token_env: false,
467 poll_interval_secs: Some(DEFAULT_POLL_INTERVAL_SECS),
468 refresh_on_request: true,
469 trust_exhaustion_for_routing: true,
470 headers: BTreeMap::new(),
471 variables: BTreeMap::new(),
472 extract: UsageProviderExtractConfig::default(),
473 };
474 if matches!(kind, ProviderKind::RightCodeAccountSummary) {
475 provider.trust_exhaustion_for_routing = false;
476 }
477 provider
478}
479
480fn auto_target_matches_provider_id_filter(
481 target: &UsageProviderTarget,
482 provider_id_filter: Option<&str>,
483) -> bool {
484 match provider_id_filter {
485 Some(filter) => auto_provider_id(target) == filter,
486 None => true,
487 }
488}
489
490fn first_auto_probe_kind(target: &UsageProviderTarget) -> ProviderKind {
491 if is_rightcode_base_url(&target.base_url) {
492 ProviderKind::RightCodeAccountSummary
493 } else {
494 ProviderKind::Sub2ApiUsage
495 }
496}
497
498fn auto_openai_official_provider(target: &UsageProviderTarget) -> UsageProviderConfig {
499 let mut provider = auto_usage_provider(target, ProviderKind::OpenAiOrganizationCosts);
500 provider.token_env = Some("OPENAI_ADMIN_KEY".to_string());
501 provider.require_token_env = true;
502 provider.refresh_on_request = false;
503 provider.trust_exhaustion_for_routing = false;
504 provider
505}
506
507fn default_providers() -> UsageProvidersFile {
508 let openrouter_extract = UsageProviderExtractConfig {
509 monthly_budget_paths: vec!["data.total_credits".to_string()],
510 monthly_spent_paths: vec!["data.total_usage".to_string()],
511 derive_remaining_from_budget_and_spent: true,
512 ..Default::default()
513 };
514
515 let novita_extract = UsageProviderExtractConfig {
516 remaining_balance_paths: vec!["availableBalance".to_string()],
517 remaining_divisor: Some(10_000),
518 ..Default::default()
519 };
520
521 let mut openai_official = default_provider_config(
522 "openai-official-costs",
523 ProviderKind::OpenAiOrganizationCosts,
524 vec!["api.openai.com"],
525 "https://api.openai.com/v1/organization/costs?start_time={{unix_days_ago:30}}&limit=30",
526 UsageProviderExtractConfig::default(),
527 );
528 openai_official.token_env = Some("OPENAI_ADMIN_KEY".to_string());
529 openai_official.require_token_env = true;
530 openai_official.refresh_on_request = false;
531 openai_official.trust_exhaustion_for_routing = false;
532
533 UsageProvidersFile {
534 providers: vec![
535 default_rightcode_provider_config("rightcode"),
536 default_provider_config(
537 "packycode",
538 ProviderKind::BudgetHttpJson,
539 vec!["packycode.com"],
540 "https://www.packycode.com/api/backend/users/info",
541 UsageProviderExtractConfig::default(),
542 ),
543 default_provider_config(
544 "yescode",
545 ProviderKind::YescodeProfile,
546 vec!["yes.vg"],
548 "https://co.yes.vg/api/v1/auth/profile",
549 UsageProviderExtractConfig::default(),
550 ),
551 default_provider_config(
552 "deepseek",
553 ProviderKind::OpenAiBalanceHttpJson,
554 vec!["api.deepseek.com"],
555 "https://api.deepseek.com/user/balance",
556 UsageProviderExtractConfig::default(),
557 ),
558 default_provider_config(
559 "stepfun",
560 ProviderKind::OpenAiBalanceHttpJson,
561 vec!["api.stepfun.ai", "api.stepfun.com"],
562 "https://api.stepfun.com/v1/accounts",
563 UsageProviderExtractConfig::default(),
564 ),
565 default_provider_config(
566 "siliconflow",
567 ProviderKind::OpenAiBalanceHttpJson,
568 vec!["api.siliconflow.cn", "api.siliconflow.com"],
569 "{{base_url}}/v1/user/info",
570 UsageProviderExtractConfig::default(),
571 ),
572 default_provider_config(
573 "openrouter",
574 ProviderKind::OpenAiBalanceHttpJson,
575 vec!["openrouter.ai"],
576 "https://openrouter.ai/api/v1/credits",
577 openrouter_extract,
578 ),
579 default_provider_config(
580 "novita",
581 ProviderKind::OpenAiBalanceHttpJson,
582 vec!["api.novita.ai"],
583 "https://api.novita.ai/v3/user/balance",
584 novita_extract,
585 ),
586 openai_official,
587 ],
588 }
589}
590
591fn load_providers() -> UsageProvidersFile {
592 let path = usage_providers_path();
593 if let Ok(text) = std::fs::read_to_string(&path)
594 && let Ok(file) = serde_json::from_str::<UsageProvidersFile>(&text)
595 {
596 return file;
597 }
598
599 let default = default_providers();
601 if let Ok(text) = serde_json::to_string_pretty(&default) {
602 if let Some(parent) = path.parent() {
603 let _ = std::fs::create_dir_all(parent);
604 }
605 let _ = std::fs::write(&path, text);
606 }
607 default
608}
609
610fn domain_matches(base_url: &str, domains: &[String]) -> bool {
611 let url = match reqwest::Url::parse(base_url) {
612 Ok(u) => u,
613 Err(_) => return false,
614 };
615 let host = match url.host_str() {
616 Some(h) => h,
617 None => return false,
618 };
619 let host = host.to_ascii_lowercase();
620 for d in domains {
621 let domain = d.trim().to_ascii_lowercase();
622 if host == domain || host.ends_with(&format!(".{}", domain)) {
623 return true;
624 }
625 }
626 false
627}
628
629fn matching_provider_targets(
630 cfg: &ProxyConfig,
631 service_name: &str,
632 provider: &UsageProviderConfig,
633 station_name_filter: Option<&str>,
634) -> Vec<UsageProviderTarget> {
635 let mut stations: Vec<_> = service_manager(cfg, service_name)
636 .stations()
637 .iter()
638 .collect();
639 stations.sort_by_key(|(name, _)| name.as_str());
640
641 let mut targets = Vec::new();
642 for (station_name, service) in stations {
643 if station_name_filter.is_some_and(|filter| filter != station_name.as_str()) {
644 continue;
645 }
646 for (index, upstream) in service.upstreams.iter().enumerate() {
647 if domain_matches(&upstream.base_url, &provider.domains) {
648 targets.push(UsageProviderTarget {
649 upstream: UpstreamRef {
650 station_name: station_name.clone(),
651 index,
652 provider_endpoint: upstream.provider_endpoint_key(service_name),
653 },
654 base_url: upstream.base_url.clone(),
655 provider_id: upstream.tags.get("provider_id").cloned(),
656 });
657 }
658 }
659 }
660
661 targets
662}
663
664fn usage_provider_targets(
665 cfg: &ProxyConfig,
666 service_name: &str,
667 station_name_filter: Option<&str>,
668) -> Vec<UsageProviderTarget> {
669 let mut stations: Vec<_> = service_manager(cfg, service_name)
670 .stations()
671 .iter()
672 .collect();
673 stations.sort_by_key(|(name, _)| name.as_str());
674
675 let mut targets = Vec::new();
676 for (station_name, service) in stations {
677 if station_name_filter.is_some_and(|filter| filter != station_name.as_str()) {
678 continue;
679 }
680 for (index, upstream) in service.upstreams.iter().enumerate() {
681 targets.push(UsageProviderTarget {
682 upstream: UpstreamRef {
683 station_name: station_name.clone(),
684 index,
685 provider_endpoint: upstream.provider_endpoint_key(service_name),
686 },
687 base_url: upstream.base_url.clone(),
688 provider_id: upstream.tags.get("provider_id").cloned(),
689 });
690 }
691 }
692
693 targets
694}
695
696fn target_key(target: &UsageProviderTarget) -> UsageProviderTargetKey {
697 UsageProviderTargetKey {
698 station_name: target.upstream.station_name.clone(),
699 upstream_index: target.upstream.index,
700 }
701}
702
703fn enqueue_request_balance_refresh(key: RequestBalanceQueueKey) -> Option<Duration> {
704 let now = Instant::now();
705 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
706 let mut queue = match queue.lock() {
707 Ok(queue) => queue,
708 Err(_) => return None,
709 };
710
711 match queue.get(&key).copied() {
712 Some(due_at) if due_at > now => None,
713 Some(_) => Some(Duration::ZERO),
714 None => {
715 queue.insert(key, now + REQUEST_BALANCE_REFRESH_DELAY);
716 Some(REQUEST_BALANCE_REFRESH_DELAY)
717 }
718 }
719}
720
721fn schedule_request_balance_refresh_at(key: RequestBalanceQueueKey, due_at: Instant) {
722 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
723 if let Ok(mut queue) = queue.lock() {
724 queue.insert(key, due_at);
725 }
726}
727
728fn take_request_balance_refresh_if_due(key: &RequestBalanceQueueKey) -> RequestBalanceQueueDue {
729 let now = Instant::now();
730 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
731 let mut queue = match queue.lock() {
732 Ok(queue) => queue,
733 Err(_) => return RequestBalanceQueueDue::Missing,
734 };
735
736 match queue.get(key).copied() {
737 Some(due_at) if due_at <= now => {
738 queue.remove(key);
739 RequestBalanceQueueDue::Due
740 }
741 Some(due_at) => RequestBalanceQueueDue::NotDue(due_at.saturating_duration_since(now)),
742 None => RequestBalanceQueueDue::Missing,
743 }
744}
745
746#[cfg(test)]
747pub fn request_balance_refresh_queued_for_provider_endpoint(
748 provider_endpoint: &ProviderEndpointKey,
749) -> bool {
750 let Some(queue) = REQUEST_BALANCE_QUEUE.get() else {
751 return false;
752 };
753 match queue.lock() {
754 Ok(guard) => guard.contains_key(&RequestBalanceQueueKey::ProviderEndpoint(
755 provider_endpoint.clone(),
756 )),
757 Err(error) => error
758 .into_inner()
759 .contains_key(&RequestBalanceQueueKey::ProviderEndpoint(
760 provider_endpoint.clone(),
761 )),
762 }
763}
764
765fn usage_provider_target_for_legacy_upstream(
766 cfg: &ProxyConfig,
767 service_name: &str,
768 station_name: &str,
769 upstream_index: usize,
770) -> Option<UsageProviderTarget> {
771 let current_service = service_manager(cfg, service_name).station(station_name)?;
772 let current_upstream = current_service.upstreams.get(upstream_index)?;
773 Some(UsageProviderTarget {
774 upstream: UpstreamRef {
775 station_name: station_name.to_string(),
776 index: upstream_index,
777 provider_endpoint: current_upstream.provider_endpoint_key(service_name),
778 },
779 base_url: current_upstream.base_url.clone(),
780 provider_id: current_upstream.tags.get("provider_id").cloned(),
781 })
782}
783
784fn usage_provider_target_for_provider_endpoint(
785 cfg: &ProxyConfig,
786 service_name: &str,
787 provider_endpoint: &ProviderEndpointKey,
788) -> Option<UsageProviderTarget> {
789 service_manager(cfg, service_name)
790 .stations()
791 .iter()
792 .filter_map(|(station_name, service)| {
793 service
794 .upstreams
795 .iter()
796 .enumerate()
797 .find_map(|(index, upstream)| {
798 let upstream_endpoint = upstream.provider_endpoint_key(service_name)?;
799 if upstream_endpoint != *provider_endpoint {
800 return None;
801 }
802 Some(UsageProviderTarget {
803 upstream: UpstreamRef {
804 station_name: station_name.clone(),
805 index,
806 provider_endpoint: Some(upstream_endpoint),
807 },
808 base_url: upstream.base_url.clone(),
809 provider_id: upstream.tags.get("provider_id").cloned(),
810 })
811 })
812 })
813 .next()
814}
815
816trait UsageProviderUpstreamIdentityExt {
817 fn provider_endpoint_key(&self, service_name: &str) -> Option<ProviderEndpointKey>;
818}
819
820impl UsageProviderUpstreamIdentityExt for crate::config::UpstreamConfig {
821 fn provider_endpoint_key(&self, service_name: &str) -> Option<ProviderEndpointKey> {
822 let provider_id = self.tags.get("provider_id")?.trim();
823 let endpoint_id = self.tags.get("endpoint_id")?.trim();
824 if provider_id.is_empty() || endpoint_id.is_empty() {
825 return None;
826 }
827 Some(ProviderEndpointKey::new(
828 service_name.to_string(),
829 provider_id.to_string(),
830 endpoint_id.to_string(),
831 ))
832 }
833}
834
835#[cfg(test)]
836fn configured_target_keys(
837 cfg: &ProxyConfig,
838 service_name: &str,
839 providers: &[UsageProviderConfig],
840 station_name_filter: Option<&str>,
841) -> HashSet<UsageProviderTargetKey> {
842 providers
843 .iter()
844 .flat_map(|provider| {
845 matching_provider_targets(cfg, service_name, provider, station_name_filter)
846 })
847 .map(|target| target_key(&target))
848 .collect()
849}
850
851fn resolve_token(
852 provider: &UsageProviderConfig,
853 upstreams: &[UpstreamRef],
854 cfg: &ProxyConfig,
855 service_name: &str,
856) -> Option<String> {
857 if let Some(env_name) = &provider.token_env
859 && let Ok(v) = std::env::var(env_name)
860 && !v.trim().is_empty()
861 {
862 return Some(v);
863 }
864
865 if provider.require_token_env {
866 return None;
867 }
868
869 for uref in upstreams {
871 if let Some(service) = service_manager(cfg, service_name).station(&uref.station_name)
872 && let Some(up) = service.upstreams.get(uref.index)
873 {
874 if let Some(token) = up.auth.resolve_auth_token() {
875 return Some(token);
876 }
877 if let Some(token) = up.auth.resolve_api_key() {
878 return Some(token);
879 }
880 }
881 }
882 None
883}
884
885fn normalized_balance_base_url(base_url: &str) -> Option<String> {
886 let mut url = reqwest::Url::parse(base_url).ok()?;
887 url.set_query(None);
888 url.set_fragment(None);
889 let path = url.path().trim_end_matches('/').to_string();
890 if path.eq_ignore_ascii_case("/v1") {
891 url.set_path("");
892 } else if path.to_ascii_lowercase().ends_with("/v1") {
893 let new_path = &path[..path.len().saturating_sub(3)];
894 url.set_path(if new_path.is_empty() { "/" } else { new_path });
895 }
896 Some(url.as_str().trim_end_matches('/').to_string())
897}
898
899fn base_path_prefixes(base_url: &str) -> Vec<String> {
900 let Some(normalized) = normalized_balance_base_url(base_url) else {
901 return Vec::new();
902 };
903 let Ok(url) = reqwest::Url::parse(&normalized) else {
904 return Vec::new();
905 };
906 let parts = url
907 .path_segments()
908 .map(|segments| {
909 segments
910 .filter(|segment| !segment.is_empty())
911 .collect::<Vec<_>>()
912 })
913 .unwrap_or_default();
914 let mut prefixes = Vec::new();
915 for len in (1..=parts.len()).rev() {
916 prefixes.push(format!("/{}", parts[..len].join("/")));
917 }
918 if prefixes.is_empty() {
919 prefixes.push("/".to_string());
920 }
921 prefixes
922}
923
924fn path_prefixes_match(provider_prefixes: &[String], available_prefixes: &[String]) -> bool {
925 if provider_prefixes.is_empty() || available_prefixes.is_empty() {
926 return false;
927 }
928 provider_prefixes.iter().any(|provider_prefix| {
929 available_prefixes.iter().any(|available_prefix| {
930 provider_prefix == available_prefix
931 || provider_prefix
932 .strip_prefix(available_prefix)
933 .is_some_and(|suffix| suffix.starts_with('/'))
934 })
935 })
936}
937
938fn render_provider_template(
939 template: &str,
940 base_url: &str,
941 upstream_base_url: &str,
942 token: &str,
943 variables: &BTreeMap<String, String>,
944) -> String {
945 let mut out = template
946 .replace("{{baseUrl}}", base_url)
947 .replace("{{base_url}}", base_url)
948 .replace("{{upstreamBaseUrl}}", upstream_base_url)
949 .replace("{{upstream_base_url}}", upstream_base_url)
950 .replace("{{apiKey}}", token)
951 .replace("{{accessToken}}", token)
952 .replace("{{token}}", token);
953
954 out = out
955 .replace("{{unix_now}}", &unix_now_secs().to_string())
956 .replace("{{unix_now_ms}}", &unix_now_ms().to_string());
957
958 while let Some(start) = out.find("{{unix_days_ago:") {
959 let Some(end_offset) = out[start..].find("}}") else {
960 break;
961 };
962 let end = start + end_offset + 2;
963 let days_str = out[start + "{{unix_days_ago:".len()..end - 2].trim();
964 let replacement = days_str
965 .parse::<u64>()
966 .ok()
967 .map(|days| unix_now_secs().saturating_sub(days.saturating_mul(24 * 60 * 60)))
968 .map(|secs| secs.to_string())
969 .unwrap_or_default();
970 out.replace_range(start..end, &replacement);
971 }
972
973 while let Some(start) = out.find("{{env:") {
974 let Some(end_offset) = out[start..].find("}}") else {
975 break;
976 };
977 let end = start + end_offset + 2;
978 let env_name = out[start + 6..end - 2].trim();
979 let value = std::env::var(env_name).unwrap_or_default();
980 out.replace_range(start..end, &value);
981 }
982
983 for (name, value_template) in variables {
984 let value = render_provider_template(
985 value_template,
986 base_url,
987 upstream_base_url,
988 token,
989 &BTreeMap::new(),
990 );
991 out = out.replace(&format!("{{{{{name}}}}}"), &value);
992 }
993
994 out
995}
996
997fn resolve_endpoint(
998 provider: &UsageProviderConfig,
999 upstream_base_url: &str,
1000 token: &str,
1001) -> Result<String> {
1002 let base_url = normalized_balance_base_url(upstream_base_url)
1003 .ok_or_else(|| anyhow::anyhow!("invalid upstream base_url for balance endpoint"))?;
1004 let endpoint = if provider.endpoint.trim().is_empty() {
1005 provider
1006 .kind
1007 .default_endpoint()
1008 .unwrap_or_default()
1009 .to_string()
1010 } else {
1011 provider.endpoint.trim().to_string()
1012 };
1013 if endpoint.is_empty() {
1014 anyhow::bail!(
1015 "usage provider '{}' has no endpoint and kind {:?} has no default endpoint",
1016 provider.id,
1017 provider.kind
1018 );
1019 }
1020
1021 let rendered = render_provider_template(
1022 &endpoint,
1023 &base_url,
1024 upstream_base_url,
1025 token,
1026 &provider.variables,
1027 );
1028 if rendered.starts_with("http://") || rendered.starts_with("https://") {
1029 return Ok(rendered);
1030 }
1031
1032 let path = if rendered.starts_with('/') {
1033 rendered
1034 } else {
1035 format!("/{rendered}")
1036 };
1037 Ok(format!("{base_url}{path}"))
1038}
1039
1040fn endpoint_origin(endpoint: &str) -> String {
1041 reqwest::Url::parse(endpoint)
1042 .ok()
1043 .and_then(|url| {
1044 let host = url.host_str()?;
1045 let origin = match url.port() {
1046 Some(port) => format!("{}://{}:{}", url.scheme(), host, port),
1047 None => format!("{}://{}", url.scheme(), host),
1048 };
1049 Some(origin)
1050 })
1051 .unwrap_or_else(|| "unknown-origin".to_string())
1052}
1053
1054async fn poll_provider_http_json(
1055 client: &Client,
1056 provider: &UsageProviderConfig,
1057 upstream_base_url: &str,
1058 token: &str,
1059) -> Result<serde_json::Value> {
1060 let endpoint = resolve_endpoint(provider, upstream_base_url, token)?;
1061 let origin = endpoint_origin(&endpoint);
1062 let base_url = normalized_balance_base_url(upstream_base_url).unwrap_or_default();
1063 let mut req = client
1064 .get(endpoint)
1065 .timeout(BALANCE_HTTP_REQUEST_TIMEOUT)
1066 .header("Accept", "application/json")
1067 .header(
1068 "User-Agent",
1069 concat!("codex-helper/", env!("CARGO_PKG_VERSION")),
1070 );
1071
1072 match provider.kind {
1073 ProviderKind::YescodeProfile => {
1074 req = req.header("X-API-Key", token);
1075 }
1076 _ => {
1077 req = req.header("Authorization", format!("Bearer {}", token));
1078 }
1079 }
1080
1081 for (name, template) in &provider.headers {
1082 let value = render_provider_template(
1083 template,
1084 &base_url,
1085 upstream_base_url,
1086 token,
1087 &provider.variables,
1088 );
1089 if !value.trim().is_empty() {
1090 req = req.header(name.as_str(), value);
1091 }
1092 }
1093
1094 let resp = req.send().await.with_context(|| {
1095 format!(
1096 "usage provider request failed for {} via {:?}",
1097 origin, provider.kind
1098 )
1099 })?;
1100
1101 if !resp.status().is_success() {
1102 anyhow::bail!(
1103 "usage provider HTTP {} from {} via {:?}",
1104 resp.status(),
1105 origin,
1106 provider.kind
1107 );
1108 }
1109 let content_type = resp
1110 .headers()
1111 .get(reqwest::header::CONTENT_TYPE)
1112 .and_then(|value| value.to_str().ok())
1113 .map(str::to_string)
1114 .unwrap_or_else(|| "unknown".to_string());
1115 let text = resp.text().await.with_context(|| {
1116 format!(
1117 "usage provider response read failed from {} via {:?}",
1118 origin, provider.kind
1119 )
1120 })?;
1121 serde_json::from_str(&text).with_context(|| {
1122 format!(
1123 "usage provider returned non-JSON response from {} via {:?} (content-type {}, {} bytes)",
1124 origin,
1125 provider.kind,
1126 content_type,
1127 text.len()
1128 )
1129 })
1130}
1131
1132fn amount_from_json(value: &serde_json::Value) -> Option<UsdAmount> {
1133 let raw = match value {
1134 serde_json::Value::Number(number) => number.to_string(),
1135 serde_json::Value::String(text) => text.trim().to_string(),
1136 _ => return None,
1137 };
1138 UsdAmount::from_decimal_str(raw.as_str())
1139}
1140
1141fn decimal_string_from_json(value: &serde_json::Value) -> Option<String> {
1142 match value {
1143 serde_json::Value::Number(number) => Some(number.to_string()),
1144 serde_json::Value::String(text) => {
1145 let text = text.trim();
1146 if text.is_empty() {
1147 None
1148 } else {
1149 Some(text.to_string())
1150 }
1151 }
1152 _ => None,
1153 }
1154}
1155
1156fn amount_from_json_with_divisor(
1157 value: &serde_json::Value,
1158 divisor: Option<u64>,
1159) -> Option<UsdAmount> {
1160 let amount = amount_from_json(value)?;
1161 match divisor {
1162 Some(divisor) => amount.checked_div_u64(divisor),
1163 None => Some(amount),
1164 }
1165}
1166
1167fn json_value_at_path<'a>(
1168 value: &'a serde_json::Value,
1169 path: &str,
1170) -> Option<&'a serde_json::Value> {
1171 let mut current = value;
1172 for segment in path
1173 .split('.')
1174 .map(str::trim)
1175 .filter(|segment| !segment.is_empty())
1176 {
1177 current = match current {
1178 serde_json::Value::Array(items) => {
1179 let index = segment.parse::<usize>().ok()?;
1180 items.get(index)?
1181 }
1182 _ => current.get(segment)?,
1183 };
1184 }
1185 Some(current)
1186}
1187
1188fn first_amount_from_paths(
1189 value: &serde_json::Value,
1190 custom_paths: &[String],
1191 default_paths: &[&str],
1192 divisor: Option<u64>,
1193) -> Option<UsdAmount> {
1194 custom_paths
1195 .iter()
1196 .map(String::as_str)
1197 .chain(default_paths.iter().copied())
1198 .find_map(|path| {
1199 json_value_at_path(value, path)
1200 .and_then(|value| amount_from_json_with_divisor(value, divisor))
1201 })
1202}
1203
1204fn bool_from_json(value: &serde_json::Value) -> Option<bool> {
1205 match value {
1206 serde_json::Value::Bool(value) => Some(*value),
1207 serde_json::Value::Number(number) => number.as_i64().map(|value| value != 0),
1208 serde_json::Value::String(text) => match text.trim().to_ascii_lowercase().as_str() {
1209 "true" | "yes" | "1" | "exhausted" => Some(true),
1210 "false" | "no" | "0" | "ok" => Some(false),
1211 _ => None,
1212 },
1213 _ => None,
1214 }
1215}
1216
1217fn first_bool_from_paths(
1218 value: &serde_json::Value,
1219 custom_paths: &[String],
1220 default_paths: &[&str],
1221) -> Option<bool> {
1222 custom_paths
1223 .iter()
1224 .map(String::as_str)
1225 .chain(default_paths.iter().copied())
1226 .find_map(|path| json_value_at_path(value, path).and_then(bool_from_json))
1227}
1228
1229fn first_decimal_string_from_paths(
1230 value: &serde_json::Value,
1231 default_paths: &[&str],
1232) -> Option<String> {
1233 default_paths
1234 .iter()
1235 .copied()
1236 .find_map(|path| json_value_at_path(value, path).and_then(decimal_string_from_json))
1237}
1238
1239fn string_from_json(value: &serde_json::Value) -> Option<String> {
1240 match value {
1241 serde_json::Value::String(text) => {
1242 let text = text.trim();
1243 if text.is_empty() {
1244 None
1245 } else {
1246 Some(text.to_string())
1247 }
1248 }
1249 _ => None,
1250 }
1251}
1252
1253fn first_string_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<String> {
1254 default_paths
1255 .iter()
1256 .copied()
1257 .find_map(|path| json_value_at_path(value, path).and_then(string_from_json))
1258}
1259
1260fn u64_from_json(value: &serde_json::Value) -> Option<u64> {
1261 match value {
1262 serde_json::Value::Number(number) => number.as_u64(),
1263 serde_json::Value::String(text) => {
1264 let text = text.trim();
1265 if text.is_empty() {
1266 None
1267 } else {
1268 text.parse::<u64>().ok()
1269 }
1270 }
1271 _ => None,
1272 }
1273}
1274
1275fn seconds_from_json(value: &serde_json::Value) -> Option<u64> {
1276 match value {
1277 serde_json::Value::Number(number) => number.as_f64().map(|value| value.max(0.0) as u64),
1278 serde_json::Value::String(text) => {
1279 let text = text.trim();
1280 if text.is_empty() {
1281 None
1282 } else if let Ok(value) = text.parse::<f64>() {
1283 Some(value.max(0.0) as u64)
1284 } else {
1285 parse_timestamp_secs(text)
1286 }
1287 }
1288 _ => None,
1289 }
1290}
1291
1292fn first_secs_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<u64> {
1293 default_paths
1294 .iter()
1295 .copied()
1296 .find_map(|path| json_value_at_path(value, path).and_then(seconds_from_json))
1297}
1298
1299fn parse_timestamp_secs(value: &str) -> Option<u64> {
1300 parse_rfc3339_like_secs(value).or_else(|| {
1301 httpdate::parse_http_date(value).ok().and_then(|time| {
1302 time.duration_since(std::time::UNIX_EPOCH)
1303 .ok()
1304 .map(|duration| duration.as_secs())
1305 })
1306 })
1307}
1308
1309fn parse_rfc3339_like_secs(value: &str) -> Option<u64> {
1310 let value = value.trim();
1311 let datetime_sep = value.find('T').or_else(|| value.find(' '))?;
1312 let (datetime, offset_secs) = if let Some(datetime) = value.strip_suffix('Z') {
1313 (datetime, 0_i64)
1314 } else {
1315 let offset_pos = value[datetime_sep + 1..]
1316 .rfind(['+', '-'])
1317 .map(|pos| datetime_sep + 1 + pos)?;
1318 let (datetime, offset) = value.split_at(offset_pos);
1319 (datetime, parse_rfc3339_offset_secs(offset)?)
1320 };
1321
1322 let (date, time) = datetime.split_at(datetime_sep);
1323 let time = time.get(1..)?;
1324 let mut date_parts = date.split('-');
1325 let year = date_parts.next()?.parse::<i32>().ok()?;
1326 let month = date_parts.next()?.parse::<u32>().ok()?;
1327 let day = date_parts.next()?.parse::<u32>().ok()?;
1328 if date_parts.next().is_some() {
1329 return None;
1330 }
1331
1332 let mut time_parts = time.split(':');
1333 let hour = time_parts.next()?.parse::<u32>().ok()?;
1334 let minute = time_parts.next()?.parse::<u32>().ok()?;
1335 let second_raw = time_parts.next().unwrap_or("0");
1336 if time_parts.next().is_some() {
1337 return None;
1338 }
1339 let second = second_raw
1340 .split('.')
1341 .next()
1342 .and_then(|value| value.parse::<u32>().ok())?;
1343 if !(1..=12).contains(&month) || day == 0 || hour > 23 || minute > 59 || second > 60 {
1344 return None;
1345 }
1346
1347 let local_secs = days_from_civil(year, month, day)
1348 .checked_mul(86_400)?
1349 .checked_add(i64::from(hour) * 3_600 + i64::from(minute) * 60 + i64::from(second))?;
1350 local_secs
1351 .checked_sub(offset_secs)
1352 .and_then(|utc_secs| u64::try_from(utc_secs).ok())
1353}
1354
1355fn parse_rfc3339_offset_secs(offset: &str) -> Option<i64> {
1356 let sign = match offset.as_bytes().first().copied()? {
1357 b'+' => 1_i64,
1358 b'-' => -1_i64,
1359 _ => return None,
1360 };
1361 let raw = offset.get(1..)?;
1362 let (hours, minutes) = raw
1363 .split_once(':')
1364 .unwrap_or_else(|| raw.split_at(raw.len().min(2)));
1365 let hours = hours.parse::<i64>().ok()?;
1366 let minutes = minutes.parse::<i64>().ok()?;
1367 if hours > 23 || minutes > 59 {
1368 return None;
1369 }
1370 Some(sign * (hours * 3_600 + minutes * 60))
1371}
1372
1373fn days_from_civil(year: i32, month: u32, day: u32) -> i64 {
1374 let year = i64::from(year) - if month <= 2 { 1 } else { 0 };
1375 let era = if year >= 0 { year } else { year - 399 } / 400;
1376 let yoe = year - era * 400;
1377 let month = i64::from(month);
1378 let doy = (153 * (month + if month > 2 { -3 } else { 9 }) + 2) / 5 + i64::from(day) - 1;
1379 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
1380 era * 146_097 + doe - 719_468
1381}
1382
1383fn first_u64_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<u64> {
1384 default_paths
1385 .iter()
1386 .copied()
1387 .find_map(|path| json_value_at_path(value, path).and_then(u64_from_json))
1388}
1389
1390fn array_from_json_path<'a>(
1391 value: &'a serde_json::Value,
1392 path: &str,
1393) -> Option<&'a Vec<serde_json::Value>> {
1394 json_value_at_path(value, path).and_then(|value| value.as_array())
1395}
1396
1397fn amount_to_string(amount: UsdAmount) -> String {
1398 amount.format_usd()
1399}
1400
1401#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1402struct QuotaWindowSnapshot {
1403 period: &'static str,
1404 remaining: UsdAmount,
1405 used: UsdAmount,
1406 limit: UsdAmount,
1407}
1408
1409fn base_snapshot(
1410 provider: &UsageProviderConfig,
1411 upstream: &UpstreamRef,
1412 fetched_at_ms: u64,
1413 stale_after_ms: Option<u64>,
1414) -> ProviderBalanceSnapshot {
1415 let mut snapshot = ProviderBalanceSnapshot::new(
1416 provider.id.clone(),
1417 upstream.station_name.clone(),
1418 upstream.index,
1419 provider.kind.source_name(),
1420 fetched_at_ms,
1421 stale_after_ms,
1422 );
1423 snapshot.exhaustion_affects_routing = provider.trust_exhaustion_for_routing;
1424 snapshot
1425}
1426
1427fn snapshot_error(
1428 provider: &UsageProviderConfig,
1429 upstream: &UpstreamRef,
1430 fetched_at_ms: u64,
1431 stale_after_ms: Option<u64>,
1432 message: impl Into<String>,
1433) -> ProviderBalanceSnapshot {
1434 base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms).with_error(message)
1435}
1436
1437fn budget_snapshot_from_json(
1438 provider: &UsageProviderConfig,
1439 upstream: &UpstreamRef,
1440 value: &serde_json::Value,
1441 fetched_at_ms: u64,
1442 stale_after_ms: Option<u64>,
1443) -> ProviderBalanceSnapshot {
1444 let monthly_budget = first_amount_from_paths(
1445 value,
1446 &provider.extract.monthly_budget_paths,
1447 &["monthly_budget_usd", "data.monthly_budget_usd"],
1448 provider.extract.monthly_budget_divisor,
1449 );
1450 let monthly_spent = first_amount_from_paths(
1451 value,
1452 &provider.extract.monthly_spent_paths,
1453 &["monthly_spent_usd", "data.monthly_spent_usd"],
1454 provider.extract.monthly_spent_divisor,
1455 );
1456 let exhausted = match (monthly_budget, monthly_spent) {
1457 (Some(budget), Some(spent)) if !budget.is_zero() => Some(spent >= budget),
1458 (Some(_), Some(_)) => Some(false),
1459 _ => None,
1460 };
1461
1462 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1463 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
1464 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
1465 snapshot.exhausted = exhausted;
1466 snapshot.refresh_status(fetched_at_ms);
1467 snapshot
1468}
1469
1470fn yescode_snapshot_from_json(
1471 provider: &UsageProviderConfig,
1472 upstream: &UpstreamRef,
1473 value: &serde_json::Value,
1474 fetched_at_ms: u64,
1475 stale_after_ms: Option<u64>,
1476) -> ProviderBalanceSnapshot {
1477 let subscription_balance = first_amount_from_paths(
1478 value,
1479 &provider.extract.subscription_balance_paths,
1480 &["subscription_balance", "data.subscription_balance"],
1481 provider.extract.remaining_divisor,
1482 );
1483 let paygo_balance = first_amount_from_paths(
1484 value,
1485 &provider.extract.paygo_balance_paths,
1486 &[
1487 "pay_as_you_go_balance",
1488 "paygo_balance",
1489 "data.pay_as_you_go_balance",
1490 "data.paygo_balance",
1491 ],
1492 provider.extract.remaining_divisor,
1493 );
1494 let total_balance = match (subscription_balance, paygo_balance) {
1495 (Some(subscription), Some(paygo)) => Some(subscription.saturating_add(paygo)),
1496 (Some(subscription), None) => Some(subscription),
1497 (None, Some(paygo)) => Some(paygo),
1498 (None, None) => None,
1499 };
1500
1501 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1502 snapshot.total_balance_usd = total_balance.map(amount_to_string);
1503 snapshot.subscription_balance_usd = subscription_balance.map(amount_to_string);
1504 snapshot.paygo_balance_usd = paygo_balance.map(amount_to_string);
1505 snapshot.exhausted = total_balance.map(UsdAmount::is_zero);
1506 snapshot.refresh_status(fetched_at_ms);
1507 snapshot
1508}
1509
1510fn balance_http_snapshot_from_json(
1511 provider: &UsageProviderConfig,
1512 upstream: &UpstreamRef,
1513 value: &serde_json::Value,
1514 fetched_at_ms: u64,
1515 stale_after_ms: Option<u64>,
1516) -> ProviderBalanceSnapshot {
1517 let remaining_balance = first_amount_from_paths(
1518 value,
1519 &provider.extract.remaining_balance_paths,
1520 &[
1521 "balance",
1522 "remaining",
1523 "remain",
1524 "available",
1525 "available_balance",
1526 "credit",
1527 "credits",
1528 "total_balance",
1529 "total_balance_usd",
1530 "totalBalance",
1531 "availableBalance",
1532 "available_balance_usd",
1533 "balance_infos.0.total_balance",
1534 "data.balance",
1535 "data.remaining",
1536 "data.available",
1537 "data.available_balance",
1538 "data.credit",
1539 "data.credits",
1540 "data.total_balance",
1541 "data.totalBalance",
1542 "data.availableBalance",
1543 ],
1544 provider.extract.remaining_divisor,
1545 );
1546 let subscription_balance = first_amount_from_paths(
1547 value,
1548 &provider.extract.subscription_balance_paths,
1549 &[
1550 "subscription_balance",
1551 "subscription_balance_usd",
1552 "subscriptionBalance",
1553 "data.subscription_balance",
1554 "data.subscription_balance_usd",
1555 "data.subscriptionBalance",
1556 ],
1557 provider.extract.remaining_divisor,
1558 );
1559 let paygo_balance = first_amount_from_paths(
1560 value,
1561 &provider.extract.paygo_balance_paths,
1562 &[
1563 "pay_as_you_go_balance",
1564 "paygo_balance",
1565 "paygo",
1566 "paygoBalance",
1567 "chargeBalance",
1568 "voucherBalance",
1569 "data.pay_as_you_go_balance",
1570 "data.paygo_balance",
1571 "data.paygo",
1572 "data.paygoBalance",
1573 "data.chargeBalance",
1574 "data.voucherBalance",
1575 ],
1576 provider.extract.remaining_divisor,
1577 );
1578 let component_remaining = match (subscription_balance, paygo_balance) {
1579 (Some(subscription), Some(paygo)) => Some(subscription.saturating_add(paygo)),
1580 (Some(subscription), None) => Some(subscription),
1581 (None, Some(paygo)) => Some(paygo),
1582 (None, None) => None,
1583 };
1584 let monthly_spent = first_amount_from_paths(
1585 value,
1586 &provider.extract.monthly_spent_paths,
1587 &[
1588 "monthly_spent_usd",
1589 "spent",
1590 "used",
1591 "used_balance",
1592 "usedBalance",
1593 "total_usage",
1594 "data.monthly_spent_usd",
1595 "data.spent",
1596 "data.used",
1597 "data.used_balance",
1598 "data.usedBalance",
1599 "data.total_usage",
1600 ],
1601 provider.extract.monthly_spent_divisor,
1602 );
1603 let monthly_budget = first_amount_from_paths(
1604 value,
1605 &provider.extract.monthly_budget_paths,
1606 &[
1607 "monthly_budget_usd",
1608 "budget",
1609 "limit",
1610 "quota_total",
1611 "creditLimit",
1612 "total_credits",
1613 "data.monthly_budget_usd",
1614 "data.budget",
1615 "data.limit",
1616 "data.quota_total",
1617 "data.creditLimit",
1618 "data.total_credits",
1619 ],
1620 provider.extract.monthly_budget_divisor,
1621 )
1622 .or_else(|| {
1623 if provider.extract.derive_budget_from_remaining_and_spent {
1624 match (remaining_balance.or(component_remaining), monthly_spent) {
1625 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
1626 _ => None,
1627 }
1628 } else {
1629 None
1630 }
1631 });
1632 let total_balance = remaining_balance.or(component_remaining).or_else(|| {
1633 match (
1634 provider.extract.derive_remaining_from_budget_and_spent,
1635 monthly_budget,
1636 monthly_spent,
1637 ) {
1638 (true, Some(budget), Some(spent)) => Some(budget.saturating_sub(spent)),
1639 _ => None,
1640 }
1641 });
1642 let exhausted = first_bool_from_paths(
1643 value,
1644 &provider.extract.exhausted_paths,
1645 &[
1646 "exhausted",
1647 "quota_exhausted",
1648 "balance_exhausted",
1649 "data.exhausted",
1650 "data.quota_exhausted",
1651 "data.balance_exhausted",
1652 ],
1653 )
1654 .or_else(|| total_balance.map(UsdAmount::is_zero));
1655
1656 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1657 snapshot.total_balance_usd = total_balance.map(amount_to_string);
1658 snapshot.subscription_balance_usd = subscription_balance.map(amount_to_string);
1659 snapshot.paygo_balance_usd = paygo_balance.map(amount_to_string);
1660 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
1661 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
1662 snapshot.exhausted = exhausted;
1663 snapshot.refresh_status(fetched_at_ms);
1664 snapshot
1665}
1666
1667fn has_any_json_path(value: &serde_json::Value, paths: &[&str]) -> bool {
1668 paths
1669 .iter()
1670 .any(|path| json_value_at_path(value, path).is_some())
1671}
1672
1673fn populate_sub2api_usage_fields(
1674 snapshot: &mut ProviderBalanceSnapshot,
1675 value: &serde_json::Value,
1676) {
1677 snapshot.plan_name = first_string_from_paths(
1678 value,
1679 &["planName", "plan_name", "data.planName", "data.plan_name"],
1680 );
1681 let remaining_balance = sub2api_remaining_balance(value);
1682 snapshot.total_balance_usd = snapshot
1683 .total_balance_usd
1684 .take()
1685 .or_else(|| remaining_balance.map(amount_to_string));
1686 snapshot.total_used_usd = first_amount_from_paths(
1687 value,
1688 &[],
1689 &[
1690 "usage.total.total_cost_usd",
1691 "usage.total.total_cost",
1692 "usage.total.cost",
1693 "data.usage.total.total_cost_usd",
1694 "data.usage.total.total_cost",
1695 "data.usage.total.cost",
1696 ],
1697 None,
1698 )
1699 .map(amount_to_string);
1700 snapshot.today_used_usd = first_amount_from_paths(
1701 value,
1702 &[],
1703 &[
1704 "usage.today.total_cost_usd",
1705 "usage.today.total_cost",
1706 "usage.today.cost",
1707 "data.usage.today.total_cost_usd",
1708 "data.usage.today.total_cost",
1709 "data.usage.today.cost",
1710 ],
1711 None,
1712 )
1713 .map(amount_to_string);
1714 snapshot.total_requests = first_u64_from_paths(
1715 value,
1716 &[
1717 "usage.total.request_count",
1718 "usage.total.requests",
1719 "usage.total.count",
1720 "data.usage.total.request_count",
1721 "data.usage.total.requests",
1722 "data.usage.total.count",
1723 ],
1724 );
1725 snapshot.today_requests = first_u64_from_paths(
1726 value,
1727 &[
1728 "usage.today.request_count",
1729 "usage.today.requests",
1730 "usage.today.count",
1731 "data.usage.today.request_count",
1732 "data.usage.today.requests",
1733 "data.usage.today.count",
1734 ],
1735 );
1736 snapshot.total_tokens = first_u64_from_paths(
1737 value,
1738 &[
1739 "usage.total.total_tokens",
1740 "usage.total.tokens",
1741 "usage.total.input_tokens",
1742 "usage.total.prompt_tokens",
1743 "data.usage.total.total_tokens",
1744 "data.usage.total.tokens",
1745 ],
1746 );
1747 snapshot.today_tokens = first_u64_from_paths(
1748 value,
1749 &[
1750 "usage.today.total_tokens",
1751 "usage.today.tokens",
1752 "usage.today.input_tokens",
1753 "usage.today.prompt_tokens",
1754 "data.usage.today.total_tokens",
1755 "data.usage.today.tokens",
1756 ],
1757 );
1758 snapshot.usage_rate = sub2api_usage_rate(value);
1759 snapshot.usage_windows = sub2api_usage_windows(value);
1760 snapshot.usage_model_stats = sub2api_model_stats(value);
1761 snapshot.subscription_expires_at = first_string_from_paths(
1762 value,
1763 &[
1764 "subscription.expires_at",
1765 "data.subscription.expires_at",
1766 "subscription.expiresAt",
1767 "data.subscription.expiresAt",
1768 ],
1769 );
1770 snapshot.usage_alerts = sub2api_usage_alerts(value);
1771}
1772
1773fn sub2api_remaining_balance(value: &serde_json::Value) -> Option<UsdAmount> {
1774 let remaining = first_amount_from_paths(value, &[], &["remaining", "data.remaining"], None)?;
1775 if sub2api_has_subscription_windows(value)
1776 && sub2api_window_remaining_amounts(value).contains(&remaining)
1777 {
1778 return None;
1779 }
1780 Some(remaining)
1781}
1782
1783fn sub2api_has_subscription_windows(value: &serde_json::Value) -> bool {
1784 has_any_json_path(
1785 value,
1786 &[
1787 "subscription.daily_usage_usd",
1788 "subscription.daily_limit_usd",
1789 "subscription.weekly_usage_usd",
1790 "subscription.weekly_limit_usd",
1791 "subscription.monthly_usage_usd",
1792 "subscription.monthly_limit_usd",
1793 "data.subscription.daily_usage_usd",
1794 "data.subscription.daily_limit_usd",
1795 "data.subscription.weekly_usage_usd",
1796 "data.subscription.weekly_limit_usd",
1797 "data.subscription.monthly_usage_usd",
1798 "data.subscription.monthly_limit_usd",
1799 ],
1800 )
1801}
1802
1803fn sub2api_window_remaining_amounts(value: &serde_json::Value) -> Vec<UsdAmount> {
1804 ["daily", "weekly", "monthly"]
1805 .into_iter()
1806 .filter_map(|period| {
1807 let used = first_amount_from_paths(
1808 value,
1809 &[],
1810 &[
1811 &format!("subscription.{period}_usage_usd"),
1812 &format!("data.subscription.{period}_usage_usd"),
1813 ],
1814 None,
1815 );
1816 let limit = first_amount_from_paths(
1817 value,
1818 &[],
1819 &[
1820 &format!("subscription.{period}_limit_usd"),
1821 &format!("data.subscription.{period}_limit_usd"),
1822 ],
1823 None,
1824 );
1825 match (limit, used) {
1826 (Some(limit), Some(used)) if !limit.is_zero() => Some(limit.saturating_sub(used)),
1827 _ => None,
1828 }
1829 })
1830 .collect()
1831}
1832
1833fn sub2api_usage_rate(value: &serde_json::Value) -> Option<ProviderUsageRateSnapshot> {
1834 let rate = ProviderUsageRateSnapshot {
1835 average_duration_ms: first_decimal_string_from_paths(
1836 value,
1837 &[
1838 "usage.average_duration_ms",
1839 "data.usage.average_duration_ms",
1840 "average_duration_ms",
1841 "data.average_duration_ms",
1842 ],
1843 ),
1844 rpm: first_decimal_string_from_paths(value, &["usage.rpm", "data.usage.rpm", "rpm"]),
1845 tpm: first_decimal_string_from_paths(value, &["usage.tpm", "data.usage.tpm", "tpm"]),
1846 };
1847 (!rate.is_empty()).then_some(rate)
1848}
1849
1850fn sub2api_usage_windows(value: &serde_json::Value) -> Vec<ProviderUsageWindow> {
1851 ["daily", "weekly", "monthly"]
1852 .into_iter()
1853 .filter_map(|period| {
1854 let used = first_amount_from_paths(
1855 value,
1856 &[],
1857 &[
1858 &format!("subscription.{period}_usage_usd"),
1859 &format!("data.subscription.{period}_usage_usd"),
1860 ],
1861 None,
1862 );
1863 let limit = first_amount_from_paths(
1864 value,
1865 &[],
1866 &[
1867 &format!("subscription.{period}_limit_usd"),
1868 &format!("data.subscription.{period}_limit_usd"),
1869 ],
1870 None,
1871 );
1872 if used.is_none() && limit.is_none() {
1873 return None;
1874 }
1875 let unlimited = limit.map(|limit| limit.is_zero());
1876 let remaining = match (limit, used) {
1877 (Some(limit), Some(used)) if !limit.is_zero() => Some(limit.saturating_sub(used)),
1878 _ => None,
1879 };
1880 Some(ProviderUsageWindow {
1881 period: period.to_string(),
1882 used_usd: used.map(amount_to_string),
1883 limit_usd: limit.map(amount_to_string),
1884 remaining_usd: remaining.map(amount_to_string),
1885 unlimited,
1886 })
1887 })
1888 .collect()
1889}
1890
1891fn sub2api_model_stats(value: &serde_json::Value) -> Vec<ProviderUsageModelStat> {
1892 [
1893 "model_stats",
1894 "data.model_stats",
1895 "modelStats",
1896 "data.modelStats",
1897 ]
1898 .into_iter()
1899 .find_map(|path| array_from_json_path(value, path))
1900 .map(|items| {
1901 items
1902 .iter()
1903 .filter_map(sub2api_model_stat_from_json)
1904 .collect::<Vec<_>>()
1905 })
1906 .unwrap_or_default()
1907}
1908
1909fn sub2api_model_stat_from_json(value: &serde_json::Value) -> Option<ProviderUsageModelStat> {
1910 let model = first_string_from_paths(value, &["model", "model_name", "name"])?;
1911 let input_cost = first_amount_from_paths(value, &[], &["input_cost_usd", "input_cost"], None);
1912 let output_cost =
1913 first_amount_from_paths(value, &[], &["output_cost_usd", "output_cost"], None);
1914 let total_cost =
1915 first_amount_from_paths(value, &[], &["total_cost_usd", "total_cost", "cost"], None)
1916 .or_else(|| match (input_cost, output_cost) {
1917 (Some(input), Some(output)) => Some(input.saturating_add(output)),
1918 _ => None,
1919 });
1920 let input_tokens = first_u64_from_paths(value, &["input_tokens", "prompt_tokens"]);
1921 let output_tokens = first_u64_from_paths(value, &["output_tokens", "completion_tokens"]);
1922 let total_tokens =
1923 first_u64_from_paths(value, &["total_tokens", "tokens"]).or_else(|| {
1924 match (input_tokens, output_tokens) {
1925 (Some(input), Some(output)) => input.checked_add(output),
1926 _ => None,
1927 }
1928 });
1929 Some(ProviderUsageModelStat {
1930 model,
1931 request_count: first_u64_from_paths(value, &["request_count", "requests", "count"]),
1932 input_tokens,
1933 output_tokens,
1934 total_tokens,
1935 input_cost_usd: input_cost.map(amount_to_string),
1936 output_cost_usd: output_cost.map(amount_to_string),
1937 total_cost_usd: total_cost.map(amount_to_string),
1938 })
1939}
1940
1941fn sub2api_usage_alerts(value: &serde_json::Value) -> Vec<ProviderUsageAlert> {
1942 let mut alerts = Vec::new();
1943 if let (Some(used), Some(limit)) = (
1944 first_amount_from_paths(
1945 value,
1946 &[],
1947 &[
1948 "subscription.daily_usage_usd",
1949 "data.subscription.daily_usage_usd",
1950 ],
1951 None,
1952 ),
1953 first_amount_from_paths(
1954 value,
1955 &[],
1956 &[
1957 "subscription.daily_limit_usd",
1958 "data.subscription.daily_limit_usd",
1959 ],
1960 None,
1961 ),
1962 ) && !limit.is_zero()
1963 {
1964 let used_femto = used.femto_usd();
1965 let limit_femto = limit.femto_usd();
1966 if used_femto.saturating_mul(100) >= limit_femto.saturating_mul(95) {
1967 alerts.push(ProviderUsageAlert {
1968 kind: ProviderUsageAlertKind::DailyUsage95,
1969 message: "daily usage is at or above 95%".to_string(),
1970 });
1971 } else if used_femto.saturating_mul(100) >= limit_femto.saturating_mul(80) {
1972 alerts.push(ProviderUsageAlert {
1973 kind: ProviderUsageAlertKind::DailyUsage80,
1974 message: "daily usage is at or above 80%".to_string(),
1975 });
1976 }
1977 }
1978
1979 if let Some(remaining) = sub2api_remaining_balance(value)
1980 && let Some(threshold) = UsdAmount::from_decimal_str(LOW_BALANCE_ALERT_THRESHOLD_USD)
1981 && remaining <= threshold
1982 {
1983 alerts.push(ProviderUsageAlert {
1984 kind: ProviderUsageAlertKind::LowBalance,
1985 message: "remaining balance is low".to_string(),
1986 });
1987 }
1988
1989 if let Some(expires_at_secs) = first_secs_from_paths(
1990 value,
1991 &[
1992 "subscription.expires_at",
1993 "data.subscription.expires_at",
1994 "subscription.expiresAt",
1995 "data.subscription.expiresAt",
1996 ],
1997 ) {
1998 let now = unix_now_secs();
1999 if expires_at_secs <= now {
2000 alerts.push(ProviderUsageAlert {
2001 kind: ProviderUsageAlertKind::SubscriptionExpired,
2002 message: "subscription has expired".to_string(),
2003 });
2004 } else if expires_at_secs <= now.saturating_add(EXPIRING_SOON_WINDOW_SECS) {
2005 alerts.push(ProviderUsageAlert {
2006 kind: ProviderUsageAlertKind::SubscriptionExpiringSoon,
2007 message: "subscription expires within 7 days".to_string(),
2008 });
2009 }
2010 }
2011
2012 alerts.sort_by_key(|alert| alert.kind);
2013 alerts.dedup_by_key(|alert| alert.kind);
2014 alerts
2015}
2016
2017fn sub2api_subscription_limit_snapshot(
2018 value: &serde_json::Value,
2019 period: &'static str,
2020 limit_paths: &[&str],
2021 usage_paths: &[&str],
2022) -> Option<QuotaWindowSnapshot> {
2023 let budget = first_amount_from_paths(value, &[], limit_paths, None)?;
2024 if budget.is_zero() {
2025 return None;
2026 }
2027 let spent = first_amount_from_paths(value, &[], usage_paths, None).unwrap_or(UsdAmount::ZERO);
2028 let remaining = budget.saturating_sub(spent);
2029 Some(QuotaWindowSnapshot {
2030 period,
2031 remaining,
2032 used: spent,
2033 limit: budget,
2034 })
2035}
2036
2037fn sub2api_limiting_subscription_window(value: &serde_json::Value) -> Option<QuotaWindowSnapshot> {
2038 let windows = [
2039 sub2api_subscription_limit_snapshot(
2040 value,
2041 "daily",
2042 &[
2043 "subscription.daily_limit_usd",
2044 "data.subscription.daily_limit_usd",
2045 ],
2046 &[
2047 "subscription.daily_usage_usd",
2048 "data.subscription.daily_usage_usd",
2049 ],
2050 ),
2051 sub2api_subscription_limit_snapshot(
2052 value,
2053 "weekly",
2054 &[
2055 "subscription.weekly_limit_usd",
2056 "data.subscription.weekly_limit_usd",
2057 ],
2058 &[
2059 "subscription.weekly_usage_usd",
2060 "data.subscription.weekly_usage_usd",
2061 ],
2062 ),
2063 sub2api_subscription_limit_snapshot(
2064 value,
2065 "monthly",
2066 &[
2067 "subscription.monthly_limit_usd",
2068 "data.subscription.monthly_limit_usd",
2069 ],
2070 &[
2071 "subscription.monthly_usage_usd",
2072 "data.subscription.monthly_usage_usd",
2073 ],
2074 ),
2075 ];
2076
2077 windows
2078 .into_iter()
2079 .flatten()
2080 .min_by_key(|window| window.remaining)
2081}
2082
2083fn sub2api_usage_snapshot_from_json(
2084 provider: &UsageProviderConfig,
2085 upstream: &UpstreamRef,
2086 value: &serde_json::Value,
2087 fetched_at_ms: u64,
2088 stale_after_ms: Option<u64>,
2089) -> ProviderBalanceSnapshot {
2090 if json_value_at_path(value, "isValid").and_then(bool_from_json) == Some(false) {
2091 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2092 .with_error("sub2api usage response reported invalid API key");
2093 }
2094
2095 let mode = first_string_from_paths(value, &["mode", "data.mode"]);
2096 let has_subscription = has_any_json_path(value, &["subscription", "data.subscription"]);
2097
2098 if mode.as_deref() == Some("quota_limited") {
2099 let quota_remaining = first_amount_from_paths(
2100 value,
2101 &provider.extract.remaining_balance_paths,
2102 &[
2103 "quota.remaining",
2104 "data.quota.remaining",
2105 "remaining",
2106 "data.remaining",
2107 ],
2108 provider.extract.remaining_divisor,
2109 );
2110 let quota_limit = first_amount_from_paths(
2111 value,
2112 &provider.extract.monthly_budget_paths,
2113 &["quota.limit", "data.quota.limit"],
2114 provider.extract.monthly_budget_divisor,
2115 );
2116 let quota_used = first_amount_from_paths(
2117 value,
2118 &provider.extract.monthly_spent_paths,
2119 &["quota.used", "data.quota.used"],
2120 provider.extract.monthly_spent_divisor,
2121 );
2122 let exhausted = first_bool_from_paths(
2123 value,
2124 &provider.extract.exhausted_paths,
2125 &[
2126 "exhausted",
2127 "data.exhausted",
2128 "quota_exhausted",
2129 "data.quota_exhausted",
2130 ],
2131 )
2132 .or_else(|| quota_remaining.map(UsdAmount::is_zero));
2133
2134 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2135 snapshot.quota_period = Some("quota".to_string());
2136 snapshot.quota_remaining_usd = quota_remaining.map(amount_to_string);
2137 snapshot.quota_limit_usd = quota_limit.map(amount_to_string);
2138 snapshot.quota_used_usd = quota_used.map(amount_to_string);
2139 snapshot.monthly_budget_usd = quota_limit.map(amount_to_string);
2140 snapshot.monthly_spent_usd = quota_used.map(amount_to_string);
2141 snapshot.exhausted = exhausted;
2142 populate_sub2api_usage_fields(&mut snapshot, value);
2143 snapshot.refresh_status(fetched_at_ms);
2144 return snapshot;
2145 }
2146
2147 if mode.as_deref() == Some("unrestricted") && has_subscription {
2148 let limiting_window = sub2api_limiting_subscription_window(value);
2149 let exhausted = first_bool_from_paths(
2150 value,
2151 &provider.extract.exhausted_paths,
2152 &[
2153 "exhausted",
2154 "data.exhausted",
2155 "quota_exhausted",
2156 "data.quota_exhausted",
2157 ],
2158 )
2159 .or_else(|| limiting_window.map(|window| window.remaining.is_zero()))
2160 .or(Some(false));
2161
2162 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2163 if let Some(window) = limiting_window {
2164 snapshot.quota_period = Some(window.period.to_string());
2165 snapshot.quota_remaining_usd = Some(amount_to_string(window.remaining));
2166 snapshot.quota_limit_usd = Some(amount_to_string(window.limit));
2167 snapshot.quota_used_usd = Some(amount_to_string(window.used));
2168 snapshot.monthly_budget_usd = Some(amount_to_string(window.limit));
2169 snapshot.monthly_spent_usd = Some(amount_to_string(window.used));
2170 }
2171 snapshot.exhaustion_affects_routing = false;
2172 snapshot.exhausted = exhausted;
2173 populate_sub2api_usage_fields(&mut snapshot, value);
2174 snapshot.refresh_status(fetched_at_ms);
2175 return snapshot;
2176 }
2177
2178 let mut snapshot =
2179 balance_http_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms);
2180 populate_sub2api_usage_fields(&mut snapshot, value);
2181 snapshot.refresh_status(fetched_at_ms);
2182 snapshot
2183}
2184
2185fn sub2api_auth_me_snapshot_from_json(
2186 provider: &UsageProviderConfig,
2187 upstream: &UpstreamRef,
2188 value: &serde_json::Value,
2189 fetched_at_ms: u64,
2190 stale_after_ms: Option<u64>,
2191) -> ProviderBalanceSnapshot {
2192 if json_value_at_path(value, "code")
2193 .and_then(|value| value.as_i64())
2194 .is_some_and(|code| code != 0)
2195 {
2196 let message = json_value_at_path(value, "message")
2197 .and_then(|value| value.as_str())
2198 .unwrap_or("sub2api auth/me response reported failure");
2199 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2200 .with_error(message.to_string());
2201 }
2202
2203 balance_http_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
2204}
2205
2206fn rightcode_available_prefixes(value: &serde_json::Value) -> Vec<String> {
2207 array_from_json_path(value, "available_prefixes")
2208 .into_iter()
2209 .flatten()
2210 .filter_map(|value| value.as_str())
2211 .map(str::trim)
2212 .filter(|value| !value.is_empty())
2213 .map(ToOwned::to_owned)
2214 .collect()
2215}
2216
2217fn rightcode_subscription_window(value: &serde_json::Value) -> Option<QuotaWindowSnapshot> {
2218 let limit = json_value_at_path(value, "total_quota").and_then(amount_from_json)?;
2219 if limit < UsdAmount::from_decimal_str("10").unwrap_or(UsdAmount::ZERO) {
2220 return None;
2221 }
2222 let raw_remaining = json_value_at_path(value, "remaining_quota").and_then(amount_from_json)?;
2223 let reset_today = json_value_at_path(value, "reset_today").and_then(bool_from_json);
2224 let remaining = if reset_today == Some(true) {
2225 raw_remaining
2226 } else {
2227 raw_remaining.saturating_add(limit)
2228 };
2229 let used = limit.saturating_sub(remaining);
2230 Some(QuotaWindowSnapshot {
2231 period: "daily",
2232 remaining,
2233 used,
2234 limit,
2235 })
2236}
2237
2238fn rightcode_account_summary_snapshot_from_json(
2239 provider: &UsageProviderConfig,
2240 upstream: &UpstreamRef,
2241 value: &serde_json::Value,
2242 upstream_base_url: &str,
2243 fetched_at_ms: u64,
2244 stale_after_ms: Option<u64>,
2245) -> ProviderBalanceSnapshot {
2246 let balance = json_value_at_path(value, "balance").and_then(amount_from_json);
2247 let provider_prefixes = base_path_prefixes(upstream_base_url);
2248 let mut matched_windows = Vec::new();
2249 let mut matched_plan_names = Vec::new();
2250
2251 if let Some(subscriptions) = array_from_json_path(value, "subscriptions") {
2252 for subscription in subscriptions {
2253 let available_prefixes = rightcode_available_prefixes(subscription);
2254 if !path_prefixes_match(&provider_prefixes, &available_prefixes) {
2255 continue;
2256 }
2257 let Some(window) = rightcode_subscription_window(subscription) else {
2258 continue;
2259 };
2260 matched_windows.push(window);
2261 if let Some(name) = json_value_at_path(subscription, "name").and_then(string_from_json)
2262 {
2263 matched_plan_names.push(name);
2264 }
2265 }
2266 }
2267
2268 if balance.is_none() && matched_windows.is_empty() {
2269 return snapshot_error(
2270 provider,
2271 upstream,
2272 fetched_at_ms,
2273 stale_after_ms,
2274 "rightcode account summary missing balance and matching subscription quota fields",
2275 );
2276 }
2277
2278 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2279 snapshot.total_balance_usd = balance.map(amount_to_string);
2280
2281 if !matched_windows.is_empty() {
2282 let mut remaining = UsdAmount::ZERO;
2283 let mut used = UsdAmount::ZERO;
2284 let mut limit = UsdAmount::ZERO;
2285 for window in matched_windows {
2286 remaining = remaining.saturating_add(window.remaining);
2287 used = used.saturating_add(window.used);
2288 limit = limit.saturating_add(window.limit);
2289 }
2290 snapshot.quota_period = Some("daily".to_string());
2291 snapshot.quota_remaining_usd = Some(amount_to_string(remaining));
2292 snapshot.quota_used_usd = Some(amount_to_string(used));
2293 snapshot.quota_limit_usd = Some(amount_to_string(limit));
2294 if !matched_plan_names.is_empty() {
2295 matched_plan_names.sort();
2296 matched_plan_names.dedup();
2297 snapshot.plan_name = Some(matched_plan_names.join(", "));
2298 }
2299 snapshot.exhausted = Some(remaining.is_zero() && balance.is_none_or(UsdAmount::is_zero));
2300 } else {
2301 snapshot.exhausted = balance.map(UsdAmount::is_zero);
2302 }
2303
2304 snapshot.refresh_status(fetched_at_ms);
2305 snapshot
2306}
2307
2308fn new_api_token_usage_snapshot_from_json(
2309 provider: &UsageProviderConfig,
2310 upstream: &UpstreamRef,
2311 value: &serde_json::Value,
2312 fetched_at_ms: u64,
2313 stale_after_ms: Option<u64>,
2314) -> ProviderBalanceSnapshot {
2315 if json_value_at_path(value, "success").and_then(bool_from_json) == Some(false)
2316 || json_value_at_path(value, "code").and_then(bool_from_json) == Some(false)
2317 {
2318 let message = json_value_at_path(value, "message")
2319 .and_then(|value| value.as_str())
2320 .unwrap_or("new api token usage response reported failure");
2321 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2322 .with_error(message.to_string());
2323 }
2324
2325 let mut effective = provider.extract.clone();
2326 if effective.remaining_balance_paths.is_empty() {
2327 effective.remaining_balance_paths = vec![
2328 "data.total_available".to_string(),
2329 "data.remain_quota".to_string(),
2330 "total_available".to_string(),
2331 "remain_quota".to_string(),
2332 ];
2333 }
2334 if effective.monthly_spent_paths.is_empty() {
2335 effective.monthly_spent_paths = vec![
2336 "data.total_used".to_string(),
2337 "data.used_quota".to_string(),
2338 "total_used".to_string(),
2339 "used_quota".to_string(),
2340 ];
2341 }
2342 if effective.monthly_budget_paths.is_empty() {
2343 effective.monthly_budget_paths = vec![
2344 "data.total_granted".to_string(),
2345 "total_granted".to_string(),
2346 ];
2347 }
2348 effective.remaining_divisor = effective.remaining_divisor.or(Some(500_000));
2349 effective.monthly_spent_divisor = effective.monthly_spent_divisor.or(Some(500_000));
2350 effective.monthly_budget_divisor = effective.monthly_budget_divisor.or(Some(500_000));
2351
2352 let unlimited_quota =
2353 first_bool_from_paths(value, &[], &["data.unlimited_quota", "unlimited_quota"])
2354 == Some(true);
2355 let remaining_balance = first_amount_from_paths(
2356 value,
2357 &effective.remaining_balance_paths,
2358 &[],
2359 effective.remaining_divisor,
2360 );
2361 let monthly_spent = first_amount_from_paths(
2362 value,
2363 &effective.monthly_spent_paths,
2364 &[],
2365 effective.monthly_spent_divisor,
2366 );
2367 let monthly_budget = first_amount_from_paths(
2368 value,
2369 &effective.monthly_budget_paths,
2370 &[],
2371 effective.monthly_budget_divisor,
2372 )
2373 .or_else(|| match (remaining_balance, monthly_spent) {
2374 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
2375 _ => None,
2376 });
2377 let exhausted = if unlimited_quota {
2378 Some(false)
2379 } else {
2380 first_bool_from_paths(
2381 value,
2382 &effective.exhausted_paths,
2383 &["data.exhausted", "exhausted"],
2384 )
2385 .or_else(|| remaining_balance.map(UsdAmount::is_zero))
2386 };
2387
2388 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2389 snapshot.plan_name = first_string_from_paths(value, &["data.name", "name"]);
2390 snapshot.unlimited_quota = Some(unlimited_quota);
2391 if !unlimited_quota {
2392 snapshot.quota_period = Some("token".to_string());
2393 snapshot.quota_remaining_usd = remaining_balance.map(amount_to_string);
2394 snapshot.quota_limit_usd = monthly_budget.map(amount_to_string);
2395 snapshot.quota_used_usd = monthly_spent.map(amount_to_string);
2396 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
2397 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2398 } else {
2399 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2400 }
2401 snapshot.exhausted = exhausted;
2402 snapshot.refresh_status(fetched_at_ms);
2403 snapshot
2404}
2405
2406fn new_api_snapshot_from_json(
2407 provider: &UsageProviderConfig,
2408 upstream: &UpstreamRef,
2409 value: &serde_json::Value,
2410 fetched_at_ms: u64,
2411 stale_after_ms: Option<u64>,
2412) -> ProviderBalanceSnapshot {
2413 if json_value_at_path(value, "success").and_then(bool_from_json) == Some(false) {
2414 let message = json_value_at_path(value, "message")
2415 .and_then(|value| value.as_str())
2416 .unwrap_or("new api balance response reported failure");
2417 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2418 .with_error(message.to_string());
2419 }
2420
2421 let mut effective = provider.extract.clone();
2422 if effective.remaining_balance_paths.is_empty() {
2423 effective.remaining_balance_paths = vec!["data.quota".to_string(), "quota".to_string()];
2424 }
2425 if effective.monthly_spent_paths.is_empty() {
2426 effective.monthly_spent_paths =
2427 vec!["data.used_quota".to_string(), "used_quota".to_string()];
2428 }
2429 effective.remaining_divisor = effective.remaining_divisor.or(Some(500_000));
2430 effective.monthly_spent_divisor = effective.monthly_spent_divisor.or(Some(500_000));
2431 effective.monthly_budget_divisor = effective.monthly_budget_divisor.or(Some(500_000));
2432
2433 let remaining_balance = first_amount_from_paths(
2434 value,
2435 &effective.remaining_balance_paths,
2436 &[],
2437 effective.remaining_divisor,
2438 );
2439 let monthly_spent = first_amount_from_paths(
2440 value,
2441 &effective.monthly_spent_paths,
2442 &[],
2443 effective.monthly_spent_divisor,
2444 );
2445 let monthly_budget = first_amount_from_paths(
2446 value,
2447 &effective.monthly_budget_paths,
2448 &["data.total_quota", "total_quota"],
2449 effective.monthly_budget_divisor,
2450 )
2451 .or_else(|| match (remaining_balance, monthly_spent) {
2452 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
2453 _ => None,
2454 });
2455 let unlimited_quota =
2456 first_bool_from_paths(value, &[], &["data.unlimited_quota", "unlimited_quota"])
2457 == Some(true);
2458 let exhausted = if unlimited_quota {
2459 Some(false)
2460 } else {
2461 first_bool_from_paths(
2462 value,
2463 &effective.exhausted_paths,
2464 &["data.exhausted", "exhausted"],
2465 )
2466 .or_else(|| remaining_balance.map(UsdAmount::is_zero))
2467 };
2468
2469 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2470 snapshot.unlimited_quota = Some(unlimited_quota);
2471 if !unlimited_quota {
2472 snapshot.quota_period = Some("quota".to_string());
2473 snapshot.quota_remaining_usd = remaining_balance.map(amount_to_string);
2474 snapshot.quota_limit_usd = monthly_budget.map(amount_to_string);
2475 snapshot.quota_used_usd = monthly_spent.map(amount_to_string);
2476 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
2477 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2478 } else {
2479 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2480 }
2481 snapshot.exhausted = exhausted;
2482 snapshot.refresh_status(fetched_at_ms);
2483 snapshot
2484}
2485
2486fn openai_cost_result_usd_amount(result: &serde_json::Value) -> Option<UsdAmount> {
2487 let amount = json_value_at_path(result, "amount.value").and_then(amount_from_json)?;
2488 let currency = json_value_at_path(result, "amount.currency").and_then(|value| value.as_str());
2489 match currency {
2490 Some(currency) if currency.eq_ignore_ascii_case("usd") => Some(amount),
2491 None => Some(amount),
2492 _ => None,
2493 }
2494}
2495
2496fn openai_organization_costs_total(value: &serde_json::Value) -> Option<UsdAmount> {
2497 let buckets = json_value_at_path(value, "data")?.as_array()?;
2498 let mut total = UsdAmount::ZERO;
2499
2500 for bucket in buckets {
2501 let Some(results) =
2502 json_value_at_path(bucket, "results").and_then(|value| value.as_array())
2503 else {
2504 continue;
2505 };
2506 for result in results {
2507 if let Some(amount) = openai_cost_result_usd_amount(result) {
2508 total = total.saturating_add(amount);
2509 }
2510 }
2511 }
2512
2513 Some(total)
2514}
2515
2516fn openai_organization_costs_snapshot_from_json(
2517 provider: &UsageProviderConfig,
2518 upstream: &UpstreamRef,
2519 value: &serde_json::Value,
2520 fetched_at_ms: u64,
2521 stale_after_ms: Option<u64>,
2522) -> ProviderBalanceSnapshot {
2523 let spent = first_amount_from_paths(
2524 value,
2525 &provider.extract.monthly_spent_paths,
2526 &[],
2527 provider.extract.monthly_spent_divisor,
2528 )
2529 .or_else(|| openai_organization_costs_total(value));
2530
2531 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2532 snapshot.monthly_spent_usd = spent.map(amount_to_string);
2533 snapshot.exhausted = None;
2534 snapshot.exhaustion_affects_routing = false;
2535 snapshot.refresh_status(fetched_at_ms);
2536 snapshot
2537}
2538
2539fn balance_exhaustion_policy_action(
2540 endpoint_key: ProviderEndpointKey,
2541 observed_at_ms: u64,
2542) -> Option<PolicyAction> {
2543 const BALANCE_EXHAUSTION_COOLDOWN_SECS: u64 = 24 * 60 * 60;
2544 let mut signal = ProviderSignal::high_confidence_route_facing(
2545 ProviderSignalKind::Balance,
2546 ProviderSignalSource::BalanceSnapshot,
2547 ProviderSignalTarget::ProviderEndpoint {
2548 provider_endpoint_key: endpoint_key,
2549 },
2550 observed_at_ms,
2551 );
2552 signal.reset_after_secs = Some(BALANCE_EXHAUSTION_COOLDOWN_SECS);
2553 signal.reason = Some("balance_exhausted".to_string());
2554 PolicyAction::cooldown_from_signal(signal, observed_at_ms, 0, observed_at_ms)
2555}
2556
2557async fn sync_balance_policy_action_for_endpoint(
2558 state: &Arc<ProxyState>,
2559 service_name: &str,
2560 endpoint_key: ProviderEndpointKey,
2561 exhausted: bool,
2562) {
2563 if exhausted {
2564 if let Some(action) =
2565 balance_exhaustion_policy_action(endpoint_key, crate::logging::now_ms())
2566 {
2567 state.upsert_owned_policy_action(service_name, action).await;
2568 }
2569 } else {
2570 state
2571 .clear_owned_policy_action(
2572 service_name,
2573 &endpoint_key,
2574 PolicyActionKind::Cooldown,
2575 ProviderSignalKind::Balance,
2576 ProviderSignalSource::BalanceSnapshot,
2577 )
2578 .await;
2579 }
2580}
2581
2582async fn update_usage_exhausted(
2583 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
2584 state: &Arc<ProxyState>,
2585 cfg: &ProxyConfig,
2586 service_name: &str,
2587 upstreams: &[UpstreamRef],
2588 exhausted: bool,
2589) {
2590 if let Ok(mut map) = lb_states.lock() {
2591 for uref in upstreams {
2592 let service = match service_manager(cfg, service_name).station(&uref.station_name) {
2593 Some(s) => s,
2594 None => continue,
2595 };
2596
2597 let entry = map
2598 .entry(uref.station_name.clone())
2599 .or_insert_with(LbState::default);
2600 entry.ensure_layout(service.name.as_str(), &service.upstreams);
2601 if uref.index < entry.usage_exhausted.len() {
2602 entry.usage_exhausted[uref.index] = exhausted;
2603 }
2604 }
2605 }
2606
2607 for uref in upstreams {
2608 if let Some(endpoint_key) = uref.provider_endpoint.clone() {
2609 sync_balance_policy_action_for_endpoint(
2610 state,
2611 service_name,
2612 endpoint_key.clone(),
2613 exhausted,
2614 )
2615 .await;
2616 state
2617 .set_provider_endpoint_usage_exhausted(service_name, endpoint_key, exhausted)
2618 .await;
2619 }
2620 }
2621}
2622
2623fn provider_hosts_for_diagnostics(
2624 cfg: &ProxyConfig,
2625 service_name: &str,
2626 provider: &UsageProviderConfig,
2627) -> Vec<String> {
2628 let mut hosts: Vec<String> = Vec::new();
2629 for service in service_manager(cfg, service_name).stations().values() {
2630 for upstream in &service.upstreams {
2631 if domain_matches(&upstream.base_url, &provider.domains)
2632 && let Ok(url) = reqwest::Url::parse(&upstream.base_url)
2633 && let Some(host) = url.host_str()
2634 {
2635 hosts.push(host.to_string());
2636 }
2637 }
2638 }
2639 hosts.sort();
2640 hosts.dedup();
2641 hosts
2642}
2643
2644fn warn_if_provider_spans_hosts(
2645 cfg: &ProxyConfig,
2646 service_name: &str,
2647 provider: &UsageProviderConfig,
2648) {
2649 let hosts = provider_hosts_for_diagnostics(cfg, service_name, provider);
2650 if hosts.len() > 1 {
2651 warn!(
2652 "usage provider '{}' is associated with multiple hosts: {:?}; \
2653将按统一额度处理这些 upstream,如需区分配额请拆分为多个 provider 配置",
2654 provider.id, hosts
2655 );
2656 }
2657}
2658
2659fn snapshot_from_provider_json(
2660 provider: &UsageProviderConfig,
2661 upstream: &UpstreamRef,
2662 value: &serde_json::Value,
2663 upstream_base_url: &str,
2664 fetched_at_ms: u64,
2665 stale_after_ms: Option<u64>,
2666) -> ProviderBalanceSnapshot {
2667 match provider.kind {
2668 ProviderKind::BudgetHttpJson => {
2669 budget_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
2670 }
2671 ProviderKind::YescodeProfile => {
2672 yescode_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
2673 }
2674 ProviderKind::OpenAiBalanceHttpJson => balance_http_snapshot_from_json(
2675 provider,
2676 upstream,
2677 value,
2678 fetched_at_ms,
2679 stale_after_ms,
2680 ),
2681 ProviderKind::Sub2ApiUsage => sub2api_usage_snapshot_from_json(
2682 provider,
2683 upstream,
2684 value,
2685 fetched_at_ms,
2686 stale_after_ms,
2687 ),
2688 ProviderKind::Sub2ApiAuthMe => sub2api_auth_me_snapshot_from_json(
2689 provider,
2690 upstream,
2691 value,
2692 fetched_at_ms,
2693 stale_after_ms,
2694 ),
2695 ProviderKind::NewApiTokenUsage => new_api_token_usage_snapshot_from_json(
2696 provider,
2697 upstream,
2698 value,
2699 fetched_at_ms,
2700 stale_after_ms,
2701 ),
2702 ProviderKind::NewApiUserSelf => {
2703 new_api_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
2704 }
2705 ProviderKind::RightCodeAccountSummary => rightcode_account_summary_snapshot_from_json(
2706 provider,
2707 upstream,
2708 value,
2709 upstream_base_url,
2710 fetched_at_ms,
2711 stale_after_ms,
2712 ),
2713 ProviderKind::OpenAiOrganizationCosts => openai_organization_costs_snapshot_from_json(
2714 provider,
2715 upstream,
2716 value,
2717 fetched_at_ms,
2718 stale_after_ms,
2719 ),
2720 }
2721}
2722
2723async fn refresh_provider_target(
2724 params: RefreshProviderTargetParams<'_>,
2725) -> UsageProviderRefreshOutcome {
2726 let RefreshProviderTargetParams {
2727 client,
2728 provider,
2729 target,
2730 cfg,
2731 lb_states,
2732 state,
2733 service_name,
2734 interval_secs,
2735 } = params;
2736
2737 let upstreams = vec![target.upstream.clone()];
2738 let fetched_at_ms = unix_now_ms();
2739 let stale_after_ms = stale_after_ms(fetched_at_ms, interval_secs);
2740
2741 let Some(token) = resolve_token(provider, &upstreams, cfg, service_name) else {
2742 let snapshot = if provider.kind == ProviderKind::OpenAiOrganizationCosts {
2743 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
2744 } else {
2745 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
2746 .with_error("no usable token; checked provider token_env and upstream auth")
2747 };
2748 state
2749 .record_provider_balance_snapshot(service_name, snapshot)
2750 .await;
2751 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
2752 if provider.kind == ProviderKind::OpenAiOrganizationCosts {
2753 warn!(
2754 "usage provider '{}' is missing OPENAI_ADMIN_KEY; OpenAI official costs stay unknown",
2755 provider.id
2756 );
2757 } else {
2758 warn!(
2759 "usage provider '{}' has no usable token (checked token_env and associated upstream auth_token); \
2760跳过本次用量查询,请检查 usage_providers.json 和 ~/.codex-helper/config.json",
2761 provider.id
2762 );
2763 }
2764 return UsageProviderRefreshOutcome::MissingToken;
2765 };
2766
2767 match poll_provider_http_json(client, provider, &target.base_url, &token).await {
2768 Ok(value) => {
2769 let snapshot = snapshot_from_provider_json(
2770 provider,
2771 &upstreams[0],
2772 &value,
2773 &target.base_url,
2774 fetched_at_ms,
2775 stale_after_ms,
2776 );
2777 let exhausted_for_lb = snapshot.routing_exhausted();
2778 update_usage_exhausted(
2779 lb_states,
2780 state,
2781 cfg,
2782 service_name,
2783 &upstreams,
2784 exhausted_for_lb,
2785 )
2786 .await;
2787 state
2788 .record_provider_balance_snapshot(service_name, snapshot)
2789 .await;
2790 info!(
2791 "usage provider '{}' refreshed {}[{}], exhausted = {}, routing_trusted = {}",
2792 provider.id,
2793 target.upstream.station_name,
2794 target.upstream.index,
2795 exhausted_for_lb,
2796 provider.trust_exhaustion_for_routing
2797 );
2798 UsageProviderRefreshOutcome::Refreshed
2799 }
2800 Err(err) => {
2801 state
2802 .record_provider_balance_snapshot(
2803 service_name,
2804 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
2805 .with_error(err.to_string()),
2806 )
2807 .await;
2808 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
2809 warn!(
2810 "usage provider '{}' poll failed for {}[{}]: {}",
2811 provider.id, target.upstream.station_name, target.upstream.index, err
2812 );
2813 UsageProviderRefreshOutcome::Failed
2814 }
2815 }
2816}
2817
2818struct ConfiguredRefreshJob<'a> {
2819 provider: &'a UsageProviderConfig,
2820 target: UsageProviderTarget,
2821 interval_secs: u64,
2822}
2823
2824struct AutoRefreshJob {
2825 target: UsageProviderTarget,
2826}
2827
2828async fn run_configured_refresh_job<'a>(
2829 client: &'a Client,
2830 job: ConfiguredRefreshJob<'a>,
2831 cfg: &'a ProxyConfig,
2832 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
2833 state: &'a Arc<ProxyState>,
2834 service_name: &'a str,
2835) -> (String, UsageProviderRefreshOutcome) {
2836 let provider_id = job.provider.id.clone();
2837 let outcome = refresh_provider_target(RefreshProviderTargetParams {
2838 client,
2839 provider: job.provider,
2840 target: &job.target,
2841 cfg,
2842 lb_states,
2843 state,
2844 service_name,
2845 interval_secs: job.interval_secs,
2846 })
2847 .await;
2848 (provider_id, outcome)
2849}
2850
2851async fn run_auto_refresh_job(
2852 client: &Client,
2853 job: AutoRefreshJob,
2854 cfg: &ProxyConfig,
2855 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
2856 state: &Arc<ProxyState>,
2857 service_name: &str,
2858) -> UsageProviderRefreshOutcome {
2859 auto_probe_provider_target(client, &job.target, cfg, lb_states, state, service_name).await
2860}
2861
2862async fn run_configured_refresh_jobs<'a>(
2863 client: &'a Client,
2864 jobs: Vec<ConfiguredRefreshJob<'a>>,
2865 cfg: &'a ProxyConfig,
2866 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
2867 state: &'a Arc<ProxyState>,
2868 service_name: &'a str,
2869) -> Vec<(String, UsageProviderRefreshOutcome)> {
2870 let mut pending = jobs.into_iter();
2871 let mut running = FuturesUnordered::new();
2872 let mut results = Vec::new();
2873 let concurrency = BALANCE_REFRESH_CONCURRENCY.max(1);
2874
2875 for _ in 0..concurrency {
2876 let Some(job) = pending.next() else {
2877 break;
2878 };
2879 running.push(run_configured_refresh_job(
2880 client,
2881 job,
2882 cfg,
2883 lb_states,
2884 state,
2885 service_name,
2886 ));
2887 }
2888
2889 while let Some(result) = running.next().await {
2890 results.push(result);
2891 if let Some(job) = pending.next() {
2892 running.push(run_configured_refresh_job(
2893 client,
2894 job,
2895 cfg,
2896 lb_states,
2897 state,
2898 service_name,
2899 ));
2900 }
2901 }
2902
2903 results
2904}
2905
2906async fn run_auto_refresh_jobs(
2907 client: &Client,
2908 jobs: Vec<AutoRefreshJob>,
2909 cfg: &ProxyConfig,
2910 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
2911 state: &Arc<ProxyState>,
2912 service_name: &str,
2913) -> Vec<UsageProviderRefreshOutcome> {
2914 let mut pending = jobs.into_iter();
2915 let mut running = FuturesUnordered::new();
2916 let mut results = Vec::new();
2917 let concurrency = BALANCE_REFRESH_CONCURRENCY.max(1);
2918
2919 for _ in 0..concurrency {
2920 let Some(job) = pending.next() else {
2921 break;
2922 };
2923 running.push(run_auto_refresh_job(
2924 client,
2925 job,
2926 cfg,
2927 lb_states,
2928 state,
2929 service_name,
2930 ));
2931 }
2932
2933 while let Some(result) = running.next().await {
2934 results.push(result);
2935 if let Some(job) = pending.next() {
2936 running.push(run_auto_refresh_job(
2937 client,
2938 job,
2939 cfg,
2940 lb_states,
2941 state,
2942 service_name,
2943 ));
2944 }
2945 }
2946
2947 results
2948}
2949
2950fn auto_snapshot_is_usable(snapshot: &ProviderBalanceSnapshot) -> bool {
2951 snapshot.error.is_none()
2952 && matches!(
2953 snapshot.status,
2954 BalanceSnapshotStatus::Ok | BalanceSnapshotStatus::Exhausted
2955 )
2956}
2957
2958async fn auto_probe_provider_target(
2959 client: &Client,
2960 target: &UsageProviderTarget,
2961 cfg: &ProxyConfig,
2962 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
2963 state: &Arc<ProxyState>,
2964 service_name: &str,
2965) -> UsageProviderRefreshOutcome {
2966 let upstreams = vec![target.upstream.clone()];
2967 let fetched_at_ms = unix_now_ms();
2968 let interval_secs = DEFAULT_POLL_INTERVAL_SECS;
2969 let stale_after_ms = stale_after_ms(fetched_at_ms, interval_secs);
2970
2971 if is_official_openai_base_url(&target.base_url) {
2972 let provider = auto_openai_official_provider(target);
2973 let Some(token) = resolve_token(&provider, &upstreams, cfg, service_name) else {
2974 state
2975 .record_provider_balance_snapshot(
2976 service_name,
2977 base_snapshot(&provider, &target.upstream, fetched_at_ms, stale_after_ms),
2978 )
2979 .await;
2980 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
2981 warn!(
2982 "OpenAI organization costs require OPENAI_ADMIN_KEY; balance stays unknown for {}[{}]",
2983 target.upstream.station_name, target.upstream.index
2984 );
2985 return UsageProviderRefreshOutcome::MissingToken;
2986 };
2987
2988 return match poll_provider_http_json(client, &provider, &target.base_url, &token).await {
2989 Ok(value) => {
2990 let snapshot = snapshot_from_provider_json(
2991 &provider,
2992 &upstreams[0],
2993 &value,
2994 &target.base_url,
2995 fetched_at_ms,
2996 stale_after_ms,
2997 );
2998 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false)
2999 .await;
3000 state
3001 .record_provider_balance_snapshot(service_name, snapshot)
3002 .await;
3003 UsageProviderRefreshOutcome::Refreshed
3004 }
3005 Err(err) => {
3006 state
3007 .record_provider_balance_snapshot(
3008 service_name,
3009 base_snapshot(&provider, &upstreams[0], fetched_at_ms, stale_after_ms)
3010 .with_error(err.to_string()),
3011 )
3012 .await;
3013 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false)
3014 .await;
3015 warn!(
3016 "OpenAI organization costs poll failed for {}[{}]: {}",
3017 target.upstream.station_name, target.upstream.index, err
3018 );
3019 UsageProviderRefreshOutcome::Failed
3020 }
3021 };
3022 }
3023
3024 let first_provider = auto_usage_provider(target, first_auto_probe_kind(target));
3025
3026 let Some(token) = resolve_token(&first_provider, &upstreams, cfg, service_name) else {
3027 state
3028 .record_provider_balance_snapshot(
3029 service_name,
3030 base_snapshot(
3031 &first_provider,
3032 &target.upstream,
3033 fetched_at_ms,
3034 stale_after_ms,
3035 )
3036 .with_error("no usable token; checked upstream auth"),
3037 )
3038 .await;
3039 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3040 return UsageProviderRefreshOutcome::MissingToken;
3041 };
3042
3043 let mut last_error: Option<String> = None;
3044 for kind in AUTO_PROBE_KINDS {
3045 if kind == ProviderKind::RightCodeAccountSummary && !is_rightcode_base_url(&target.base_url)
3046 {
3047 continue;
3048 }
3049 let provider = auto_usage_provider(target, kind);
3050 match poll_provider_http_json(client, &provider, &target.base_url, &token).await {
3051 Ok(value) => {
3052 let snapshot = snapshot_from_provider_json(
3053 &provider,
3054 &upstreams[0],
3055 &value,
3056 &target.base_url,
3057 fetched_at_ms,
3058 stale_after_ms,
3059 );
3060 if auto_snapshot_is_usable(&snapshot) {
3061 let exhausted_for_lb = snapshot.routing_exhausted();
3062 update_usage_exhausted(
3063 lb_states,
3064 state,
3065 cfg,
3066 service_name,
3067 &upstreams,
3068 exhausted_for_lb,
3069 )
3070 .await;
3071 state
3072 .record_provider_balance_snapshot(service_name, snapshot)
3073 .await;
3074 info!(
3075 "auto usage provider '{}' refreshed {}[{}] via {:?}, exhausted = {}",
3076 provider.id,
3077 target.upstream.station_name,
3078 target.upstream.index,
3079 kind,
3080 exhausted_for_lb
3081 );
3082 return UsageProviderRefreshOutcome::Refreshed;
3083 }
3084 last_error = snapshot.error.or_else(|| {
3085 Some(format!(
3086 "auto probe {:?} returned no usable balance fields",
3087 kind
3088 ))
3089 });
3090 }
3091 Err(err) => {
3092 last_error = Some(err.to_string());
3093 }
3094 }
3095 }
3096
3097 if let Some(error) = last_error {
3098 warn!(
3099 "auto usage provider '{}' found no usable balance endpoint for {}[{}]: {}",
3100 first_provider.id, target.upstream.station_name, target.upstream.index, error
3101 );
3102 state
3103 .record_provider_balance_snapshot(
3104 service_name,
3105 base_snapshot(
3106 &first_provider,
3107 &target.upstream,
3108 fetched_at_ms,
3109 stale_after_ms,
3110 )
3111 .with_error(error),
3112 )
3113 .await;
3114 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3115 }
3116 UsageProviderRefreshOutcome::Failed
3117}
3118
3119pub async fn refresh_balances_for_service(
3120 client: &Client,
3121 cfg: Arc<ProxyConfig>,
3122 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
3123 state: Arc<ProxyState>,
3124 service_name: &str,
3125 station_name_filter: Option<&str>,
3126 provider_id_filter: Option<&str>,
3127) -> UsageProviderRefreshSummary {
3128 if cfg!(test) {
3130 return UsageProviderRefreshSummary::default();
3131 }
3132
3133 let station_name_filter = station_name_filter
3134 .map(str::trim)
3135 .filter(|value| !value.is_empty());
3136 let provider_id_filter = provider_id_filter
3137 .map(str::trim)
3138 .filter(|value| !value.is_empty());
3139 let providers_file = load_providers();
3140 let mut summary = UsageProviderRefreshSummary {
3141 providers_configured: providers_file.providers.len(),
3142 ..UsageProviderRefreshSummary::default()
3143 };
3144
3145 let poll_map = LAST_USAGE_POLL.get_or_init(|| Mutex::new(HashMap::new()));
3146 let mut configured_jobs = Vec::new();
3147 let mut configured_job_keys = HashSet::new();
3148 for provider in &providers_file.providers {
3149 if provider_id_filter.is_some_and(|filter| filter != provider.id.as_str()) {
3150 continue;
3151 }
3152
3153 let targets = matching_provider_targets(&cfg, service_name, provider, station_name_filter);
3154 if targets.is_empty() {
3155 continue;
3156 }
3157
3158 summary.providers_matched += 1;
3159 summary.upstreams_matched += targets.len();
3160 warn_if_provider_spans_hosts(&cfg, service_name, provider);
3161
3162 let interval_secs = snapshot_refresh_interval_secs(provider);
3163 for target in targets {
3164 summary.attempted += 1;
3165 configured_job_keys.insert(target_key(&target));
3166 configured_jobs.push(ConfiguredRefreshJob {
3167 provider,
3168 target,
3169 interval_secs,
3170 });
3171 }
3172 }
3173
3174 let mut refreshed_provider_ids = HashSet::new();
3175 if !configured_jobs.is_empty() {
3176 for (provider_id, outcome) in run_configured_refresh_jobs(
3177 client,
3178 configured_jobs,
3179 &cfg,
3180 &lb_states,
3181 &state,
3182 service_name,
3183 )
3184 .await
3185 {
3186 match outcome {
3187 UsageProviderRefreshOutcome::Refreshed => {
3188 summary.refreshed += 1;
3189 refreshed_provider_ids.insert(provider_id);
3190 }
3191 UsageProviderRefreshOutcome::Failed => summary.failed += 1,
3192 UsageProviderRefreshOutcome::MissingToken => summary.missing_token += 1,
3193 }
3194 }
3195 }
3196
3197 let mut auto_jobs = Vec::new();
3198 for target in usage_provider_targets(&cfg, service_name, station_name_filter) {
3199 if configured_job_keys.contains(&target_key(&target)) {
3200 continue;
3201 }
3202 if !auto_target_matches_provider_id_filter(&target, provider_id_filter) {
3203 continue;
3204 }
3205
3206 summary.attempted += 1;
3207 summary.auto_attempted += 1;
3208 auto_jobs.push(AutoRefreshJob { target });
3209 }
3210
3211 if !auto_jobs.is_empty() {
3212 for outcome in
3213 run_auto_refresh_jobs(client, auto_jobs, &cfg, &lb_states, &state, service_name).await
3214 {
3215 match outcome {
3216 UsageProviderRefreshOutcome::Refreshed => {
3217 summary.refreshed += 1;
3218 summary.auto_refreshed += 1;
3219 }
3220 UsageProviderRefreshOutcome::Failed => {
3221 summary.failed += 1;
3222 summary.auto_failed += 1;
3223 }
3224 UsageProviderRefreshOutcome::MissingToken => {
3225 summary.missing_token += 1;
3226 }
3227 }
3228 }
3229 }
3230
3231 if !refreshed_provider_ids.is_empty()
3232 && let Ok(mut map) = poll_map.lock()
3233 {
3234 let now = Instant::now();
3235 for provider_id in refreshed_provider_ids {
3236 map.insert(provider_id, now);
3237 }
3238 }
3239
3240 summary
3241}
3242
3243pub fn enqueue_poll_for_codex_upstream(
3246 client: Client,
3247 cfg: Arc<ProxyConfig>,
3248 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
3249 state: Arc<ProxyState>,
3250 service_name: &str,
3251 station_name: &str,
3252 upstream_index: usize,
3253) {
3254 let key = RequestBalanceQueueKey::LegacyUpstream {
3255 service_name: service_name.to_string(),
3256 station_name: station_name.to_string(),
3257 upstream_index,
3258 };
3259 let Some(initial_sleep_for) = enqueue_request_balance_refresh(key.clone()) else {
3260 return;
3261 };
3262
3263 let service_name = service_name.to_string();
3264 let station_name = station_name.to_string();
3265 tokio::spawn(async move {
3266 let mut sleep_for = initial_sleep_for;
3267 loop {
3268 tokio::time::sleep(sleep_for).await;
3269 match take_request_balance_refresh_if_due(&key) {
3270 RequestBalanceQueueDue::Due => {}
3271 RequestBalanceQueueDue::NotDue(delay) => {
3272 sleep_for = delay;
3273 continue;
3274 }
3275 RequestBalanceQueueDue::Missing => return,
3276 }
3277
3278 let current_target = usage_provider_target_for_legacy_upstream(
3279 &cfg,
3280 service_name.as_str(),
3281 &station_name,
3282 upstream_index,
3283 );
3284 match poll_for_codex_target(
3285 client.clone(),
3286 cfg.clone(),
3287 lb_states.clone(),
3288 state.clone(),
3289 service_name.as_str(),
3290 current_target,
3291 )
3292 .await
3293 {
3294 RequestBalancePollOutcome::Attempted | RequestBalancePollOutcome::Skipped => {
3295 return;
3296 }
3297 RequestBalancePollOutcome::Deferred(delay) => {
3298 schedule_request_balance_refresh_at(key.clone(), Instant::now() + delay);
3299 sleep_for = delay;
3300 }
3301 }
3302 }
3303 });
3304}
3305
3306pub fn enqueue_poll_for_codex_provider_endpoint(
3310 client: Client,
3311 cfg: Arc<ProxyConfig>,
3312 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
3313 state: Arc<ProxyState>,
3314 service_name: &str,
3315 provider_endpoint: ProviderEndpointKey,
3316) {
3317 let key = RequestBalanceQueueKey::ProviderEndpoint(provider_endpoint.clone());
3318 let Some(initial_sleep_for) = enqueue_request_balance_refresh(key.clone()) else {
3319 return;
3320 };
3321
3322 let service_name = service_name.to_string();
3323 tokio::spawn(async move {
3324 let mut sleep_for = initial_sleep_for;
3325 loop {
3326 tokio::time::sleep(sleep_for).await;
3327 match take_request_balance_refresh_if_due(&key) {
3328 RequestBalanceQueueDue::Due => {}
3329 RequestBalanceQueueDue::NotDue(delay) => {
3330 sleep_for = delay;
3331 continue;
3332 }
3333 RequestBalanceQueueDue::Missing => return,
3334 }
3335
3336 let current_target = usage_provider_target_for_provider_endpoint(
3337 &cfg,
3338 service_name.as_str(),
3339 &provider_endpoint,
3340 );
3341 match poll_for_codex_target(
3342 client.clone(),
3343 cfg.clone(),
3344 lb_states.clone(),
3345 state.clone(),
3346 service_name.as_str(),
3347 current_target,
3348 )
3349 .await
3350 {
3351 RequestBalancePollOutcome::Attempted | RequestBalancePollOutcome::Skipped => {
3352 return;
3353 }
3354 RequestBalancePollOutcome::Deferred(delay) => {
3355 schedule_request_balance_refresh_at(key.clone(), Instant::now() + delay);
3356 sleep_for = delay;
3357 }
3358 }
3359 }
3360 });
3361}
3362
3363async fn poll_for_codex_target(
3364 client: Client,
3365 cfg: Arc<ProxyConfig>,
3366 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
3367 state: Arc<ProxyState>,
3368 service_name: &str,
3369 current_target: Option<UsageProviderTarget>,
3370) -> RequestBalancePollOutcome {
3371 if cfg!(test) {
3374 return RequestBalancePollOutcome::Skipped;
3375 }
3376
3377 let providers_file = load_providers();
3378 let Some(current_target) = current_target else {
3379 return RequestBalancePollOutcome::Skipped;
3380 };
3381
3382 let now = Instant::now();
3383 let poll_map = LAST_USAGE_POLL.get_or_init(|| Mutex::new(HashMap::new()));
3384 let mut matched_configured_provider = false;
3385 let mut configured_jobs = Vec::new();
3386 let mut next_cooldown = None::<Duration>;
3387
3388 for provider in &providers_file.providers {
3389 if !domain_matches(¤t_target.base_url, &provider.domains) {
3390 continue;
3391 }
3392 matched_configured_provider = true;
3393
3394 let Some(interval_secs) = effective_poll_interval_secs(provider) else {
3395 continue;
3396 };
3397
3398 {
3399 let mut map = match poll_map.lock() {
3400 Ok(m) => m,
3401 Err(_) => continue,
3402 };
3403 if let Some(last) = map.get(&provider.id)
3404 && let Some(cooldown) = remaining_poll_cooldown(*last, interval_secs, now)
3405 {
3406 next_cooldown =
3407 Some(next_cooldown.map_or(cooldown, |existing| existing.min(cooldown)));
3408 continue;
3409 }
3410 map.insert(provider.id.clone(), now);
3411 }
3412
3413 warn_if_provider_spans_hosts(&cfg, service_name, provider);
3414 configured_jobs.push(ConfiguredRefreshJob {
3415 provider,
3416 target: current_target.clone(),
3417 interval_secs,
3418 });
3419 }
3420
3421 if !configured_jobs.is_empty() {
3422 let _ = run_configured_refresh_jobs(
3423 &client,
3424 configured_jobs,
3425 &cfg,
3426 &lb_states,
3427 &state,
3428 service_name,
3429 )
3430 .await;
3431 return RequestBalancePollOutcome::Attempted;
3432 }
3433
3434 if matched_configured_provider {
3435 return next_cooldown
3436 .map(RequestBalancePollOutcome::Deferred)
3437 .unwrap_or(RequestBalancePollOutcome::Skipped);
3438 }
3439
3440 let auto_provider = if is_official_openai_base_url(¤t_target.base_url) {
3441 auto_openai_official_provider(¤t_target)
3442 } else {
3443 auto_usage_provider(¤t_target, first_auto_probe_kind(¤t_target))
3444 };
3445 let Some(interval_secs) = effective_poll_interval_secs(&auto_provider) else {
3446 return RequestBalancePollOutcome::Skipped;
3447 };
3448
3449 {
3450 let mut map = match poll_map.lock() {
3451 Ok(m) => m,
3452 Err(_) => return RequestBalancePollOutcome::Skipped,
3453 };
3454 if let Some(last) = map.get(&auto_provider.id)
3455 && let Some(cooldown) = remaining_poll_cooldown(*last, interval_secs, now)
3456 {
3457 return RequestBalancePollOutcome::Deferred(cooldown);
3458 }
3459 map.insert(auto_provider.id.clone(), now);
3460 }
3461
3462 let _ = auto_probe_provider_target(
3463 &client,
3464 ¤t_target,
3465 &cfg,
3466 &lb_states,
3467 &state,
3468 service_name,
3469 )
3470 .await;
3471 RequestBalancePollOutcome::Attempted
3472}
3473
3474#[cfg(test)]
3475mod tests {
3476 use super::*;
3477
3478 use crate::balance::BalanceSnapshotStatus;
3479 use crate::config::{ServiceConfig, UpstreamAuth, UpstreamConfig};
3480
3481 fn provider(id: &str, kind: ProviderKind) -> UsageProviderConfig {
3482 UsageProviderConfig {
3483 id: id.to_string(),
3484 kind,
3485 domains: vec!["example.com".to_string()],
3486 endpoint: "https://example.com/usage".to_string(),
3487 token_env: None,
3488 require_token_env: false,
3489 poll_interval_secs: Some(60),
3490 refresh_on_request: true,
3491 trust_exhaustion_for_routing: true,
3492 headers: BTreeMap::new(),
3493 variables: BTreeMap::new(),
3494 extract: UsageProviderExtractConfig::default(),
3495 }
3496 }
3497
3498 fn upstream() -> UpstreamRef {
3499 UpstreamRef {
3500 station_name: "right".to_string(),
3501 index: 1,
3502 provider_endpoint: None,
3503 }
3504 }
3505
3506 fn endpoint_upstream() -> UpstreamRef {
3507 UpstreamRef {
3508 station_name: "right".to_string(),
3509 index: 1,
3510 provider_endpoint: Some(ProviderEndpointKey::new("codex", "right", "default")),
3511 }
3512 }
3513
3514 fn upstream_config(base_url: &str) -> UpstreamConfig {
3515 UpstreamConfig {
3516 base_url: base_url.to_string(),
3517 auth: UpstreamAuth::default(),
3518 tags: HashMap::new(),
3519 supported_models: HashMap::new(),
3520 model_mapping: HashMap::new(),
3521 }
3522 }
3523
3524 fn endpoint_upstream_config(
3525 base_url: &str,
3526 provider_id: &str,
3527 endpoint_id: &str,
3528 ) -> UpstreamConfig {
3529 let mut upstream = upstream_config(base_url);
3530 upstream
3531 .tags
3532 .insert("provider_id".to_string(), provider_id.to_string());
3533 upstream
3534 .tags
3535 .insert("endpoint_id".to_string(), endpoint_id.to_string());
3536 upstream
3537 }
3538
3539 fn service_config(name: &str, upstreams: Vec<UpstreamConfig>) -> ServiceConfig {
3540 ServiceConfig {
3541 name: name.to_string(),
3542 alias: None,
3543 enabled: true,
3544 level: 1,
3545 upstreams,
3546 }
3547 }
3548
3549 fn proxy_config(stations: Vec<ServiceConfig>) -> ProxyConfig {
3550 let mut cfg = ProxyConfig::default();
3551 cfg.codex.configs = stations
3552 .into_iter()
3553 .map(|station| (station.name.clone(), station))
3554 .collect();
3555 cfg
3556 }
3557
3558 #[test]
3559 fn budget_snapshot_reports_monthly_budget_and_exhaustion() {
3560 let snapshot = budget_snapshot_from_json(
3561 &provider("packycode", ProviderKind::BudgetHttpJson),
3562 &upstream(),
3563 &serde_json::json!({
3564 "monthly_budget_usd": "10.50",
3565 "monthly_spent_usd": 10.5
3566 }),
3567 100,
3568 Some(1_000),
3569 );
3570
3571 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
3572 assert_eq!(snapshot.exhausted, Some(true));
3573 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10.5"));
3574 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("10.5"));
3575 }
3576
3577 #[test]
3578 fn budget_snapshot_keeps_missing_amounts_unknown() {
3579 let snapshot = budget_snapshot_from_json(
3580 &provider("packycode", ProviderKind::BudgetHttpJson),
3581 &upstream(),
3582 &serde_json::json!({}),
3583 100,
3584 Some(1_000),
3585 );
3586
3587 assert_eq!(snapshot.status, BalanceSnapshotStatus::Unknown);
3588 assert_eq!(snapshot.exhausted, None);
3589 }
3590
3591 #[test]
3592 fn yescode_snapshot_sums_subscription_and_paygo_balances() {
3593 let snapshot = yescode_snapshot_from_json(
3594 &provider("yescode", ProviderKind::YescodeProfile),
3595 &upstream(),
3596 &serde_json::json!({
3597 "subscription_balance": "1.25",
3598 "pay_as_you_go_balance": 2.5
3599 }),
3600 100,
3601 Some(1_000),
3602 );
3603
3604 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
3605 assert_eq!(snapshot.exhausted, Some(false));
3606 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.75"));
3607 assert_eq!(snapshot.subscription_balance_usd.as_deref(), Some("1.25"));
3608 assert_eq!(snapshot.paygo_balance_usd.as_deref(), Some("2.5"));
3609 }
3610
3611 #[test]
3612 fn openai_balance_endpoint_defaults_to_base_user_balance_without_v1() {
3613 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
3614 provider.endpoint.clear();
3615
3616 let endpoint =
3617 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
3618
3619 assert_eq!(endpoint, "https://relay.example.com/user/balance");
3620 }
3621
3622 #[test]
3623 fn sub2api_usage_endpoint_defaults_to_upstream_usage_under_v1() {
3624 let mut provider = provider("sub2api", ProviderKind::Sub2ApiUsage);
3625 provider.endpoint.clear();
3626
3627 let endpoint =
3628 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
3629
3630 assert_eq!(endpoint, "https://relay.example.com/v1/usage");
3631 }
3632
3633 #[test]
3634 fn sub2api_auth_me_endpoint_defaults_to_dashboard_path_without_v1() {
3635 let mut provider = provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe);
3636 provider.endpoint.clear();
3637
3638 let endpoint =
3639 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
3640
3641 assert_eq!(endpoint, "https://relay.example.com/api/v1/auth/me");
3642 }
3643
3644 #[test]
3645 fn provider_templates_support_variables_for_custom_headers_or_queries() {
3646 let mut provider = provider("newapi", ProviderKind::NewApiUserSelf);
3647 provider.endpoint = "{{base_url}}/api/user/self?user={{userId}}".to_string();
3648 provider
3649 .variables
3650 .insert("userId".to_string(), "42".to_string());
3651
3652 let endpoint = resolve_endpoint(&provider, "https://newapi.example.com/v1", "token")
3653 .expect("endpoint");
3654
3655 assert_eq!(endpoint, "https://newapi.example.com/api/user/self?user=42");
3656 }
3657
3658 #[test]
3659 fn new_api_token_usage_endpoint_defaults_to_model_key_usage_path() {
3660 let mut provider = provider("newapi-token", ProviderKind::NewApiTokenUsage);
3661 provider.endpoint.clear();
3662
3663 let endpoint = resolve_endpoint(&provider, "https://newapi.example.com/v1", "token")
3664 .expect("endpoint");
3665
3666 assert_eq!(endpoint, "https://newapi.example.com/api/usage/token/");
3667 }
3668
3669 #[test]
3670 fn openai_organization_costs_endpoint_defaults_to_official_v1_costs_window() {
3671 let mut provider = provider("openai", ProviderKind::OpenAiOrganizationCosts);
3672 provider.endpoint.clear();
3673
3674 let endpoint =
3675 resolve_endpoint(&provider, "https://api.openai.com/v1", "token").expect("endpoint");
3676
3677 assert!(endpoint.starts_with("https://api.openai.com/v1/organization/costs?start_time="));
3678 assert!(endpoint.ends_with("&limit=30"));
3679 let start_time = endpoint
3680 .split("start_time=")
3681 .nth(1)
3682 .and_then(|value| value.split('&').next())
3683 .and_then(|value| value.parse::<u64>().ok())
3684 .expect("numeric start_time");
3685 assert!(start_time > 0);
3686 }
3687
3688 #[test]
3689 fn require_token_env_prevents_upstream_model_key_fallback() {
3690 let mut cfg = proxy_config(vec![service_config(
3691 "right",
3692 vec![upstream_config("https://api.openai.com/v1")],
3693 )]);
3694 cfg.codex
3695 .configs
3696 .get_mut("right")
3697 .expect("station")
3698 .upstreams[0]
3699 .auth
3700 .auth_token = Some("model-key".to_string());
3701
3702 let mut provider = provider("openai", ProviderKind::OpenAiOrganizationCosts);
3703 provider.token_env = Some("__CODEX_HELPER_TEST_MISSING_TOKEN_ENV__".to_string());
3704 provider.require_token_env = true;
3705 let upstreams = [UpstreamRef {
3706 station_name: "right".to_string(),
3707 index: 0,
3708 provider_endpoint: None,
3709 }];
3710
3711 assert_eq!(resolve_token(&provider, &upstreams, &cfg, "codex"), None);
3712
3713 provider.require_token_env = false;
3714 assert_eq!(
3715 resolve_token(&provider, &upstreams, &cfg, "codex").as_deref(),
3716 Some("model-key")
3717 );
3718 }
3719
3720 #[test]
3721 fn effective_poll_interval_respects_disable_flag_zero_and_minimum() {
3722 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
3723
3724 provider.poll_interval_secs = Some(0);
3725 assert_eq!(effective_poll_interval_secs(&provider), None);
3726
3727 provider.poll_interval_secs = Some(10);
3728 assert_eq!(
3729 effective_poll_interval_secs(&provider),
3730 Some(MIN_POLL_INTERVAL_SECS)
3731 );
3732
3733 provider.poll_interval_secs = None;
3734 assert_eq!(
3735 effective_poll_interval_secs(&provider),
3736 Some(DEFAULT_POLL_INTERVAL_SECS)
3737 );
3738
3739 provider.refresh_on_request = false;
3740 assert_eq!(effective_poll_interval_secs(&provider), None);
3741 }
3742
3743 #[test]
3744 fn auto_provider_uses_stable_target_id_across_probe_kinds() {
3745 let target = UsageProviderTarget {
3746 upstream: UpstreamRef {
3747 station_name: "input/sub".to_string(),
3748 index: 2,
3749 provider_endpoint: None,
3750 },
3751 base_url: "https://ai.input.im/v1".to_string(),
3752 provider_id: None,
3753 };
3754
3755 let sub2api = auto_usage_provider(&target, ProviderKind::Sub2ApiUsage);
3756 let newapi_token = auto_usage_provider(&target, ProviderKind::NewApiTokenUsage);
3757 let newapi = auto_usage_provider(&target, ProviderKind::NewApiUserSelf);
3758
3759 assert_eq!(sub2api.id, "auto:balance:input-sub:2");
3760 assert_eq!(sub2api.id, newapi_token.id);
3761 assert_eq!(sub2api.id, newapi.id);
3762 assert_eq!(sub2api.domains, vec!["ai.input.im".to_string()]);
3763 assert_eq!(
3764 resolve_endpoint(&sub2api, &target.base_url, "token").unwrap(),
3765 "https://ai.input.im/v1/usage"
3766 );
3767 assert_eq!(
3768 resolve_endpoint(&newapi_token, &target.base_url, "token").unwrap(),
3769 "https://ai.input.im/api/usage/token/"
3770 );
3771 }
3772
3773 #[test]
3774 fn auto_probe_prefers_rightcode_adapter_for_rightcode_hosts() {
3775 let target = UsageProviderTarget {
3776 upstream: UpstreamRef {
3777 station_name: "right".to_string(),
3778 index: 0,
3779 provider_endpoint: None,
3780 },
3781 base_url: "https://www.right.codes/codex/v1".to_string(),
3782 provider_id: Some("right".to_string()),
3783 };
3784
3785 assert_eq!(
3786 first_auto_probe_kind(&target),
3787 ProviderKind::RightCodeAccountSummary
3788 );
3789 assert_eq!(
3790 resolve_endpoint(
3791 &auto_usage_provider(&target, ProviderKind::RightCodeAccountSummary),
3792 &target.base_url,
3793 "token"
3794 )
3795 .unwrap(),
3796 "https://www.right.codes/account/summary"
3797 );
3798 assert_eq!(
3799 auto_usage_provider(&target, ProviderKind::RightCodeAccountSummary)
3800 .token_env
3801 .as_deref(),
3802 None
3803 );
3804 }
3805
3806 #[test]
3807 fn auto_provider_id_prefers_runtime_provider_tag() {
3808 let target = UsageProviderTarget {
3809 upstream: UpstreamRef {
3810 station_name: "routing".to_string(),
3811 index: 0,
3812 provider_endpoint: Some(ProviderEndpointKey::new("codex", "input", "default")),
3813 },
3814 base_url: "https://ai.input.im/v1".to_string(),
3815 provider_id: Some("input".to_string()),
3816 };
3817
3818 let provider = auto_usage_provider(&target, ProviderKind::Sub2ApiUsage);
3819
3820 assert_eq!(provider.id, "input");
3821 }
3822
3823 #[test]
3824 fn auto_target_provider_id_filter_matches_runtime_provider_tag() {
3825 let target = UsageProviderTarget {
3826 upstream: UpstreamRef {
3827 station_name: "routing".to_string(),
3828 index: 6,
3829 provider_endpoint: Some(ProviderEndpointKey::new("codex", "input6", "default")),
3830 },
3831 base_url: "https://input.9z1.me/v1".to_string(),
3832 provider_id: Some("input6".to_string()),
3833 };
3834
3835 assert!(auto_target_matches_provider_id_filter(&target, None));
3836 assert!(auto_target_matches_provider_id_filter(
3837 &target,
3838 Some("input6")
3839 ));
3840 assert!(!auto_target_matches_provider_id_filter(
3841 &target,
3842 Some("input5")
3843 ));
3844 }
3845
3846 #[test]
3847 fn auto_target_provider_id_filter_matches_generated_auto_id() {
3848 let target = UsageProviderTarget {
3849 upstream: UpstreamRef {
3850 station_name: "routing".to_string(),
3851 index: 6,
3852 provider_endpoint: None,
3853 },
3854 base_url: "https://input.9z1.me/v1".to_string(),
3855 provider_id: None,
3856 };
3857
3858 assert!(auto_target_matches_provider_id_filter(
3859 &target,
3860 Some("auto:balance:routing:6")
3861 ));
3862 assert!(!auto_target_matches_provider_id_filter(
3863 &target,
3864 Some("input6")
3865 ));
3866 }
3867
3868 #[test]
3869 fn provider_endpoint_target_lookup_uses_endpoint_identity() {
3870 let cfg = proxy_config(vec![service_config(
3871 "routing",
3872 vec![
3873 endpoint_upstream_config("https://input.example/v1", "input", "default"),
3874 endpoint_upstream_config("https://right.example/v1", "right", "default"),
3875 ],
3876 )]);
3877
3878 let target = usage_provider_target_for_provider_endpoint(
3879 &cfg,
3880 "codex",
3881 &ProviderEndpointKey::new("codex", "right", "default"),
3882 )
3883 .expect("provider endpoint target");
3884
3885 assert_eq!(target.upstream.station_name, "routing");
3886 assert_eq!(target.upstream.index, 1);
3887 assert_eq!(
3888 target
3889 .upstream
3890 .provider_endpoint
3891 .as_ref()
3892 .map(ProviderEndpointKey::stable_key)
3893 .as_deref(),
3894 Some("codex/right/default")
3895 );
3896 assert_eq!(target.base_url, "https://right.example/v1");
3897 assert_eq!(target.provider_id.as_deref(), Some("right"));
3898 }
3899
3900 #[test]
3901 fn request_balance_queue_deduplicates_until_due() {
3902 let key = RequestBalanceQueueKey::ProviderEndpoint(ProviderEndpointKey::new(
3903 "codex", "input", "default",
3904 ));
3905 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
3906 {
3907 let mut queue = queue.lock().expect("queue");
3908 queue.remove(&key);
3909 }
3910
3911 assert_eq!(
3912 enqueue_request_balance_refresh(key.clone()),
3913 Some(REQUEST_BALANCE_REFRESH_DELAY)
3914 );
3915 assert_eq!(enqueue_request_balance_refresh(key.clone()), None);
3916 assert!(matches!(
3917 take_request_balance_refresh_if_due(&key),
3918 RequestBalanceQueueDue::NotDue(_)
3919 ));
3920
3921 {
3922 let mut queue = queue.lock().expect("queue");
3923 queue.insert(key.clone(), Instant::now() - Duration::from_secs(1));
3924 }
3925
3926 assert_eq!(
3927 take_request_balance_refresh_if_due(&key),
3928 RequestBalanceQueueDue::Due
3929 );
3930 assert_eq!(
3931 take_request_balance_refresh_if_due(&key),
3932 RequestBalanceQueueDue::Missing
3933 );
3934 assert_eq!(
3935 enqueue_request_balance_refresh(key.clone()),
3936 Some(REQUEST_BALANCE_REFRESH_DELAY)
3937 );
3938
3939 queue.lock().expect("queue").remove(&key);
3940 }
3941
3942 #[test]
3943 fn request_balance_queue_does_not_extend_due_refresh() {
3944 let key = RequestBalanceQueueKey::ProviderEndpoint(ProviderEndpointKey::new(
3945 "codex", "input", "default",
3946 ));
3947 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
3948 {
3949 let mut queue = queue.lock().expect("queue");
3950 queue.insert(key.clone(), Instant::now() - Duration::from_secs(1));
3951 }
3952
3953 assert_eq!(
3954 enqueue_request_balance_refresh(key.clone()),
3955 Some(Duration::ZERO)
3956 );
3957 assert_eq!(
3958 take_request_balance_refresh_if_due(&key),
3959 RequestBalanceQueueDue::Due
3960 );
3961 assert_eq!(
3962 take_request_balance_refresh_if_due(&key),
3963 RequestBalanceQueueDue::Missing
3964 );
3965
3966 queue.lock().expect("queue").remove(&key);
3967 }
3968
3969 #[test]
3970 fn remaining_poll_cooldown_returns_only_unexpired_window() {
3971 let now = Instant::now();
3972 assert_eq!(
3973 remaining_poll_cooldown(now - Duration::from_secs(30), 60, now),
3974 Some(Duration::from_secs(30))
3975 );
3976 assert_eq!(
3977 remaining_poll_cooldown(now - Duration::from_secs(60), 60, now),
3978 None
3979 );
3980 assert_eq!(
3981 remaining_poll_cooldown(now - Duration::from_secs(61), 60, now),
3982 None
3983 );
3984 }
3985
3986 #[test]
3987 fn request_balance_queue_distinguishes_legacy_upstreams() {
3988 let key = RequestBalanceQueueKey::LegacyUpstream {
3989 service_name: "codex".to_string(),
3990 station_name: "routing".to_string(),
3991 upstream_index: 2,
3992 };
3993 let same_key = RequestBalanceQueueKey::LegacyUpstream {
3994 service_name: "codex".to_string(),
3995 station_name: "routing".to_string(),
3996 upstream_index: 2,
3997 };
3998 let other_key = RequestBalanceQueueKey::LegacyUpstream {
3999 service_name: "codex".to_string(),
4000 station_name: "routing".to_string(),
4001 upstream_index: 3,
4002 };
4003
4004 assert_eq!(key, same_key);
4005 assert_ne!(key, other_key);
4006 }
4007
4008 #[test]
4009 fn configured_target_keys_prevent_auto_probe_for_explicit_balance_domains() {
4010 let cfg = proxy_config(vec![
4011 service_config("explicit", vec![upstream_config("https://example.com/v1")]),
4012 service_config("auto", vec![upstream_config("https://ai.input.im/v1")]),
4013 ]);
4014 let configured = configured_target_keys(
4015 &cfg,
4016 "codex",
4017 &[provider("relay", ProviderKind::OpenAiBalanceHttpJson)],
4018 None,
4019 );
4020 let auto_targets = usage_provider_targets(&cfg, "codex", None)
4021 .into_iter()
4022 .filter(|target| !configured.contains(&target_key(target)))
4023 .map(|target| target.upstream.station_name)
4024 .collect::<Vec<_>>();
4025
4026 assert_eq!(auto_targets, vec!["auto".to_string()]);
4027 }
4028
4029 #[test]
4030 fn auto_probe_accepts_only_usable_balance_snapshots() {
4031 let usable = sub2api_usage_snapshot_from_json(
4032 &provider("auto", ProviderKind::Sub2ApiUsage),
4033 &upstream(),
4034 &serde_json::json!({
4035 "isValid": true,
4036 "remaining": 1
4037 }),
4038 100,
4039 Some(1_000),
4040 );
4041 let unusable = balance_http_snapshot_from_json(
4042 &provider("auto", ProviderKind::OpenAiBalanceHttpJson),
4043 &upstream(),
4044 &serde_json::json!({ "ok": true }),
4045 100,
4046 Some(1_000),
4047 );
4048
4049 assert!(auto_snapshot_is_usable(&usable));
4050 assert!(!auto_snapshot_is_usable(&unusable));
4051 }
4052
4053 #[test]
4054 fn openai_balance_snapshot_reads_common_sub2api_balance_shape() {
4055 let snapshot = balance_http_snapshot_from_json(
4056 &provider("sub2api", ProviderKind::OpenAiBalanceHttpJson),
4057 &upstream(),
4058 &serde_json::json!({
4059 "balance": "1.25"
4060 }),
4061 100,
4062 Some(1_000),
4063 );
4064
4065 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4066 assert_eq!(snapshot.exhausted, Some(false));
4067 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("1.25"));
4068 }
4069
4070 #[test]
4071 fn json_path_supports_array_indices_for_official_balance_shapes() {
4072 let value = serde_json::json!({
4073 "balance_infos": [
4074 { "currency": "CNY", "total_balance": "3.25" }
4075 ]
4076 });
4077
4078 assert_eq!(
4079 json_value_at_path(&value, "balance_infos.0.total_balance")
4080 .and_then(|value| value.as_str()),
4081 Some("3.25")
4082 );
4083 }
4084
4085 #[test]
4086 fn openai_balance_snapshot_reads_cc_switch_official_balance_shapes() {
4087 let snapshot = balance_http_snapshot_from_json(
4088 &provider("deepseek", ProviderKind::OpenAiBalanceHttpJson),
4089 &upstream(),
4090 &serde_json::json!({
4091 "balance_infos": [
4092 { "currency": "CNY", "total_balance": "3.25" }
4093 ],
4094 "is_available": true
4095 }),
4096 100,
4097 Some(1_000),
4098 );
4099
4100 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4101 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.25"));
4102
4103 let snapshot = balance_http_snapshot_from_json(
4104 &provider("siliconflow", ProviderKind::OpenAiBalanceHttpJson),
4105 &upstream(),
4106 &serde_json::json!({
4107 "code": 20000,
4108 "data": {
4109 "totalBalance": "8.5",
4110 "chargeBalance": "2.5"
4111 }
4112 }),
4113 100,
4114 Some(1_000),
4115 );
4116
4117 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4118 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("8.5"));
4119 assert_eq!(snapshot.paygo_balance_usd.as_deref(), Some("2.5"));
4120 }
4121
4122 #[test]
4123 fn openai_balance_snapshot_can_derive_remaining_from_total_and_used() {
4124 let mut provider = provider("openrouter", ProviderKind::OpenAiBalanceHttpJson);
4125 provider.extract.monthly_budget_paths = vec!["data.total_credits".to_string()];
4126 provider.extract.monthly_spent_paths = vec!["data.total_usage".to_string()];
4127 provider.extract.derive_remaining_from_budget_and_spent = true;
4128
4129 let snapshot = balance_http_snapshot_from_json(
4130 &provider,
4131 &upstream(),
4132 &serde_json::json!({
4133 "data": {
4134 "total_credits": "10",
4135 "total_usage": "4"
4136 }
4137 }),
4138 100,
4139 Some(1_000),
4140 );
4141
4142 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4143 assert_eq!(snapshot.exhausted, Some(false));
4144 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("6"));
4145 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10"));
4146 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("4"));
4147 }
4148
4149 #[test]
4150 fn openai_balance_snapshot_supports_divisor_for_minor_units() {
4151 let mut provider = provider("novita", ProviderKind::OpenAiBalanceHttpJson);
4152 provider.extract.remaining_balance_paths = vec!["availableBalance".to_string()];
4153 provider.extract.remaining_divisor = Some(10_000);
4154
4155 let snapshot = balance_http_snapshot_from_json(
4156 &provider,
4157 &upstream(),
4158 &serde_json::json!({
4159 "availableBalance": 12345
4160 }),
4161 100,
4162 Some(1_000),
4163 );
4164
4165 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4166 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("1.2345"));
4167 }
4168
4169 #[test]
4170 fn sub2api_usage_snapshot_reads_all_api_hub_usage_shape() {
4171 let snapshot = sub2api_usage_snapshot_from_json(
4172 &provider("sub2api", ProviderKind::Sub2ApiUsage),
4173 &upstream(),
4174 &serde_json::json!({
4175 "isValid": true,
4176 "mode": "unrestricted",
4177 "planName": "CodeX Air",
4178 "remaining": 165.0877165,
4179 "usage": {
4180 "today": {
4181 "cost": 0,
4182 "requests": 0,
4183 "total_tokens": 0
4184 },
4185 "total": {
4186 "cost": 354.194748,
4187 "requests": 2691,
4188 "total_tokens": 384084697
4189 }
4190 }
4191 }),
4192 100,
4193 Some(1_000),
4194 );
4195
4196 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4197 assert_eq!(snapshot.exhausted, Some(false));
4198 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Air"));
4199 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("165.0877165"));
4200 assert_eq!(snapshot.total_used_usd.as_deref(), Some("354.194748"));
4201 assert_eq!(snapshot.today_used_usd.as_deref(), Some("0"));
4202 assert_eq!(snapshot.total_requests, Some(2691));
4203 assert_eq!(snapshot.today_requests, Some(0));
4204 assert_eq!(snapshot.total_tokens, Some(384084697));
4205 assert_eq!(snapshot.today_tokens, Some(0));
4206 }
4207
4208 #[test]
4209 fn sub2api_usage_snapshot_reads_rates_model_stats_windows_and_alerts() {
4210 let snapshot = sub2api_usage_snapshot_from_json(
4211 &provider("sub2api", ProviderKind::Sub2ApiUsage),
4212 &upstream(),
4213 &serde_json::json!({
4214 "isValid": true,
4215 "mode": "unrestricted",
4216 "plan_name": "CodeX Pro",
4217 "remaining": 9,
4218 "subscription": {
4219 "daily_usage_usd": 95,
4220 "daily_limit_usd": 100,
4221 "weekly_usage_usd": "120.5",
4222 "weekly_limit_usd": 0,
4223 "monthly_usage_usd": 300.25,
4224 "monthly_limit_usd": 1000,
4225 "expires_at": "2026-05-09T12:00:00.000Z"
4226 },
4227 "usage": {
4228 "today": {
4229 "request_count": "7",
4230 "input_tokens": 100,
4231 "output_tokens": 25,
4232 "total_cost_usd": "1.5"
4233 },
4234 "total": {
4235 "requests": 42,
4236 "tokens": 1234,
4237 "cost": 9.25
4238 },
4239 "average_duration_ms": "842.7",
4240 "rpm": "0.7",
4241 "tpm": 85.3
4242 },
4243 "model_stats": [
4244 {
4245 "model": "gpt-4o-mini",
4246 "request_count": "7",
4247 "prompt_tokens": 100,
4248 "completion_tokens": 25,
4249 "input_cost": "0.12",
4250 "output_cost": "0.34"
4251 }
4252 ]
4253 }),
4254 100,
4255 Some(1_000),
4256 );
4257
4258 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4259 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Pro"));
4260 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("9"));
4261 assert_eq!(snapshot.today_requests, Some(7));
4262 assert_eq!(snapshot.today_tokens, Some(100));
4263 assert_eq!(snapshot.today_used_usd.as_deref(), Some("1.5"));
4264 assert_eq!(snapshot.total_requests, Some(42));
4265 assert_eq!(snapshot.total_tokens, Some(1234));
4266 assert_eq!(snapshot.total_used_usd.as_deref(), Some("9.25"));
4267 let rate = snapshot.usage_rate.expect("rate");
4268 assert_eq!(rate.average_duration_ms.as_deref(), Some("842.7"));
4269 assert_eq!(rate.rpm.as_deref(), Some("0.7"));
4270 assert_eq!(rate.tpm.as_deref(), Some("85.3"));
4271 assert_eq!(snapshot.usage_windows.len(), 3);
4272 assert_eq!(snapshot.usage_windows[0].period, "daily");
4273 assert_eq!(
4274 snapshot.usage_windows[0].remaining_usd.as_deref(),
4275 Some("5")
4276 );
4277 assert_eq!(snapshot.usage_windows[1].unlimited, Some(true));
4278 assert_eq!(snapshot.usage_model_stats.len(), 1);
4279 assert_eq!(snapshot.usage_model_stats[0].model, "gpt-4o-mini");
4280 assert_eq!(snapshot.usage_model_stats[0].request_count, Some(7));
4281 assert_eq!(snapshot.usage_model_stats[0].total_tokens, Some(125));
4282 assert_eq!(
4283 snapshot.usage_model_stats[0].total_cost_usd.as_deref(),
4284 Some("0.46")
4285 );
4286 assert_eq!(
4287 snapshot
4288 .usage_alerts
4289 .iter()
4290 .map(|alert| alert.kind)
4291 .collect::<Vec<_>>(),
4292 vec![
4293 ProviderUsageAlertKind::DailyUsage95,
4294 ProviderUsageAlertKind::LowBalance,
4295 ProviderUsageAlertKind::SubscriptionExpired,
4296 ]
4297 );
4298 }
4299
4300 #[test]
4301 fn sub2api_subscription_zero_remaining_is_display_only_period_capacity_exhaustion() {
4302 let snapshot = sub2api_usage_snapshot_from_json(
4303 &provider("sub2api", ProviderKind::Sub2ApiUsage),
4304 &upstream(),
4305 &serde_json::json!({
4306 "isValid": true,
4307 "mode": "unrestricted",
4308 "planName": "CodeX Lite 年度",
4309 "remaining": 0,
4310 "subscription": {
4311 "daily_usage_usd": 100.468025,
4312 "daily_limit_usd": 100,
4313 "weekly_usage_usd": 401.441684,
4314 "weekly_limit_usd": 0,
4315 "monthly_usage_usd": 401.441684,
4316 "monthly_limit_usd": 0
4317 },
4318 "usage": {
4319 "today": { "cost": 0, "requests": 0, "total_tokens": 0 },
4320 "total": { "cost": 702.492098, "requests": 42, "total_tokens": 1234 }
4321 }
4322 }),
4323 100,
4324 Some(1_000),
4325 );
4326
4327 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
4328 assert_eq!(snapshot.exhausted, Some(true));
4329 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Lite 年度"));
4330 assert_eq!(snapshot.total_balance_usd, None);
4331 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
4332 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
4333 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("100"));
4334 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("100.468025"));
4335 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("100"));
4336 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("100.468025"));
4337 assert_eq!(snapshot.total_used_usd.as_deref(), Some("702.492098"));
4338 assert_eq!(snapshot.today_used_usd.as_deref(), Some("0"));
4339 assert!(
4340 !snapshot.routing_exhausted(),
4341 "sub2api /v1/usage skips billing checks; subscription windows are reset lazily on real requests"
4342 );
4343 }
4344
4345 #[test]
4346 fn sub2api_quota_limited_zero_remaining_still_marks_exhausted() {
4347 let snapshot = sub2api_usage_snapshot_from_json(
4348 &provider("sub2api", ProviderKind::Sub2ApiUsage),
4349 &upstream(),
4350 &serde_json::json!({
4351 "isValid": true,
4352 "mode": "quota_limited",
4353 "quota": {
4354 "limit": 10,
4355 "used": 10,
4356 "remaining": 0,
4357 "unit": "USD"
4358 }
4359 }),
4360 100,
4361 Some(1_000),
4362 );
4363
4364 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
4365 assert_eq!(snapshot.exhausted, Some(true));
4366 assert_eq!(snapshot.total_balance_usd, None);
4367 assert_eq!(snapshot.quota_period.as_deref(), Some("quota"));
4368 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
4369 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("10"));
4370 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("10"));
4371 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10"));
4372 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("10"));
4373 assert!(snapshot.routing_exhausted());
4374 }
4375
4376 #[test]
4377 fn sub2api_usage_snapshot_marks_invalid_key_as_error() {
4378 let snapshot = sub2api_usage_snapshot_from_json(
4379 &provider("sub2api", ProviderKind::Sub2ApiUsage),
4380 &upstream(),
4381 &serde_json::json!({
4382 "isValid": false
4383 }),
4384 100,
4385 Some(1_000),
4386 );
4387
4388 assert_eq!(snapshot.status, BalanceSnapshotStatus::Error);
4389 assert_eq!(
4390 snapshot.error.as_deref(),
4391 Some("sub2api usage response reported invalid API key")
4392 );
4393 }
4394
4395 #[test]
4396 fn sub2api_auth_me_snapshot_reads_dashboard_balance_envelope() {
4397 let snapshot = sub2api_auth_me_snapshot_from_json(
4398 &provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe),
4399 &upstream(),
4400 &serde_json::json!({
4401 "code": 0,
4402 "message": "ok",
4403 "data": {
4404 "id": 42,
4405 "username": "demo",
4406 "balance": "12.5"
4407 }
4408 }),
4409 100,
4410 Some(1_000),
4411 );
4412
4413 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4414 assert_eq!(snapshot.exhausted, Some(false));
4415 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("12.5"));
4416 }
4417
4418 #[test]
4419 fn rightcode_endpoint_defaults_to_account_summary() {
4420 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
4421 provider.endpoint.clear();
4422
4423 let endpoint = resolve_endpoint(&provider, "https://www.right.codes/codex/v1", "token")
4424 .expect("endpoint");
4425
4426 assert_eq!(endpoint, "https://www.right.codes/account/summary");
4427 }
4428
4429 #[test]
4430 fn rightcode_account_summary_reads_matching_subscription_and_balance() {
4431 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
4432 provider.trust_exhaustion_for_routing = false;
4433
4434 let snapshot = rightcode_account_summary_snapshot_from_json(
4435 &provider,
4436 &upstream(),
4437 &serde_json::json!({
4438 "balance": 3.25,
4439 "subscriptions": [
4440 {
4441 "name": "Daily",
4442 "total_quota": 20,
4443 "remaining_quota": 7.5,
4444 "reset_today": true,
4445 "available_prefixes": ["/codex"]
4446 },
4447 {
4448 "name": "Other",
4449 "total_quota": 99,
4450 "remaining_quota": 99,
4451 "reset_today": true,
4452 "available_prefixes": ["/claude"]
4453 },
4454 {
4455 "name": "Badge",
4456 "total_quota": 5,
4457 "remaining_quota": 5,
4458 "reset_today": true,
4459 "available_prefixes": ["/codex"]
4460 }
4461 ]
4462 }),
4463 "https://www.right.codes/codex/v1",
4464 100,
4465 Some(1_000),
4466 );
4467
4468 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4469 assert_eq!(snapshot.exhausted, Some(false));
4470 assert!(!snapshot.routing_exhausted());
4471 assert_eq!(snapshot.plan_name.as_deref(), Some("Daily"));
4472 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.25"));
4473 assert_eq!(snapshot.paygo_balance_usd, None);
4474 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
4475 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("7.5"));
4476 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("20"));
4477 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("12.5"));
4478 }
4479
4480 #[test]
4481 fn rightcode_account_summary_accounts_for_not_reset_today() {
4482 let provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
4483
4484 let snapshot = rightcode_account_summary_snapshot_from_json(
4485 &provider,
4486 &upstream(),
4487 &serde_json::json!({
4488 "subscriptions": [
4489 {
4490 "name": "Daily",
4491 "total_quota": 20,
4492 "remaining_quota": 7.5,
4493 "reset_today": false,
4494 "available_prefixes": ["/codex"]
4495 }
4496 ]
4497 }),
4498 "https://www.right.codes/codex/v1",
4499 100,
4500 Some(1_000),
4501 );
4502
4503 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4504 assert_eq!(snapshot.exhausted, Some(false));
4505 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("27.5"));
4506 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("20"));
4507 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("0"));
4508 }
4509
4510 #[test]
4511 fn rightcode_zero_daily_quota_without_balance_is_display_only_exhaustion_by_default() {
4512 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
4513 provider.trust_exhaustion_for_routing = false;
4514
4515 let snapshot = rightcode_account_summary_snapshot_from_json(
4516 &provider,
4517 &upstream(),
4518 &serde_json::json!({
4519 "balance": 0,
4520 "subscriptions": [
4521 {
4522 "name": "Daily",
4523 "total_quota": 20,
4524 "remaining_quota": 0,
4525 "reset_today": true,
4526 "available_prefixes": ["/codex"]
4527 }
4528 ]
4529 }),
4530 "https://www.right.codes/codex/v1",
4531 100,
4532 Some(1_000),
4533 );
4534
4535 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
4536 assert_eq!(snapshot.exhausted, Some(true));
4537 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
4538 assert!(!snapshot.routing_exhausted());
4539 assert!(snapshot.routing_ignored_exhaustion());
4540 }
4541
4542 #[test]
4543 fn sub2api_auth_me_snapshot_marks_business_error() {
4544 let snapshot = sub2api_auth_me_snapshot_from_json(
4545 &provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe),
4546 &upstream(),
4547 &serde_json::json!({
4548 "code": 401,
4549 "message": "login required"
4550 }),
4551 100,
4552 Some(1_000),
4553 );
4554
4555 assert_eq!(snapshot.status, BalanceSnapshotStatus::Error);
4556 assert_eq!(snapshot.error.as_deref(), Some("login required"));
4557 }
4558
4559 #[test]
4560 fn provider_can_disable_routing_trust_for_exhausted_balance() {
4561 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
4562 provider.trust_exhaustion_for_routing = false;
4563
4564 let snapshot = balance_http_snapshot_from_json(
4565 &provider,
4566 &upstream(),
4567 &serde_json::json!({
4568 "balance": "0"
4569 }),
4570 100,
4571 Some(1_000),
4572 );
4573
4574 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
4575 assert_eq!(snapshot.exhausted, Some(true));
4576 assert!(!snapshot.exhaustion_affects_routing);
4577 assert!(!snapshot.routing_exhausted());
4578 }
4579
4580 #[test]
4581 fn provider_exhaustion_trust_defaults_to_enabled_when_omitted() {
4582 let provider: UsageProviderConfig = serde_json::from_value(serde_json::json!({
4583 "id": "sub2api",
4584 "kind": "openai_balance_http_json",
4585 "domains": ["example.com"]
4586 }))
4587 .expect("provider config");
4588
4589 assert!(provider.trust_exhaustion_for_routing);
4590 }
4591
4592 #[tokio::test]
4593 async fn provider_missing_token_clears_stale_lb_exhaustion_marker() {
4594 let cfg = proxy_config(vec![service_config(
4595 "right",
4596 vec![
4597 upstream_config("https://primary.example/v1"),
4598 upstream_config("https://backup.example/v1"),
4599 ],
4600 )]);
4601 let lb_states = Arc::new(Mutex::new(HashMap::new()));
4602 let target = UsageProviderTarget {
4603 upstream: upstream(),
4604 base_url: "https://backup.example/v1".to_string(),
4605 provider_id: Some("right".to_string()),
4606 };
4607 let upstreams = vec![target.upstream.clone()];
4608 let state = ProxyState::new();
4609 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, true).await;
4610 {
4611 let guard = lb_states.lock().expect("lb states");
4612 assert!(
4613 guard
4614 .get("right")
4615 .and_then(|entry| entry.usage_exhausted.get(1))
4616 .copied()
4617 .unwrap_or(false)
4618 );
4619 }
4620
4621 let outcome = refresh_provider_target(RefreshProviderTargetParams {
4622 client: &Client::new(),
4623 provider: &provider("sub2api", ProviderKind::OpenAiBalanceHttpJson),
4624 target: &target,
4625 cfg: &cfg,
4626 lb_states: &lb_states,
4627 state: &state,
4628 service_name: "codex",
4629 interval_secs: 60,
4630 })
4631 .await;
4632
4633 assert_eq!(outcome, UsageProviderRefreshOutcome::MissingToken);
4634 let guard = lb_states.lock().expect("lb states");
4635 assert!(
4636 !guard
4637 .get("right")
4638 .and_then(|entry| entry.usage_exhausted.get(1))
4639 .copied()
4640 .unwrap_or(true)
4641 );
4642 }
4643
4644 #[tokio::test]
4645 async fn usage_exhaustion_syncs_owned_balance_policy_action() {
4646 let cfg = proxy_config(vec![service_config(
4647 "right",
4648 vec![
4649 upstream_config("https://primary.example/v1"),
4650 upstream_config("https://backup.example/v1"),
4651 ],
4652 )]);
4653 let lb_states = Arc::new(Mutex::new(HashMap::new()));
4654 let upstreams = vec![endpoint_upstream()];
4655 let state = ProxyState::new();
4656
4657 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, true).await;
4658 let actions = state.list_policy_actions("codex").await;
4659 assert_eq!(actions.len(), 1);
4660 assert_eq!(actions[0].source_signal.kind, ProviderSignalKind::Balance);
4661 assert_eq!(actions[0].reason, "balance_exhausted");
4662
4663 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, false).await;
4664 assert!(state.list_policy_actions("codex").await.is_empty());
4665 }
4666
4667 #[test]
4668 fn openai_balance_snapshot_supports_custom_paths_and_derived_budget() {
4669 let mut provider = provider("custom", ProviderKind::OpenAiBalanceHttpJson);
4670 provider.extract.remaining_balance_paths = vec!["payload.remaining_usd".to_string()];
4671 provider.extract.monthly_spent_paths = vec!["payload.used_usd".to_string()];
4672 provider.extract.derive_budget_from_remaining_and_spent = true;
4673
4674 let snapshot = balance_http_snapshot_from_json(
4675 &provider,
4676 &upstream(),
4677 &serde_json::json!({
4678 "payload": {
4679 "remaining_usd": "2",
4680 "used_usd": "0.5"
4681 }
4682 }),
4683 100,
4684 Some(1_000),
4685 );
4686
4687 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("2"));
4688 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
4689 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("2.5"));
4690 assert_eq!(snapshot.exhausted, Some(false));
4691 }
4692
4693 #[test]
4694 fn new_api_snapshot_converts_quota_units_like_cc_switch_template() {
4695 let snapshot = new_api_snapshot_from_json(
4696 &provider("newapi", ProviderKind::NewApiUserSelf),
4697 &upstream(),
4698 &serde_json::json!({
4699 "success": true,
4700 "data": {
4701 "quota": 500000,
4702 "used_quota": 250000
4703 }
4704 }),
4705 100,
4706 Some(1_000),
4707 );
4708
4709 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4710 assert_eq!(snapshot.exhausted, Some(false));
4711 assert_eq!(snapshot.total_balance_usd, None);
4712 assert_eq!(snapshot.quota_period.as_deref(), Some("quota"));
4713 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("1"));
4714 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("1.5"));
4715 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("0.5"));
4716 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
4717 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("1.5"));
4718 }
4719
4720 #[test]
4721 fn new_api_user_self_honors_unlimited_quota_flag() {
4722 let snapshot = new_api_snapshot_from_json(
4723 &provider("newapi", ProviderKind::NewApiUserSelf),
4724 &upstream(),
4725 &serde_json::json!({
4726 "success": true,
4727 "data": {
4728 "quota": 0,
4729 "used_quota": 250000,
4730 "unlimited_quota": true
4731 }
4732 }),
4733 100,
4734 Some(1_000),
4735 );
4736
4737 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4738 assert_eq!(snapshot.exhausted, Some(false));
4739 assert_eq!(snapshot.total_balance_usd, None);
4740 assert_eq!(snapshot.monthly_budget_usd, None);
4741 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
4742 assert_eq!(snapshot.unlimited_quota, Some(true));
4743 }
4744
4745 #[test]
4746 fn new_api_token_usage_honors_unlimited_quota_flag() {
4747 let snapshot = new_api_token_usage_snapshot_from_json(
4748 &provider("newapi-token", ProviderKind::NewApiTokenUsage),
4749 &upstream(),
4750 &serde_json::json!({
4751 "code": true,
4752 "message": "ok",
4753 "data": {
4754 "object": "token_usage",
4755 "name": "demo-token",
4756 "total_granted": 0,
4757 "total_used": 250000,
4758 "total_available": 0,
4759 "unlimited_quota": true
4760 }
4761 }),
4762 100,
4763 Some(1_000),
4764 );
4765
4766 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4767 assert_eq!(snapshot.exhausted, Some(false));
4768 assert_eq!(snapshot.plan_name.as_deref(), Some("demo-token"));
4769 assert_eq!(snapshot.total_balance_usd, None);
4770 assert_eq!(snapshot.monthly_budget_usd, None);
4771 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
4772 assert_eq!(snapshot.unlimited_quota, Some(true));
4773 }
4774
4775 #[test]
4776 fn openai_organization_costs_sums_official_cost_buckets_without_exhaustion() {
4777 let snapshot = openai_organization_costs_snapshot_from_json(
4778 &provider("openai", ProviderKind::OpenAiOrganizationCosts),
4779 &upstream(),
4780 &serde_json::json!({
4781 "object": "page",
4782 "data": [
4783 {
4784 "object": "bucket",
4785 "start_time": 1710000000,
4786 "end_time": 1710086400,
4787 "results": [
4788 {
4789 "object": "organization.costs.result",
4790 "amount": { "value": 1.25, "currency": "usd" }
4791 },
4792 {
4793 "object": "organization.costs.result",
4794 "amount": { "value": "2.5", "currency": "usd" }
4795 },
4796 {
4797 "object": "organization.costs.result",
4798 "amount": { "value": 99, "currency": "eur" }
4799 }
4800 ]
4801 },
4802 {
4803 "object": "bucket",
4804 "results": [
4805 {
4806 "object": "organization.costs.result",
4807 "amount": { "value": "0.25", "currency": "USD" }
4808 }
4809 ]
4810 }
4811 ],
4812 "has_more": false
4813 }),
4814 100,
4815 Some(1_000),
4816 );
4817
4818 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4819 assert_eq!(snapshot.exhausted, None);
4820 assert!(!snapshot.exhaustion_affects_routing);
4821 assert!(!snapshot.routing_exhausted());
4822 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("4"));
4823 assert_eq!(snapshot.total_balance_usd, None);
4824 }
4825}