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