1use std::collections::{BTreeMap, HashMap, HashSet};
2use std::sync::{Arc, Mutex, OnceLock};
3use std::time::{Duration, Instant};
4
5use anyhow::{Context, Result};
6use futures_util::stream::{FuturesUnordered, StreamExt};
7use reqwest::Client;
8use serde::{Deserialize, Serialize};
9use tracing::{info, warn};
10
11use crate::balance::{
12 BalanceSnapshotStatus, ProviderBalanceSnapshot, ProviderUsageAlert, ProviderUsageAlertKind,
13 ProviderUsageModelStat, ProviderUsageRateSnapshot, ProviderUsageWindow,
14};
15use crate::config::{ProxyConfig, ServiceConfigManager, proxy_home_dir};
16use crate::lb::LbState;
17use crate::policy_actions::{PolicyAction, PolicyActionKind};
18use crate::pricing::UsdAmount;
19use crate::provider_signals::{
20 ProviderSignal, ProviderSignalKind, ProviderSignalSource, ProviderSignalTarget,
21};
22use crate::runtime_identity::ProviderEndpointKey;
23use crate::state::ProxyState;
24use crate::usage_forecast::next_reset_at_ms;
25
26#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq, Hash)]
27#[serde(rename_all = "snake_case")]
28enum ProviderKind {
29 BudgetHttpJson,
31 YescodeProfile,
33 #[serde(
35 rename = "openai_balance_http_json",
36 alias = "open_ai_balance_http_json",
37 alias = "relay_balance_http_json"
38 )]
39 OpenAiBalanceHttpJson,
40 #[serde(rename = "sub2api_usage", alias = "sub2api_usage_http_json")]
42 Sub2ApiUsage,
43 #[serde(rename = "sub2api_auth_me", alias = "sub2api_auth_me_http_json")]
45 Sub2ApiAuthMe,
46 #[serde(
48 rename = "new_api_token_usage",
49 alias = "new_api_token_usage_http_json"
50 )]
51 NewApiTokenUsage,
52 NewApiUserSelf,
54 #[serde(
56 rename = "rightcode_account_summary",
57 alias = "right_code_account_summary",
58 alias = "rightcode"
59 )]
60 RightCodeAccountSummary,
61 #[serde(
63 rename = "openai_organization_costs",
64 alias = "openai_org_costs",
65 alias = "openai_costs"
66 )]
67 OpenAiOrganizationCosts,
68}
69
70impl ProviderKind {
71 fn source_name(&self) -> &'static str {
72 match self {
73 ProviderKind::BudgetHttpJson => "usage_provider:budget_http_json",
74 ProviderKind::YescodeProfile => "usage_provider:yescode_profile",
75 ProviderKind::OpenAiBalanceHttpJson => "usage_provider:openai_balance_http_json",
76 ProviderKind::Sub2ApiUsage => "usage_provider:sub2api_usage",
77 ProviderKind::Sub2ApiAuthMe => "usage_provider:sub2api_auth_me",
78 ProviderKind::NewApiTokenUsage => "usage_provider:new_api_token_usage",
79 ProviderKind::NewApiUserSelf => "usage_provider:new_api_user_self",
80 ProviderKind::RightCodeAccountSummary => "usage_provider:rightcode_account_summary",
81 ProviderKind::OpenAiOrganizationCosts => "usage_provider:openai_organization_costs",
82 }
83 }
84
85 fn default_endpoint(&self) -> Option<&'static str> {
86 match self {
87 ProviderKind::OpenAiBalanceHttpJson => Some("{{base_url}}/user/balance"),
88 ProviderKind::Sub2ApiUsage => Some("{{base_url}}/v1/usage"),
89 ProviderKind::Sub2ApiAuthMe => Some("{{base_url}}/api/v1/auth/me"),
90 ProviderKind::NewApiTokenUsage => Some("{{base_url}}/api/usage/token/"),
91 ProviderKind::NewApiUserSelf => Some("{{base_url}}/api/user/self"),
92 ProviderKind::RightCodeAccountSummary => {
93 Some("https://www.right.codes/account/summary")
94 }
95 ProviderKind::OpenAiOrganizationCosts => {
96 Some("{{base_url}}/v1/organization/costs?start_time={{unix_days_ago:30}}&limit=30")
97 }
98 _ => None,
99 }
100 }
101}
102
103#[derive(Debug, Deserialize, Serialize, Default, Clone)]
104#[serde(default)]
105struct UsageProviderExtractConfig {
106 #[serde(skip_serializing_if = "Vec::is_empty")]
107 remaining_balance_paths: Vec<String>,
108 #[serde(skip_serializing_if = "Vec::is_empty")]
109 subscription_balance_paths: Vec<String>,
110 #[serde(skip_serializing_if = "Vec::is_empty")]
111 paygo_balance_paths: Vec<String>,
112 #[serde(skip_serializing_if = "Vec::is_empty")]
113 monthly_budget_paths: Vec<String>,
114 #[serde(skip_serializing_if = "Vec::is_empty")]
115 monthly_spent_paths: Vec<String>,
116 #[serde(skip_serializing_if = "Vec::is_empty")]
117 exhausted_paths: Vec<String>,
118 #[serde(skip_serializing_if = "Option::is_none")]
119 remaining_divisor: Option<u64>,
120 #[serde(skip_serializing_if = "Option::is_none")]
121 monthly_budget_divisor: Option<u64>,
122 #[serde(skip_serializing_if = "Option::is_none")]
123 monthly_spent_divisor: Option<u64>,
124 #[serde(skip_serializing_if = "bool_is_false")]
125 derive_budget_from_remaining_and_spent: bool,
126 #[serde(skip_serializing_if = "bool_is_false")]
127 derive_remaining_from_budget_and_spent: bool,
128}
129
130impl UsageProviderExtractConfig {
131 fn is_empty(&self) -> bool {
132 self.remaining_balance_paths.is_empty()
133 && self.subscription_balance_paths.is_empty()
134 && self.paygo_balance_paths.is_empty()
135 && self.monthly_budget_paths.is_empty()
136 && self.monthly_spent_paths.is_empty()
137 && self.exhausted_paths.is_empty()
138 && self.remaining_divisor.is_none()
139 && self.monthly_budget_divisor.is_none()
140 && self.monthly_spent_divisor.is_none()
141 && !self.derive_budget_from_remaining_and_spent
142 && !self.derive_remaining_from_budget_and_spent
143 }
144}
145
146#[derive(Debug, Deserialize, Serialize)]
147struct UsageProviderConfig {
148 id: String,
149 kind: ProviderKind,
150 domains: Vec<String>,
151 #[serde(default)]
152 endpoint: String,
153 #[serde(default)]
154 token_env: Option<String>,
155 #[serde(default, skip_serializing_if = "bool_is_false")]
156 require_token_env: bool,
157 #[serde(default)]
158 poll_interval_secs: Option<u64>,
159 #[serde(
160 default = "default_refresh_on_request",
161 skip_serializing_if = "bool_is_true"
162 )]
163 refresh_on_request: bool,
164 #[serde(
165 default = "default_trust_exhaustion_for_routing",
166 skip_serializing_if = "bool_is_true"
167 )]
168 trust_exhaustion_for_routing: bool,
169 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
170 headers: BTreeMap<String, String>,
171 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
172 variables: BTreeMap<String, String>,
173 #[serde(default, skip_serializing_if = "UsageProviderExtractConfig::is_empty")]
174 extract: UsageProviderExtractConfig,
175}
176
177#[derive(Debug, Deserialize, Serialize, Default)]
178struct UsageProvidersFile {
179 #[serde(default)]
180 providers: Vec<UsageProviderConfig>,
181}
182
183#[derive(Debug, Clone)]
184struct UpstreamRef {
185 station_name: String,
186 index: usize,
187 provider_endpoint: Option<ProviderEndpointKey>,
188}
189
190#[derive(Debug, Clone)]
191struct UsageProviderTarget {
192 upstream: UpstreamRef,
193 base_url: String,
194 provider_id: Option<String>,
195}
196
197#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
198struct UsageProviderTargetKey {
199 station_name: String,
200 upstream_index: usize,
201}
202
203#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
204pub struct UsageProviderRefreshSummary {
205 pub providers_configured: usize,
206 pub providers_matched: usize,
207 pub upstreams_matched: usize,
208 pub attempted: usize,
209 pub refreshed: usize,
210 pub failed: usize,
211 pub missing_token: usize,
212 #[serde(skip_serializing_if = "usize_is_zero")]
213 pub auto_attempted: usize,
214 #[serde(skip_serializing_if = "usize_is_zero")]
215 pub auto_refreshed: usize,
216 #[serde(skip_serializing_if = "usize_is_zero")]
217 pub auto_failed: usize,
218 #[serde(skip_serializing_if = "usize_is_zero")]
219 pub deduplicated: usize,
220}
221
222#[derive(Debug, Clone, Copy, PartialEq, Eq)]
223enum UsageProviderRefreshOutcome {
224 Refreshed,
225 Failed,
226 MissingToken,
227}
228
229struct RefreshProviderTargetParams<'a> {
230 client: &'a Client,
231 provider: &'a UsageProviderConfig,
232 target: &'a UsageProviderTarget,
233 cfg: &'a ProxyConfig,
234 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
235 state: &'a Arc<ProxyState>,
236 service_name: &'a str,
237 interval_secs: u64,
238}
239
240static LAST_USAGE_POLL: OnceLock<Mutex<HashMap<String, Instant>>> = OnceLock::new();
242static REQUEST_BALANCE_QUEUE: OnceLock<Mutex<HashMap<ProviderEndpointKey, Instant>>> =
243 OnceLock::new();
244static AUTO_PROBE_KIND_HINTS: OnceLock<Mutex<HashMap<String, ProviderKind>>> = OnceLock::new();
245static AUTO_PROBE_KIND_FAILURES: OnceLock<Mutex<HashMap<AutoProbeKindFailureKey, Instant>>> =
246 OnceLock::new();
247static USAGE_PROVIDER_TARGET_SUPPRESSIONS: OnceLock<
248 Mutex<HashMap<ProviderTargetSuppressionKey, ProviderTargetSuppression>>,
249> = OnceLock::new();
250
251const DEFAULT_POLL_INTERVAL_SECS: u64 = 10 * 60;
252const MIN_POLL_INTERVAL_SECS: u64 = 2 * 60;
254pub const REQUEST_BALANCE_REFRESH_DELAY: Duration = Duration::from_secs(60);
255const BALANCE_REFRESH_CONCURRENCY: usize = 6;
256const BALANCE_HTTP_REQUEST_TIMEOUT: Duration = Duration::from_secs(6);
257const BALANCE_HTTP_ERROR_BODY_LIMIT: usize = 2_048;
258const AUTO_PROBE_KIND_FAILURE_TTL: Duration = Duration::from_secs(10 * 60);
259const USAGE_PROVIDER_TERMINAL_FAILURE_TTL: Duration = Duration::from_secs(6 * 60 * 60);
260const USAGE_PROVIDER_EXHAUSTED_SUPPRESSION_TTL: Duration = Duration::from_secs(6 * 60 * 60);
261const USAGE_PROVIDER_DAILY_RESET_SUPPRESSION_GRACE: Duration = Duration::from_secs(5 * 60);
262const LOW_BALANCE_ALERT_THRESHOLD_USD: &str = "10";
263const EXPIRING_SOON_WINDOW_SECS: u64 = 7 * 24 * 60 * 60;
264const AUTO_PROVIDER_ID_PREFIX: &str = "auto:balance:";
265const AUTO_PROBE_KINDS: [ProviderKind; 5] = [
266 ProviderKind::RightCodeAccountSummary,
267 ProviderKind::Sub2ApiUsage,
268 ProviderKind::NewApiTokenUsage,
269 ProviderKind::NewApiUserSelf,
270 ProviderKind::OpenAiBalanceHttpJson,
271];
272
273#[derive(Debug, Clone, PartialEq, Eq, Hash)]
274struct AutoProbeKindFailureKey {
275 provider_id: String,
276 target: AutoProbeTargetKey,
277 kind: ProviderKind,
278}
279
280#[derive(Debug, Clone, PartialEq, Eq, Hash)]
281struct ProviderTargetSuppressionKey {
282 provider_id: String,
283 target: AutoProbeTargetKey,
284}
285
286#[derive(Debug, Clone)]
287struct ProviderTargetSuppression {
288 until: Instant,
289 reason: String,
290 routing_exhausted: bool,
291}
292
293#[derive(Debug, Clone)]
294struct ProviderTargetSuppressionDecision {
295 reason: String,
296 routing_exhausted: bool,
297 ttl: Duration,
298}
299
300#[derive(Debug, Clone, PartialEq, Eq, Hash)]
301struct AutoProbeTargetKey {
302 station_name: String,
303 upstream_index: usize,
304 provider_endpoint_key: Option<String>,
305 base_url: String,
306}
307
308#[derive(Debug, Clone, Copy, PartialEq, Eq)]
309enum RequestBalanceQueueDue {
310 Due,
311 NotDue(Duration),
312 Missing,
313}
314
315#[derive(Debug, Clone, Copy, PartialEq, Eq)]
316enum RequestBalancePollOutcome {
317 Attempted,
318 Deferred(Duration),
319 Skipped,
320}
321
322fn bool_is_false(value: &bool) -> bool {
323 !*value
324}
325
326fn bool_is_true(value: &bool) -> bool {
327 *value
328}
329
330fn usize_is_zero(value: &usize) -> bool {
331 *value == 0
332}
333
334fn default_refresh_on_request() -> bool {
335 true
336}
337
338fn default_trust_exhaustion_for_routing() -> bool {
339 true
340}
341
342fn unix_now_ms() -> u64 {
343 std::time::SystemTime::now()
344 .duration_since(std::time::UNIX_EPOCH)
345 .map(|d| d.as_millis() as u64)
346 .unwrap_or(0)
347}
348
349fn unix_now_secs() -> u64 {
350 std::time::SystemTime::now()
351 .duration_since(std::time::UNIX_EPOCH)
352 .map(|d| d.as_secs())
353 .unwrap_or(0)
354}
355
356fn stale_after_ms(fetched_at_ms: u64, interval_secs: u64) -> Option<u64> {
357 fetched_at_ms.checked_add(interval_secs.saturating_mul(3).saturating_mul(1_000))
358}
359
360fn snapshot_refresh_interval_secs(provider: &UsageProviderConfig) -> u64 {
361 let interval_secs = provider
362 .poll_interval_secs
363 .unwrap_or(DEFAULT_POLL_INTERVAL_SECS);
364 if interval_secs == 0 {
365 DEFAULT_POLL_INTERVAL_SECS
366 } else {
367 interval_secs.max(MIN_POLL_INTERVAL_SECS)
368 }
369}
370
371fn effective_poll_interval_secs(provider: &UsageProviderConfig) -> Option<u64> {
372 if !provider.refresh_on_request {
373 return None;
374 }
375
376 let interval_secs = provider
377 .poll_interval_secs
378 .unwrap_or(DEFAULT_POLL_INTERVAL_SECS);
379 if interval_secs == 0 {
380 return None;
381 }
382 Some(interval_secs.max(MIN_POLL_INTERVAL_SECS))
383}
384
385fn remaining_poll_cooldown(last: Instant, interval_secs: u64, now: Instant) -> Option<Duration> {
386 let interval = Duration::from_secs(interval_secs);
387 let elapsed = now.saturating_duration_since(last);
388 interval.checked_sub(elapsed).filter(|d| !d.is_zero())
389}
390
391fn usage_providers_path() -> std::path::PathBuf {
392 proxy_home_dir().join("usage_providers.json")
393}
394
395fn service_manager<'a>(cfg: &'a ProxyConfig, service_name: &str) -> &'a ServiceConfigManager {
396 match service_name {
397 "claude" => &cfg.claude,
398 _ => &cfg.codex,
399 }
400}
401
402fn default_provider_config(
403 id: &str,
404 kind: ProviderKind,
405 domains: Vec<&str>,
406 endpoint: &str,
407 extract: UsageProviderExtractConfig,
408) -> UsageProviderConfig {
409 UsageProviderConfig {
410 id: id.to_string(),
411 kind,
412 domains: domains.into_iter().map(str::to_string).collect(),
413 endpoint: endpoint.to_string(),
414 token_env: None,
415 require_token_env: false,
416 poll_interval_secs: Some(DEFAULT_POLL_INTERVAL_SECS),
417 refresh_on_request: true,
418 trust_exhaustion_for_routing: true,
419 headers: BTreeMap::new(),
420 variables: BTreeMap::new(),
421 extract,
422 }
423}
424
425fn default_rightcode_provider_config(id: &str) -> UsageProviderConfig {
426 let mut provider = default_provider_config(
427 id,
428 ProviderKind::RightCodeAccountSummary,
429 vec!["www.right.codes", "right.codes"],
430 "https://www.right.codes/account/summary",
431 UsageProviderExtractConfig::default(),
432 );
433 provider.trust_exhaustion_for_routing = false;
437 provider
438}
439
440fn host_from_base_url(base_url: &str) -> Option<String> {
441 reqwest::Url::parse(base_url)
442 .ok()
443 .and_then(|url| url.host_str().map(|host| host.to_ascii_lowercase()))
444}
445
446fn is_official_openai_base_url(base_url: &str) -> bool {
447 host_from_base_url(base_url).as_deref() == Some("api.openai.com")
448}
449
450fn is_rightcode_base_url(base_url: &str) -> bool {
451 matches!(
452 host_from_base_url(base_url).as_deref(),
453 Some("www.right.codes" | "right.codes")
454 )
455}
456
457fn provider_id_component(value: &str) -> String {
458 let component = value
459 .chars()
460 .map(|ch| {
461 if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' || ch == '.' {
462 ch
463 } else {
464 '-'
465 }
466 })
467 .collect::<String>()
468 .trim_matches('-')
469 .to_string();
470 if component.is_empty() {
471 "station".to_string()
472 } else {
473 component
474 }
475}
476
477fn auto_provider_id(target: &UsageProviderTarget) -> String {
478 if let Some(provider_id) = target
479 .provider_id
480 .as_deref()
481 .map(str::trim)
482 .filter(|value| !value.is_empty())
483 {
484 return provider_id.to_string();
485 }
486 format!(
487 "{}{}:{}",
488 AUTO_PROVIDER_ID_PREFIX,
489 provider_id_component(&target.upstream.station_name),
490 target.upstream.index
491 )
492}
493
494fn auto_usage_provider(target: &UsageProviderTarget, kind: ProviderKind) -> UsageProviderConfig {
495 let mut provider = UsageProviderConfig {
496 id: auto_provider_id(target),
497 kind,
498 domains: host_from_base_url(&target.base_url)
499 .into_iter()
500 .collect::<Vec<_>>(),
501 endpoint: String::new(),
502 token_env: None,
503 require_token_env: false,
504 poll_interval_secs: Some(DEFAULT_POLL_INTERVAL_SECS),
505 refresh_on_request: true,
506 trust_exhaustion_for_routing: true,
507 headers: BTreeMap::new(),
508 variables: BTreeMap::new(),
509 extract: UsageProviderExtractConfig::default(),
510 };
511 if matches!(kind, ProviderKind::RightCodeAccountSummary) {
512 provider.trust_exhaustion_for_routing = false;
513 }
514 provider
515}
516
517fn auto_target_matches_provider_id_filter(
518 target: &UsageProviderTarget,
519 provider_id_filter: Option<&str>,
520) -> bool {
521 match provider_id_filter {
522 Some(filter) => auto_provider_id(target) == filter,
523 None => true,
524 }
525}
526
527fn first_auto_probe_kind(target: &UsageProviderTarget) -> ProviderKind {
528 if is_rightcode_base_url(&target.base_url) {
529 ProviderKind::RightCodeAccountSummary
530 } else {
531 ProviderKind::Sub2ApiUsage
532 }
533}
534
535fn auto_probe_target_key(target: &UsageProviderTarget) -> AutoProbeTargetKey {
536 AutoProbeTargetKey {
537 station_name: target.upstream.station_name.clone(),
538 upstream_index: target.upstream.index,
539 provider_endpoint_key: target
540 .upstream
541 .provider_endpoint
542 .as_ref()
543 .map(ProviderEndpointKey::stable_key),
544 base_url: normalized_balance_base_url(&target.base_url)
545 .unwrap_or_else(|| target.base_url.clone()),
546 }
547}
548
549fn auto_probe_kind_order(provider_id: &str, target: &UsageProviderTarget) -> Vec<ProviderKind> {
550 let now = Instant::now();
551 let target_key = auto_probe_target_key(target);
552 let mut ordered = Vec::new();
553 if let Some(kind) = remembered_auto_probe_kind(provider_id) {
554 ordered.push(kind);
555 }
556 ordered.push(first_auto_probe_kind(target));
557 ordered.extend(AUTO_PROBE_KINDS);
558
559 let mut seen = HashSet::new();
560 ordered
561 .into_iter()
562 .filter(|kind| {
563 if *kind == ProviderKind::RightCodeAccountSummary
564 && !is_rightcode_base_url(&target.base_url)
565 {
566 return false;
567 }
568 seen.insert(*kind)
569 && !auto_probe_kind_failure_active(provider_id, &target_key, *kind, now)
570 })
571 .collect()
572}
573
574fn remembered_auto_probe_kind(provider_id: &str) -> Option<ProviderKind> {
575 AUTO_PROBE_KIND_HINTS
576 .get_or_init(|| Mutex::new(HashMap::new()))
577 .lock()
578 .ok()
579 .and_then(|map| map.get(provider_id).copied())
580}
581
582fn remember_auto_probe_kind_success(
583 provider_id: &str,
584 target: &UsageProviderTarget,
585 kind: ProviderKind,
586) {
587 if let Ok(mut hints) = AUTO_PROBE_KIND_HINTS
588 .get_or_init(|| Mutex::new(HashMap::new()))
589 .lock()
590 {
591 hints.insert(provider_id.to_string(), kind);
592 }
593 if let Ok(mut failures) = AUTO_PROBE_KIND_FAILURES
594 .get_or_init(|| Mutex::new(HashMap::new()))
595 .lock()
596 {
597 failures.remove(&AutoProbeKindFailureKey {
598 provider_id: provider_id.to_string(),
599 target: auto_probe_target_key(target),
600 kind,
601 });
602 }
603 clear_usage_provider_target_suppression(provider_id, target);
604}
605
606fn remember_auto_probe_kind_failure(
607 provider_id: &str,
608 target: &UsageProviderTarget,
609 kind: ProviderKind,
610 now: Instant,
611) {
612 if let Ok(mut failures) = AUTO_PROBE_KIND_FAILURES
613 .get_or_init(|| Mutex::new(HashMap::new()))
614 .lock()
615 {
616 failures.insert(
617 AutoProbeKindFailureKey {
618 provider_id: provider_id.to_string(),
619 target: auto_probe_target_key(target),
620 kind,
621 },
622 now,
623 );
624 }
625}
626
627fn auto_probe_kind_failure_active(
628 provider_id: &str,
629 target: &AutoProbeTargetKey,
630 kind: ProviderKind,
631 now: Instant,
632) -> bool {
633 let key = AutoProbeKindFailureKey {
634 provider_id: provider_id.to_string(),
635 target: target.clone(),
636 kind,
637 };
638 let Ok(mut failures) = AUTO_PROBE_KIND_FAILURES
639 .get_or_init(|| Mutex::new(HashMap::new()))
640 .lock()
641 else {
642 return false;
643 };
644 let Some(failed_at) = failures.get(&key).copied() else {
645 return false;
646 };
647 if now.duration_since(failed_at) < AUTO_PROBE_KIND_FAILURE_TTL {
648 true
649 } else {
650 failures.remove(&key);
651 false
652 }
653}
654
655fn usage_provider_target_suppression_key(
656 provider_id: &str,
657 target: &UsageProviderTarget,
658) -> ProviderTargetSuppressionKey {
659 ProviderTargetSuppressionKey {
660 provider_id: provider_id.to_string(),
661 target: auto_probe_target_key(target),
662 }
663}
664
665fn remember_usage_provider_target_suppression(
666 provider_id: &str,
667 target: &UsageProviderTarget,
668 ttl: Duration,
669 reason: impl Into<String>,
670 routing_exhausted: bool,
671 now: Instant,
672) {
673 if let Ok(mut suppressions) = USAGE_PROVIDER_TARGET_SUPPRESSIONS
674 .get_or_init(|| Mutex::new(HashMap::new()))
675 .lock()
676 {
677 suppressions.insert(
678 usage_provider_target_suppression_key(provider_id, target),
679 ProviderTargetSuppression {
680 until: now + ttl,
681 reason: reason.into(),
682 routing_exhausted,
683 },
684 );
685 }
686}
687
688fn clear_usage_provider_target_suppression(provider_id: &str, target: &UsageProviderTarget) {
689 if let Ok(mut suppressions) = USAGE_PROVIDER_TARGET_SUPPRESSIONS
690 .get_or_init(|| Mutex::new(HashMap::new()))
691 .lock()
692 {
693 suppressions.remove(&usage_provider_target_suppression_key(provider_id, target));
694 }
695}
696
697#[cfg(test)]
698fn clear_usage_provider_target_suppressions_for_provider(provider_id: &str) {
699 if let Some(suppressions) = USAGE_PROVIDER_TARGET_SUPPRESSIONS.get()
700 && let Ok(mut suppressions) = suppressions.lock()
701 {
702 suppressions.retain(|key, _| key.provider_id != provider_id);
703 }
704}
705
706fn usage_provider_target_suppression_active(
707 provider_id: &str,
708 target: &UsageProviderTarget,
709 now: Instant,
710) -> Option<ProviderTargetSuppression> {
711 let key = usage_provider_target_suppression_key(provider_id, target);
712 let Ok(mut suppressions) = USAGE_PROVIDER_TARGET_SUPPRESSIONS
713 .get_or_init(|| Mutex::new(HashMap::new()))
714 .lock()
715 else {
716 return None;
717 };
718 let suppression = suppressions.get(&key).cloned()?;
719 if now < suppression.until {
720 Some(suppression)
721 } else {
722 suppressions.remove(&key);
723 None
724 }
725}
726
727fn auto_openai_official_provider(target: &UsageProviderTarget) -> UsageProviderConfig {
728 let mut provider = auto_usage_provider(target, ProviderKind::OpenAiOrganizationCosts);
729 provider.token_env = Some("OPENAI_ADMIN_KEY".to_string());
730 provider.require_token_env = true;
731 provider.refresh_on_request = false;
732 provider.trust_exhaustion_for_routing = false;
733 provider
734}
735
736fn default_providers() -> UsageProvidersFile {
737 let openrouter_extract = UsageProviderExtractConfig {
738 monthly_budget_paths: vec!["data.total_credits".to_string()],
739 monthly_spent_paths: vec!["data.total_usage".to_string()],
740 derive_remaining_from_budget_and_spent: true,
741 ..Default::default()
742 };
743
744 let novita_extract = UsageProviderExtractConfig {
745 remaining_balance_paths: vec!["availableBalance".to_string()],
746 remaining_divisor: Some(10_000),
747 ..Default::default()
748 };
749
750 let mut openai_official = default_provider_config(
751 "openai-official-costs",
752 ProviderKind::OpenAiOrganizationCosts,
753 vec!["api.openai.com"],
754 "https://api.openai.com/v1/organization/costs?start_time={{unix_days_ago:30}}&limit=30",
755 UsageProviderExtractConfig::default(),
756 );
757 openai_official.token_env = Some("OPENAI_ADMIN_KEY".to_string());
758 openai_official.require_token_env = true;
759 openai_official.refresh_on_request = false;
760 openai_official.trust_exhaustion_for_routing = false;
761
762 UsageProvidersFile {
763 providers: vec![
764 default_rightcode_provider_config("rightcode"),
765 default_provider_config(
766 "packycode",
767 ProviderKind::BudgetHttpJson,
768 vec!["packycode.com"],
769 "https://www.packycode.com/api/backend/users/info",
770 UsageProviderExtractConfig::default(),
771 ),
772 default_provider_config(
773 "yescode",
774 ProviderKind::YescodeProfile,
775 vec!["yes.vg"],
777 "https://co.yes.vg/api/v1/auth/profile",
778 UsageProviderExtractConfig::default(),
779 ),
780 default_provider_config(
781 "deepseek",
782 ProviderKind::OpenAiBalanceHttpJson,
783 vec!["api.deepseek.com"],
784 "https://api.deepseek.com/user/balance",
785 UsageProviderExtractConfig::default(),
786 ),
787 default_provider_config(
788 "stepfun",
789 ProviderKind::OpenAiBalanceHttpJson,
790 vec!["api.stepfun.ai", "api.stepfun.com"],
791 "https://api.stepfun.com/v1/accounts",
792 UsageProviderExtractConfig::default(),
793 ),
794 default_provider_config(
795 "siliconflow",
796 ProviderKind::OpenAiBalanceHttpJson,
797 vec!["api.siliconflow.cn", "api.siliconflow.com"],
798 "{{base_url}}/v1/user/info",
799 UsageProviderExtractConfig::default(),
800 ),
801 default_provider_config(
802 "openrouter",
803 ProviderKind::OpenAiBalanceHttpJson,
804 vec!["openrouter.ai"],
805 "https://openrouter.ai/api/v1/credits",
806 openrouter_extract,
807 ),
808 default_provider_config(
809 "novita",
810 ProviderKind::OpenAiBalanceHttpJson,
811 vec!["api.novita.ai"],
812 "https://api.novita.ai/v3/user/balance",
813 novita_extract,
814 ),
815 openai_official,
816 ],
817 }
818}
819
820fn load_providers() -> UsageProvidersFile {
821 let path = usage_providers_path();
822 if let Ok(text) = std::fs::read_to_string(&path)
823 && let Ok(file) = serde_json::from_str::<UsageProvidersFile>(&text)
824 {
825 return file;
826 }
827
828 let default = default_providers();
830 if let Ok(text) = serde_json::to_string_pretty(&default) {
831 if let Some(parent) = path.parent() {
832 let _ = std::fs::create_dir_all(parent);
833 }
834 let _ = std::fs::write(&path, text);
835 }
836 default
837}
838
839fn domain_matches(base_url: &str, domains: &[String]) -> bool {
840 let url = match reqwest::Url::parse(base_url) {
841 Ok(u) => u,
842 Err(_) => return false,
843 };
844 let host = match url.host_str() {
845 Some(h) => h,
846 None => return false,
847 };
848 let host = host.to_ascii_lowercase();
849 for d in domains {
850 let domain = d.trim().to_ascii_lowercase();
851 if host == domain || host.ends_with(&format!(".{}", domain)) {
852 return true;
853 }
854 }
855 false
856}
857
858fn matching_provider_targets(
859 cfg: &ProxyConfig,
860 service_name: &str,
861 provider: &UsageProviderConfig,
862 station_name_filter: Option<&str>,
863) -> Vec<UsageProviderTarget> {
864 let mut stations: Vec<_> = service_manager(cfg, service_name)
865 .stations()
866 .iter()
867 .collect();
868 stations.sort_by_key(|(name, _)| name.as_str());
869
870 let mut targets = Vec::new();
871 for (station_name, service) in stations {
872 if station_name_filter.is_some_and(|filter| filter != station_name.as_str()) {
873 continue;
874 }
875 for (index, upstream) in service.upstreams.iter().enumerate() {
876 if domain_matches(&upstream.base_url, &provider.domains) {
877 targets.push(UsageProviderTarget {
878 upstream: UpstreamRef {
879 station_name: station_name.clone(),
880 index,
881 provider_endpoint: upstream.provider_endpoint_key(service_name),
882 },
883 base_url: upstream.base_url.clone(),
884 provider_id: upstream.tags.get("provider_id").cloned(),
885 });
886 }
887 }
888 }
889
890 targets
891}
892
893fn usage_provider_targets(
894 cfg: &ProxyConfig,
895 service_name: &str,
896 station_name_filter: Option<&str>,
897) -> Vec<UsageProviderTarget> {
898 let mut stations: Vec<_> = service_manager(cfg, service_name)
899 .stations()
900 .iter()
901 .collect();
902 stations.sort_by_key(|(name, _)| name.as_str());
903
904 let mut targets = Vec::new();
905 for (station_name, service) in stations {
906 if station_name_filter.is_some_and(|filter| filter != station_name.as_str()) {
907 continue;
908 }
909 for (index, upstream) in service.upstreams.iter().enumerate() {
910 targets.push(UsageProviderTarget {
911 upstream: UpstreamRef {
912 station_name: station_name.clone(),
913 index,
914 provider_endpoint: upstream.provider_endpoint_key(service_name),
915 },
916 base_url: upstream.base_url.clone(),
917 provider_id: upstream.tags.get("provider_id").cloned(),
918 });
919 }
920 }
921
922 targets
923}
924
925fn target_key(target: &UsageProviderTarget) -> UsageProviderTargetKey {
926 UsageProviderTargetKey {
927 station_name: target.upstream.station_name.clone(),
928 upstream_index: target.upstream.index,
929 }
930}
931
932fn enqueue_request_balance_refresh(key: ProviderEndpointKey) -> Option<Duration> {
933 let now = Instant::now();
934 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
935 let mut queue = match queue.lock() {
936 Ok(queue) => queue,
937 Err(_) => return None,
938 };
939
940 match queue.get(&key).copied() {
941 Some(due_at) if due_at > now => None,
942 Some(_) => Some(Duration::ZERO),
943 None => {
944 queue.insert(key, now + REQUEST_BALANCE_REFRESH_DELAY);
945 Some(REQUEST_BALANCE_REFRESH_DELAY)
946 }
947 }
948}
949
950fn schedule_request_balance_refresh_at(key: ProviderEndpointKey, due_at: Instant) {
951 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
952 if let Ok(mut queue) = queue.lock() {
953 queue.insert(key, due_at);
954 }
955}
956
957fn take_request_balance_refresh_if_due(key: &ProviderEndpointKey) -> RequestBalanceQueueDue {
958 let now = Instant::now();
959 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
960 let mut queue = match queue.lock() {
961 Ok(queue) => queue,
962 Err(_) => return RequestBalanceQueueDue::Missing,
963 };
964
965 match queue.get(key).copied() {
966 Some(due_at) if due_at <= now => {
967 queue.remove(key);
968 RequestBalanceQueueDue::Due
969 }
970 Some(due_at) => RequestBalanceQueueDue::NotDue(due_at.saturating_duration_since(now)),
971 None => RequestBalanceQueueDue::Missing,
972 }
973}
974
975#[cfg(test)]
976pub fn request_balance_refresh_queued_for_provider_endpoint(
977 provider_endpoint: &ProviderEndpointKey,
978) -> bool {
979 let Some(queue) = REQUEST_BALANCE_QUEUE.get() else {
980 return false;
981 };
982 match queue.lock() {
983 Ok(guard) => guard.contains_key(provider_endpoint),
984 Err(error) => error.into_inner().contains_key(provider_endpoint),
985 }
986}
987
988fn usage_provider_target_for_provider_endpoint(
989 cfg: &ProxyConfig,
990 service_name: &str,
991 provider_endpoint: &ProviderEndpointKey,
992) -> Option<UsageProviderTarget> {
993 service_manager(cfg, service_name)
994 .stations()
995 .iter()
996 .filter_map(|(station_name, service)| {
997 service
998 .upstreams
999 .iter()
1000 .enumerate()
1001 .find_map(|(index, upstream)| {
1002 let upstream_endpoint = upstream.provider_endpoint_key(service_name)?;
1003 if upstream_endpoint != *provider_endpoint {
1004 return None;
1005 }
1006 Some(UsageProviderTarget {
1007 upstream: UpstreamRef {
1008 station_name: station_name.clone(),
1009 index,
1010 provider_endpoint: Some(upstream_endpoint),
1011 },
1012 base_url: upstream.base_url.clone(),
1013 provider_id: upstream.tags.get("provider_id").cloned(),
1014 })
1015 })
1016 })
1017 .next()
1018}
1019
1020trait UsageProviderUpstreamIdentityExt {
1021 fn provider_endpoint_key(&self, service_name: &str) -> Option<ProviderEndpointKey>;
1022}
1023
1024impl UsageProviderUpstreamIdentityExt for crate::config::UpstreamConfig {
1025 fn provider_endpoint_key(&self, service_name: &str) -> Option<ProviderEndpointKey> {
1026 let provider_id = self.tags.get("provider_id")?.trim();
1027 let endpoint_id = self.tags.get("endpoint_id")?.trim();
1028 if provider_id.is_empty() || endpoint_id.is_empty() {
1029 return None;
1030 }
1031 Some(ProviderEndpointKey::new(
1032 service_name.to_string(),
1033 provider_id.to_string(),
1034 endpoint_id.to_string(),
1035 ))
1036 }
1037}
1038
1039#[cfg(test)]
1040fn configured_target_keys(
1041 cfg: &ProxyConfig,
1042 service_name: &str,
1043 providers: &[UsageProviderConfig],
1044 station_name_filter: Option<&str>,
1045) -> HashSet<UsageProviderTargetKey> {
1046 providers
1047 .iter()
1048 .flat_map(|provider| {
1049 matching_provider_targets(cfg, service_name, provider, station_name_filter)
1050 })
1051 .map(|target| target_key(&target))
1052 .collect()
1053}
1054
1055fn resolve_token(
1056 provider: &UsageProviderConfig,
1057 upstreams: &[UpstreamRef],
1058 cfg: &ProxyConfig,
1059 service_name: &str,
1060) -> Option<String> {
1061 if let Some(env_name) = &provider.token_env
1063 && let Ok(v) = std::env::var(env_name)
1064 && !v.trim().is_empty()
1065 {
1066 return Some(v);
1067 }
1068
1069 if provider.require_token_env {
1070 return None;
1071 }
1072
1073 for uref in upstreams {
1075 if let Some(service) = service_manager(cfg, service_name).station(&uref.station_name)
1076 && let Some(up) = service.upstreams.get(uref.index)
1077 {
1078 if let Some(token) = up.auth.resolve_auth_token() {
1079 return Some(token);
1080 }
1081 if let Some(token) = up.auth.resolve_api_key() {
1082 return Some(token);
1083 }
1084 }
1085 }
1086 None
1087}
1088
1089fn normalized_balance_base_url(base_url: &str) -> Option<String> {
1090 let mut url = reqwest::Url::parse(base_url).ok()?;
1091 url.set_query(None);
1092 url.set_fragment(None);
1093 let path = url.path().trim_end_matches('/').to_string();
1094 if path.eq_ignore_ascii_case("/v1") {
1095 url.set_path("");
1096 } else if path.to_ascii_lowercase().ends_with("/v1") {
1097 let new_path = &path[..path.len().saturating_sub(3)];
1098 url.set_path(if new_path.is_empty() { "/" } else { new_path });
1099 }
1100 Some(url.as_str().trim_end_matches('/').to_string())
1101}
1102
1103fn base_path_prefixes(base_url: &str) -> Vec<String> {
1104 let Some(normalized) = normalized_balance_base_url(base_url) else {
1105 return Vec::new();
1106 };
1107 let Ok(url) = reqwest::Url::parse(&normalized) else {
1108 return Vec::new();
1109 };
1110 let parts = url
1111 .path_segments()
1112 .map(|segments| {
1113 segments
1114 .filter(|segment| !segment.is_empty())
1115 .collect::<Vec<_>>()
1116 })
1117 .unwrap_or_default();
1118 let mut prefixes = Vec::new();
1119 for len in (1..=parts.len()).rev() {
1120 prefixes.push(format!("/{}", parts[..len].join("/")));
1121 }
1122 if prefixes.is_empty() {
1123 prefixes.push("/".to_string());
1124 }
1125 prefixes
1126}
1127
1128fn path_prefixes_match(provider_prefixes: &[String], available_prefixes: &[String]) -> bool {
1129 if provider_prefixes.is_empty() || available_prefixes.is_empty() {
1130 return false;
1131 }
1132 provider_prefixes.iter().any(|provider_prefix| {
1133 available_prefixes.iter().any(|available_prefix| {
1134 provider_prefix == available_prefix
1135 || provider_prefix
1136 .strip_prefix(available_prefix)
1137 .is_some_and(|suffix| suffix.starts_with('/'))
1138 })
1139 })
1140}
1141
1142fn render_provider_template(
1143 template: &str,
1144 base_url: &str,
1145 upstream_base_url: &str,
1146 token: &str,
1147 variables: &BTreeMap<String, String>,
1148) -> String {
1149 let mut out = template
1150 .replace("{{baseUrl}}", base_url)
1151 .replace("{{base_url}}", base_url)
1152 .replace("{{upstreamBaseUrl}}", upstream_base_url)
1153 .replace("{{upstream_base_url}}", upstream_base_url)
1154 .replace("{{apiKey}}", token)
1155 .replace("{{accessToken}}", token)
1156 .replace("{{token}}", token);
1157
1158 out = out
1159 .replace("{{unix_now}}", &unix_now_secs().to_string())
1160 .replace("{{unix_now_ms}}", &unix_now_ms().to_string());
1161
1162 while let Some(start) = out.find("{{unix_days_ago:") {
1163 let Some(end_offset) = out[start..].find("}}") else {
1164 break;
1165 };
1166 let end = start + end_offset + 2;
1167 let days_str = out[start + "{{unix_days_ago:".len()..end - 2].trim();
1168 let replacement = days_str
1169 .parse::<u64>()
1170 .ok()
1171 .map(|days| unix_now_secs().saturating_sub(days.saturating_mul(24 * 60 * 60)))
1172 .map(|secs| secs.to_string())
1173 .unwrap_or_default();
1174 out.replace_range(start..end, &replacement);
1175 }
1176
1177 while let Some(start) = out.find("{{env:") {
1178 let Some(end_offset) = out[start..].find("}}") else {
1179 break;
1180 };
1181 let end = start + end_offset + 2;
1182 let env_name = out[start + 6..end - 2].trim();
1183 let value = std::env::var(env_name).unwrap_or_default();
1184 out.replace_range(start..end, &value);
1185 }
1186
1187 for (name, value_template) in variables {
1188 let value = render_provider_template(
1189 value_template,
1190 base_url,
1191 upstream_base_url,
1192 token,
1193 &BTreeMap::new(),
1194 );
1195 out = out.replace(&format!("{{{{{name}}}}}"), &value);
1196 }
1197
1198 out
1199}
1200
1201fn resolve_endpoint(
1202 provider: &UsageProviderConfig,
1203 upstream_base_url: &str,
1204 token: &str,
1205) -> Result<String> {
1206 let base_url = normalized_balance_base_url(upstream_base_url)
1207 .ok_or_else(|| anyhow::anyhow!("invalid upstream base_url for balance endpoint"))?;
1208 let endpoint = if provider.endpoint.trim().is_empty() {
1209 provider
1210 .kind
1211 .default_endpoint()
1212 .unwrap_or_default()
1213 .to_string()
1214 } else {
1215 provider.endpoint.trim().to_string()
1216 };
1217 if endpoint.is_empty() {
1218 anyhow::bail!(
1219 "usage provider '{}' has no endpoint and kind {:?} has no default endpoint",
1220 provider.id,
1221 provider.kind
1222 );
1223 }
1224
1225 let rendered = render_provider_template(
1226 &endpoint,
1227 &base_url,
1228 upstream_base_url,
1229 token,
1230 &provider.variables,
1231 );
1232 if rendered.starts_with("http://") || rendered.starts_with("https://") {
1233 return Ok(rendered);
1234 }
1235
1236 let path = if rendered.starts_with('/') {
1237 rendered
1238 } else {
1239 format!("/{rendered}")
1240 };
1241 Ok(format!("{base_url}{path}"))
1242}
1243
1244fn endpoint_origin(endpoint: &str) -> String {
1245 reqwest::Url::parse(endpoint)
1246 .ok()
1247 .and_then(|url| {
1248 let host = url.host_str()?;
1249 let origin = match url.port() {
1250 Some(port) => format!("{}://{}:{}", url.scheme(), host, port),
1251 None => format!("{}://{}", url.scheme(), host),
1252 };
1253 Some(origin)
1254 })
1255 .unwrap_or_else(|| "unknown-origin".to_string())
1256}
1257
1258async fn poll_provider_http_json(
1259 client: &Client,
1260 provider: &UsageProviderConfig,
1261 upstream_base_url: &str,
1262 token: &str,
1263) -> Result<serde_json::Value> {
1264 let endpoint = resolve_endpoint(provider, upstream_base_url, token)?;
1265 let origin = endpoint_origin(&endpoint);
1266 let base_url = normalized_balance_base_url(upstream_base_url).unwrap_or_default();
1267 let mut req = client
1268 .get(endpoint)
1269 .timeout(BALANCE_HTTP_REQUEST_TIMEOUT)
1270 .header("Accept", "application/json")
1271 .header(
1272 "User-Agent",
1273 concat!("codex-helper/", env!("CARGO_PKG_VERSION")),
1274 );
1275
1276 match provider.kind {
1277 ProviderKind::YescodeProfile => {
1278 req = req.header("X-API-Key", token);
1279 }
1280 _ => {
1281 req = req.header("Authorization", format!("Bearer {}", token));
1282 }
1283 }
1284
1285 for (name, template) in &provider.headers {
1286 let value = render_provider_template(
1287 template,
1288 &base_url,
1289 upstream_base_url,
1290 token,
1291 &provider.variables,
1292 );
1293 if !value.trim().is_empty() {
1294 req = req.header(name.as_str(), value);
1295 }
1296 }
1297
1298 let resp = req.send().await.with_context(|| {
1299 format!(
1300 "usage provider request failed for {} via {:?}",
1301 origin, provider.kind
1302 )
1303 })?;
1304
1305 let status = resp.status();
1306 let content_type = resp
1307 .headers()
1308 .get(reqwest::header::CONTENT_TYPE)
1309 .and_then(|value| value.to_str().ok())
1310 .map(str::to_string)
1311 .unwrap_or_else(|| "unknown".to_string());
1312 if !status.is_success() {
1313 let text = resp.text().await.with_context(|| {
1314 format!(
1315 "usage provider error response read failed from {} via {:?}",
1316 origin, provider.kind
1317 )
1318 })?;
1319 let detail = usage_provider_http_error_detail(&text)
1320 .map(|detail| format!(": {detail}"))
1321 .unwrap_or_default();
1322 anyhow::bail!(
1323 "usage provider HTTP {} from {} via {:?}{}",
1324 status,
1325 origin,
1326 provider.kind,
1327 detail
1328 );
1329 }
1330 let text = resp.text().await.with_context(|| {
1331 format!(
1332 "usage provider response read failed from {} via {:?}",
1333 origin, provider.kind
1334 )
1335 })?;
1336 serde_json::from_str(&text).with_context(|| {
1337 format!(
1338 "usage provider returned non-JSON response from {} via {:?} (content-type {}, {} bytes)",
1339 origin,
1340 provider.kind,
1341 content_type,
1342 text.len()
1343 )
1344 })
1345}
1346
1347fn truncate_error_detail(value: &str, max_chars: usize) -> String {
1348 let mut out = String::new();
1349 for (idx, ch) in value.chars().enumerate() {
1350 if idx >= max_chars {
1351 out.push_str("...");
1352 return out;
1353 }
1354 out.push(ch);
1355 }
1356 out
1357}
1358
1359fn compact_error_detail(value: &str) -> Option<String> {
1360 let compact = value.split_whitespace().collect::<Vec<_>>().join(" ");
1361 if compact.is_empty() {
1362 None
1363 } else {
1364 Some(truncate_error_detail(
1365 &compact,
1366 BALANCE_HTTP_ERROR_BODY_LIMIT,
1367 ))
1368 }
1369}
1370
1371fn json_error_detail(value: &serde_json::Value) -> Option<String> {
1372 let code = first_string_from_paths(
1373 value,
1374 &[
1375 "code",
1376 "error.code",
1377 "error.type",
1378 "type",
1379 "data.code",
1380 "data.error.code",
1381 ],
1382 );
1383 let message = first_string_from_paths(
1384 value,
1385 &[
1386 "message",
1387 "msg",
1388 "detail",
1389 "error.message",
1390 "error_description",
1391 "error",
1392 "data.message",
1393 "data.error.message",
1394 ],
1395 );
1396
1397 match (code, message) {
1398 (Some(code), Some(message)) if !message.eq_ignore_ascii_case(&code) => {
1399 Some(format!("{code}: {message}"))
1400 }
1401 (Some(code), _) => Some(code),
1402 (_, Some(message)) => Some(message),
1403 _ => None,
1404 }
1405}
1406
1407fn usage_provider_http_error_detail(text: &str) -> Option<String> {
1408 let trimmed = text.trim();
1409 if trimmed.is_empty() {
1410 return None;
1411 }
1412
1413 if let Ok(value) = serde_json::from_str::<serde_json::Value>(trimmed)
1414 && let Some(detail) = json_error_detail(&value)
1415 {
1416 return compact_error_detail(&detail);
1417 }
1418
1419 compact_error_detail(trimmed)
1420}
1421
1422fn amount_from_json(value: &serde_json::Value) -> Option<UsdAmount> {
1423 let raw = match value {
1424 serde_json::Value::Number(number) => number.to_string(),
1425 serde_json::Value::String(text) => text.trim().to_string(),
1426 _ => return None,
1427 };
1428 UsdAmount::from_decimal_str(raw.as_str())
1429}
1430
1431fn decimal_string_from_json(value: &serde_json::Value) -> Option<String> {
1432 match value {
1433 serde_json::Value::Number(number) => Some(number.to_string()),
1434 serde_json::Value::String(text) => {
1435 let text = text.trim();
1436 if text.is_empty() {
1437 None
1438 } else {
1439 Some(text.to_string())
1440 }
1441 }
1442 _ => None,
1443 }
1444}
1445
1446fn amount_from_json_with_divisor(
1447 value: &serde_json::Value,
1448 divisor: Option<u64>,
1449) -> Option<UsdAmount> {
1450 let amount = amount_from_json(value)?;
1451 match divisor {
1452 Some(divisor) => amount.checked_div_u64(divisor),
1453 None => Some(amount),
1454 }
1455}
1456
1457fn json_value_at_path<'a>(
1458 value: &'a serde_json::Value,
1459 path: &str,
1460) -> Option<&'a serde_json::Value> {
1461 let mut current = value;
1462 for segment in path
1463 .split('.')
1464 .map(str::trim)
1465 .filter(|segment| !segment.is_empty())
1466 {
1467 current = match current {
1468 serde_json::Value::Array(items) => {
1469 let index = segment.parse::<usize>().ok()?;
1470 items.get(index)?
1471 }
1472 _ => current.get(segment)?,
1473 };
1474 }
1475 Some(current)
1476}
1477
1478fn first_amount_from_paths(
1479 value: &serde_json::Value,
1480 custom_paths: &[String],
1481 default_paths: &[&str],
1482 divisor: Option<u64>,
1483) -> Option<UsdAmount> {
1484 custom_paths
1485 .iter()
1486 .map(String::as_str)
1487 .chain(default_paths.iter().copied())
1488 .find_map(|path| {
1489 json_value_at_path(value, path)
1490 .and_then(|value| amount_from_json_with_divisor(value, divisor))
1491 })
1492}
1493
1494fn bool_from_json(value: &serde_json::Value) -> Option<bool> {
1495 match value {
1496 serde_json::Value::Bool(value) => Some(*value),
1497 serde_json::Value::Number(number) => number.as_i64().map(|value| value != 0),
1498 serde_json::Value::String(text) => match text.trim().to_ascii_lowercase().as_str() {
1499 "true" | "yes" | "1" | "exhausted" => Some(true),
1500 "false" | "no" | "0" | "ok" => Some(false),
1501 _ => None,
1502 },
1503 _ => None,
1504 }
1505}
1506
1507fn first_bool_from_paths(
1508 value: &serde_json::Value,
1509 custom_paths: &[String],
1510 default_paths: &[&str],
1511) -> Option<bool> {
1512 custom_paths
1513 .iter()
1514 .map(String::as_str)
1515 .chain(default_paths.iter().copied())
1516 .find_map(|path| json_value_at_path(value, path).and_then(bool_from_json))
1517}
1518
1519fn first_decimal_string_from_paths(
1520 value: &serde_json::Value,
1521 default_paths: &[&str],
1522) -> Option<String> {
1523 default_paths
1524 .iter()
1525 .copied()
1526 .find_map(|path| json_value_at_path(value, path).and_then(decimal_string_from_json))
1527}
1528
1529fn string_from_json(value: &serde_json::Value) -> Option<String> {
1530 match value {
1531 serde_json::Value::String(text) => {
1532 let text = text.trim();
1533 if text.is_empty() {
1534 None
1535 } else {
1536 Some(text.to_string())
1537 }
1538 }
1539 _ => None,
1540 }
1541}
1542
1543fn first_string_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<String> {
1544 default_paths
1545 .iter()
1546 .copied()
1547 .find_map(|path| json_value_at_path(value, path).and_then(string_from_json))
1548}
1549
1550fn u64_from_json(value: &serde_json::Value) -> Option<u64> {
1551 match value {
1552 serde_json::Value::Number(number) => number.as_u64(),
1553 serde_json::Value::String(text) => {
1554 let text = text.trim();
1555 if text.is_empty() {
1556 None
1557 } else {
1558 text.parse::<u64>().ok()
1559 }
1560 }
1561 _ => None,
1562 }
1563}
1564
1565fn seconds_from_json(value: &serde_json::Value) -> Option<u64> {
1566 match value {
1567 serde_json::Value::Number(number) => number.as_f64().map(|value| value.max(0.0) as u64),
1568 serde_json::Value::String(text) => {
1569 let text = text.trim();
1570 if text.is_empty() {
1571 None
1572 } else if let Ok(value) = text.parse::<f64>() {
1573 Some(value.max(0.0) as u64)
1574 } else {
1575 parse_timestamp_secs(text)
1576 }
1577 }
1578 _ => None,
1579 }
1580}
1581
1582fn first_secs_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<u64> {
1583 default_paths
1584 .iter()
1585 .copied()
1586 .find_map(|path| json_value_at_path(value, path).and_then(seconds_from_json))
1587}
1588
1589fn parse_timestamp_secs(value: &str) -> Option<u64> {
1590 parse_rfc3339_like_secs(value).or_else(|| {
1591 httpdate::parse_http_date(value).ok().and_then(|time| {
1592 time.duration_since(std::time::UNIX_EPOCH)
1593 .ok()
1594 .map(|duration| duration.as_secs())
1595 })
1596 })
1597}
1598
1599fn parse_rfc3339_like_secs(value: &str) -> Option<u64> {
1600 let value = value.trim();
1601 let datetime_sep = value.find('T').or_else(|| value.find(' '))?;
1602 let (datetime, offset_secs) = if let Some(datetime) = value.strip_suffix('Z') {
1603 (datetime, 0_i64)
1604 } else {
1605 let offset_pos = value[datetime_sep + 1..]
1606 .rfind(['+', '-'])
1607 .map(|pos| datetime_sep + 1 + pos)?;
1608 let (datetime, offset) = value.split_at(offset_pos);
1609 (datetime, parse_rfc3339_offset_secs(offset)?)
1610 };
1611
1612 let (date, time) = datetime.split_at(datetime_sep);
1613 let time = time.get(1..)?;
1614 let mut date_parts = date.split('-');
1615 let year = date_parts.next()?.parse::<i32>().ok()?;
1616 let month = date_parts.next()?.parse::<u32>().ok()?;
1617 let day = date_parts.next()?.parse::<u32>().ok()?;
1618 if date_parts.next().is_some() {
1619 return None;
1620 }
1621
1622 let mut time_parts = time.split(':');
1623 let hour = time_parts.next()?.parse::<u32>().ok()?;
1624 let minute = time_parts.next()?.parse::<u32>().ok()?;
1625 let second_raw = time_parts.next().unwrap_or("0");
1626 if time_parts.next().is_some() {
1627 return None;
1628 }
1629 let second = second_raw
1630 .split('.')
1631 .next()
1632 .and_then(|value| value.parse::<u32>().ok())?;
1633 if !(1..=12).contains(&month) || day == 0 || hour > 23 || minute > 59 || second > 60 {
1634 return None;
1635 }
1636
1637 let local_secs = days_from_civil(year, month, day)
1638 .checked_mul(86_400)?
1639 .checked_add(i64::from(hour) * 3_600 + i64::from(minute) * 60 + i64::from(second))?;
1640 local_secs
1641 .checked_sub(offset_secs)
1642 .and_then(|utc_secs| u64::try_from(utc_secs).ok())
1643}
1644
1645fn parse_rfc3339_offset_secs(offset: &str) -> Option<i64> {
1646 let sign = match offset.as_bytes().first().copied()? {
1647 b'+' => 1_i64,
1648 b'-' => -1_i64,
1649 _ => return None,
1650 };
1651 let raw = offset.get(1..)?;
1652 let (hours, minutes) = raw
1653 .split_once(':')
1654 .unwrap_or_else(|| raw.split_at(raw.len().min(2)));
1655 let hours = hours.parse::<i64>().ok()?;
1656 let minutes = minutes.parse::<i64>().ok()?;
1657 if hours > 23 || minutes > 59 {
1658 return None;
1659 }
1660 Some(sign * (hours * 3_600 + minutes * 60))
1661}
1662
1663fn days_from_civil(year: i32, month: u32, day: u32) -> i64 {
1664 let year = i64::from(year) - if month <= 2 { 1 } else { 0 };
1665 let era = if year >= 0 { year } else { year - 399 } / 400;
1666 let yoe = year - era * 400;
1667 let month = i64::from(month);
1668 let doy = (153 * (month + if month > 2 { -3 } else { 9 }) + 2) / 5 + i64::from(day) - 1;
1669 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
1670 era * 146_097 + doe - 719_468
1671}
1672
1673fn first_u64_from_paths(value: &serde_json::Value, default_paths: &[&str]) -> Option<u64> {
1674 default_paths
1675 .iter()
1676 .copied()
1677 .find_map(|path| json_value_at_path(value, path).and_then(u64_from_json))
1678}
1679
1680fn array_from_json_path<'a>(
1681 value: &'a serde_json::Value,
1682 path: &str,
1683) -> Option<&'a Vec<serde_json::Value>> {
1684 json_value_at_path(value, path).and_then(|value| value.as_array())
1685}
1686
1687fn amount_to_string(amount: UsdAmount) -> String {
1688 amount.format_usd()
1689}
1690
1691#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1692struct QuotaWindowSnapshot {
1693 period: &'static str,
1694 remaining: UsdAmount,
1695 used: UsdAmount,
1696 limit: UsdAmount,
1697}
1698
1699#[derive(Debug, Clone, PartialEq, Eq)]
1700struct RateLimitWindowSnapshot {
1701 period: String,
1702 reset_at_ms: Option<u64>,
1703}
1704
1705fn base_snapshot(
1706 provider: &UsageProviderConfig,
1707 upstream: &UpstreamRef,
1708 fetched_at_ms: u64,
1709 stale_after_ms: Option<u64>,
1710) -> ProviderBalanceSnapshot {
1711 let mut snapshot = ProviderBalanceSnapshot::new(
1712 provider.id.clone(),
1713 upstream.station_name.clone(),
1714 upstream.index,
1715 provider.kind.source_name(),
1716 fetched_at_ms,
1717 stale_after_ms,
1718 );
1719 if let Some(provider_endpoint) = &upstream.provider_endpoint {
1720 snapshot.provider_endpoint_key = Some(provider_endpoint.stable_key());
1721 }
1722 snapshot.exhaustion_affects_routing = provider.trust_exhaustion_for_routing;
1723 snapshot
1724}
1725
1726fn snapshot_error(
1727 provider: &UsageProviderConfig,
1728 upstream: &UpstreamRef,
1729 fetched_at_ms: u64,
1730 stale_after_ms: Option<u64>,
1731 message: impl Into<String>,
1732) -> ProviderBalanceSnapshot {
1733 base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms).with_error(message)
1734}
1735
1736fn budget_snapshot_from_json(
1737 provider: &UsageProviderConfig,
1738 upstream: &UpstreamRef,
1739 value: &serde_json::Value,
1740 fetched_at_ms: u64,
1741 stale_after_ms: Option<u64>,
1742) -> ProviderBalanceSnapshot {
1743 let monthly_budget = first_amount_from_paths(
1744 value,
1745 &provider.extract.monthly_budget_paths,
1746 &["monthly_budget_usd", "data.monthly_budget_usd"],
1747 provider.extract.monthly_budget_divisor,
1748 );
1749 let monthly_spent = first_amount_from_paths(
1750 value,
1751 &provider.extract.monthly_spent_paths,
1752 &["monthly_spent_usd", "data.monthly_spent_usd"],
1753 provider.extract.monthly_spent_divisor,
1754 );
1755 let exhausted = match (monthly_budget, monthly_spent) {
1756 (Some(budget), Some(spent)) if !budget.is_zero() => Some(spent >= budget),
1757 (Some(_), Some(_)) => Some(false),
1758 _ => None,
1759 };
1760
1761 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1762 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
1763 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
1764 snapshot.exhausted = exhausted;
1765 snapshot.refresh_status(fetched_at_ms);
1766 snapshot
1767}
1768
1769fn yescode_snapshot_from_json(
1770 provider: &UsageProviderConfig,
1771 upstream: &UpstreamRef,
1772 value: &serde_json::Value,
1773 fetched_at_ms: u64,
1774 stale_after_ms: Option<u64>,
1775) -> ProviderBalanceSnapshot {
1776 let subscription_balance = first_amount_from_paths(
1777 value,
1778 &provider.extract.subscription_balance_paths,
1779 &["subscription_balance", "data.subscription_balance"],
1780 provider.extract.remaining_divisor,
1781 );
1782 let paygo_balance = first_amount_from_paths(
1783 value,
1784 &provider.extract.paygo_balance_paths,
1785 &[
1786 "pay_as_you_go_balance",
1787 "paygo_balance",
1788 "data.pay_as_you_go_balance",
1789 "data.paygo_balance",
1790 ],
1791 provider.extract.remaining_divisor,
1792 );
1793 let total_balance = match (subscription_balance, paygo_balance) {
1794 (Some(subscription), Some(paygo)) => Some(subscription.saturating_add(paygo)),
1795 (Some(subscription), None) => Some(subscription),
1796 (None, Some(paygo)) => Some(paygo),
1797 (None, None) => None,
1798 };
1799
1800 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1801 snapshot.total_balance_usd = total_balance.map(amount_to_string);
1802 snapshot.subscription_balance_usd = subscription_balance.map(amount_to_string);
1803 snapshot.paygo_balance_usd = paygo_balance.map(amount_to_string);
1804 snapshot.exhausted = total_balance.map(UsdAmount::is_zero);
1805 snapshot.refresh_status(fetched_at_ms);
1806 snapshot
1807}
1808
1809fn balance_http_snapshot_from_json(
1810 provider: &UsageProviderConfig,
1811 upstream: &UpstreamRef,
1812 value: &serde_json::Value,
1813 fetched_at_ms: u64,
1814 stale_after_ms: Option<u64>,
1815) -> ProviderBalanceSnapshot {
1816 let remaining_balance = first_amount_from_paths(
1817 value,
1818 &provider.extract.remaining_balance_paths,
1819 &[
1820 "balance",
1821 "remaining",
1822 "remain",
1823 "available",
1824 "available_balance",
1825 "credit",
1826 "credits",
1827 "total_balance",
1828 "total_balance_usd",
1829 "totalBalance",
1830 "availableBalance",
1831 "available_balance_usd",
1832 "balance_infos.0.total_balance",
1833 "data.balance",
1834 "data.remaining",
1835 "data.available",
1836 "data.available_balance",
1837 "data.credit",
1838 "data.credits",
1839 "data.total_balance",
1840 "data.totalBalance",
1841 "data.availableBalance",
1842 ],
1843 provider.extract.remaining_divisor,
1844 );
1845 let subscription_balance = first_amount_from_paths(
1846 value,
1847 &provider.extract.subscription_balance_paths,
1848 &[
1849 "subscription_balance",
1850 "subscription_balance_usd",
1851 "subscriptionBalance",
1852 "data.subscription_balance",
1853 "data.subscription_balance_usd",
1854 "data.subscriptionBalance",
1855 ],
1856 provider.extract.remaining_divisor,
1857 );
1858 let paygo_balance = first_amount_from_paths(
1859 value,
1860 &provider.extract.paygo_balance_paths,
1861 &[
1862 "pay_as_you_go_balance",
1863 "paygo_balance",
1864 "paygo",
1865 "paygoBalance",
1866 "chargeBalance",
1867 "voucherBalance",
1868 "data.pay_as_you_go_balance",
1869 "data.paygo_balance",
1870 "data.paygo",
1871 "data.paygoBalance",
1872 "data.chargeBalance",
1873 "data.voucherBalance",
1874 ],
1875 provider.extract.remaining_divisor,
1876 );
1877 let component_remaining = match (subscription_balance, paygo_balance) {
1878 (Some(subscription), Some(paygo)) => Some(subscription.saturating_add(paygo)),
1879 (Some(subscription), None) => Some(subscription),
1880 (None, Some(paygo)) => Some(paygo),
1881 (None, None) => None,
1882 };
1883 let monthly_spent = first_amount_from_paths(
1884 value,
1885 &provider.extract.monthly_spent_paths,
1886 &[
1887 "monthly_spent_usd",
1888 "spent",
1889 "used",
1890 "used_balance",
1891 "usedBalance",
1892 "total_usage",
1893 "data.monthly_spent_usd",
1894 "data.spent",
1895 "data.used",
1896 "data.used_balance",
1897 "data.usedBalance",
1898 "data.total_usage",
1899 ],
1900 provider.extract.monthly_spent_divisor,
1901 );
1902 let monthly_budget = first_amount_from_paths(
1903 value,
1904 &provider.extract.monthly_budget_paths,
1905 &[
1906 "monthly_budget_usd",
1907 "budget",
1908 "limit",
1909 "quota_total",
1910 "creditLimit",
1911 "total_credits",
1912 "data.monthly_budget_usd",
1913 "data.budget",
1914 "data.limit",
1915 "data.quota_total",
1916 "data.creditLimit",
1917 "data.total_credits",
1918 ],
1919 provider.extract.monthly_budget_divisor,
1920 )
1921 .or_else(|| {
1922 if provider.extract.derive_budget_from_remaining_and_spent {
1923 match (remaining_balance.or(component_remaining), monthly_spent) {
1924 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
1925 _ => None,
1926 }
1927 } else {
1928 None
1929 }
1930 });
1931 let total_balance = remaining_balance.or(component_remaining).or_else(|| {
1932 match (
1933 provider.extract.derive_remaining_from_budget_and_spent,
1934 monthly_budget,
1935 monthly_spent,
1936 ) {
1937 (true, Some(budget), Some(spent)) => Some(budget.saturating_sub(spent)),
1938 _ => None,
1939 }
1940 });
1941 let exhausted = first_bool_from_paths(
1942 value,
1943 &provider.extract.exhausted_paths,
1944 &[
1945 "exhausted",
1946 "quota_exhausted",
1947 "balance_exhausted",
1948 "data.exhausted",
1949 "data.quota_exhausted",
1950 "data.balance_exhausted",
1951 ],
1952 )
1953 .or_else(|| total_balance.map(UsdAmount::is_zero));
1954
1955 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
1956 snapshot.total_balance_usd = total_balance.map(amount_to_string);
1957 snapshot.subscription_balance_usd = subscription_balance.map(amount_to_string);
1958 snapshot.paygo_balance_usd = paygo_balance.map(amount_to_string);
1959 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
1960 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
1961 snapshot.exhausted = exhausted;
1962 snapshot.refresh_status(fetched_at_ms);
1963 snapshot
1964}
1965
1966fn has_any_json_path(value: &serde_json::Value, paths: &[&str]) -> bool {
1967 paths
1968 .iter()
1969 .any(|path| json_value_at_path(value, path).is_some())
1970}
1971
1972fn populate_sub2api_usage_fields(
1973 snapshot: &mut ProviderBalanceSnapshot,
1974 value: &serde_json::Value,
1975) {
1976 snapshot.plan_name = first_string_from_paths(
1977 value,
1978 &["planName", "plan_name", "data.planName", "data.plan_name"],
1979 );
1980 let remaining_balance = sub2api_remaining_balance(value);
1981 snapshot.total_balance_usd = snapshot
1982 .total_balance_usd
1983 .take()
1984 .or_else(|| remaining_balance.map(amount_to_string));
1985 snapshot.total_used_usd = first_amount_from_paths(
1986 value,
1987 &[],
1988 &[
1989 "usage.total.total_cost_usd",
1990 "usage.total.total_cost",
1991 "usage.total.cost",
1992 "data.usage.total.total_cost_usd",
1993 "data.usage.total.total_cost",
1994 "data.usage.total.cost",
1995 ],
1996 None,
1997 )
1998 .map(amount_to_string);
1999 snapshot.today_used_usd = first_amount_from_paths(
2000 value,
2001 &[],
2002 &[
2003 "usage.today.total_cost_usd",
2004 "usage.today.total_cost",
2005 "usage.today.cost",
2006 "data.usage.today.total_cost_usd",
2007 "data.usage.today.total_cost",
2008 "data.usage.today.cost",
2009 ],
2010 None,
2011 )
2012 .map(amount_to_string);
2013 snapshot.total_requests = first_u64_from_paths(
2014 value,
2015 &[
2016 "usage.total.request_count",
2017 "usage.total.requests",
2018 "usage.total.count",
2019 "data.usage.total.request_count",
2020 "data.usage.total.requests",
2021 "data.usage.total.count",
2022 ],
2023 );
2024 snapshot.today_requests = first_u64_from_paths(
2025 value,
2026 &[
2027 "usage.today.request_count",
2028 "usage.today.requests",
2029 "usage.today.count",
2030 "data.usage.today.request_count",
2031 "data.usage.today.requests",
2032 "data.usage.today.count",
2033 ],
2034 );
2035 snapshot.total_tokens = first_u64_from_paths(
2036 value,
2037 &[
2038 "usage.total.total_tokens",
2039 "usage.total.tokens",
2040 "usage.total.input_tokens",
2041 "usage.total.prompt_tokens",
2042 "data.usage.total.total_tokens",
2043 "data.usage.total.tokens",
2044 ],
2045 );
2046 snapshot.today_tokens = first_u64_from_paths(
2047 value,
2048 &[
2049 "usage.today.total_tokens",
2050 "usage.today.tokens",
2051 "usage.today.input_tokens",
2052 "usage.today.prompt_tokens",
2053 "data.usage.today.total_tokens",
2054 "data.usage.today.tokens",
2055 ],
2056 );
2057 snapshot.usage_rate = sub2api_usage_rate(value);
2058 snapshot.usage_windows = sub2api_usage_windows(value);
2059 snapshot.usage_model_stats = sub2api_model_stats(value);
2060 snapshot.subscription_expires_at = first_string_from_paths(
2061 value,
2062 &[
2063 "subscription.expires_at",
2064 "data.subscription.expires_at",
2065 "subscription.expiresAt",
2066 "data.subscription.expiresAt",
2067 ],
2068 );
2069 snapshot.usage_alerts = sub2api_usage_alerts(value);
2070}
2071
2072fn sub2api_remaining_balance(value: &serde_json::Value) -> Option<UsdAmount> {
2073 let remaining = first_amount_from_paths(value, &[], &["remaining", "data.remaining"], None)?;
2074 if sub2api_has_subscription_windows(value)
2075 && sub2api_window_remaining_amounts(value).contains(&remaining)
2076 {
2077 return None;
2078 }
2079 Some(remaining)
2080}
2081
2082fn sub2api_has_subscription_windows(value: &serde_json::Value) -> bool {
2083 has_any_json_path(
2084 value,
2085 &[
2086 "subscription.daily_usage_usd",
2087 "subscription.daily_limit_usd",
2088 "subscription.weekly_usage_usd",
2089 "subscription.weekly_limit_usd",
2090 "subscription.monthly_usage_usd",
2091 "subscription.monthly_limit_usd",
2092 "data.subscription.daily_usage_usd",
2093 "data.subscription.daily_limit_usd",
2094 "data.subscription.weekly_usage_usd",
2095 "data.subscription.weekly_limit_usd",
2096 "data.subscription.monthly_usage_usd",
2097 "data.subscription.monthly_limit_usd",
2098 ],
2099 )
2100}
2101
2102fn sub2api_window_remaining_amounts(value: &serde_json::Value) -> Vec<UsdAmount> {
2103 ["daily", "weekly", "monthly"]
2104 .into_iter()
2105 .filter_map(|period| {
2106 let used = first_amount_from_paths(
2107 value,
2108 &[],
2109 &[
2110 &format!("subscription.{period}_usage_usd"),
2111 &format!("data.subscription.{period}_usage_usd"),
2112 ],
2113 None,
2114 );
2115 let limit = first_amount_from_paths(
2116 value,
2117 &[],
2118 &[
2119 &format!("subscription.{period}_limit_usd"),
2120 &format!("data.subscription.{period}_limit_usd"),
2121 ],
2122 None,
2123 );
2124 match (limit, used) {
2125 (Some(limit), Some(used)) if !limit.is_zero() => Some(limit.saturating_sub(used)),
2126 _ => None,
2127 }
2128 })
2129 .collect()
2130}
2131
2132fn optional_amount_is_zero(value: Option<UsdAmount>) -> bool {
2133 value.map(UsdAmount::is_zero).unwrap_or(true)
2134}
2135
2136fn optional_u64_is_zero(value: Option<u64>) -> bool {
2137 value.unwrap_or(0) == 0
2138}
2139
2140fn sub2api_today_usage_is_zero(value: &serde_json::Value) -> bool {
2141 let today_cost = first_amount_from_paths(
2142 value,
2143 &[],
2144 &[
2145 "usage.today.actual_cost",
2146 "usage.today.total_cost_usd",
2147 "usage.today.total_cost",
2148 "usage.today.cost",
2149 "data.usage.today.actual_cost",
2150 "data.usage.today.total_cost_usd",
2151 "data.usage.today.total_cost",
2152 "data.usage.today.cost",
2153 ],
2154 None,
2155 );
2156 let today_requests = first_u64_from_paths(
2157 value,
2158 &[
2159 "usage.today.request_count",
2160 "usage.today.requests",
2161 "usage.today.count",
2162 "data.usage.today.request_count",
2163 "data.usage.today.requests",
2164 "data.usage.today.count",
2165 ],
2166 );
2167 let today_tokens = first_u64_from_paths(
2168 value,
2169 &[
2170 "usage.today.total_tokens",
2171 "usage.today.tokens",
2172 "usage.today.input_tokens",
2173 "usage.today.prompt_tokens",
2174 "data.usage.today.total_tokens",
2175 "data.usage.today.tokens",
2176 ],
2177 );
2178 let has_today_usage_data =
2179 today_cost.is_some() || today_requests.is_some() || today_tokens.is_some();
2180 has_today_usage_data
2181 && optional_amount_is_zero(today_cost)
2182 && optional_u64_is_zero(today_requests)
2183 && optional_u64_is_zero(today_tokens)
2184}
2185
2186fn sub2api_daily_subscription_usage_is_lazy_stale(value: &serde_json::Value) -> bool {
2187 if first_string_from_paths(value, &["mode", "data.mode"]).as_deref() != Some("unrestricted") {
2188 return false;
2189 }
2190
2191 let used = first_amount_from_paths(
2192 value,
2193 &[],
2194 &[
2195 "subscription.daily_usage_usd",
2196 "data.subscription.daily_usage_usd",
2197 ],
2198 None,
2199 );
2200 let limit = first_amount_from_paths(
2201 value,
2202 &[],
2203 &[
2204 "subscription.daily_limit_usd",
2205 "data.subscription.daily_limit_usd",
2206 ],
2207 None,
2208 );
2209
2210 matches!(
2211 (used, limit),
2212 (Some(used), Some(limit))
2213 if !limit.is_zero() && used >= limit && sub2api_today_usage_is_zero(value)
2214 )
2215}
2216
2217fn sub2api_usage_rate(value: &serde_json::Value) -> Option<ProviderUsageRateSnapshot> {
2218 let rate = ProviderUsageRateSnapshot {
2219 average_duration_ms: first_decimal_string_from_paths(
2220 value,
2221 &[
2222 "usage.average_duration_ms",
2223 "data.usage.average_duration_ms",
2224 "average_duration_ms",
2225 "data.average_duration_ms",
2226 ],
2227 ),
2228 rpm: first_decimal_string_from_paths(value, &["usage.rpm", "data.usage.rpm", "rpm"]),
2229 tpm: first_decimal_string_from_paths(value, &["usage.tpm", "data.usage.tpm", "tpm"]),
2230 };
2231 (!rate.is_empty()).then_some(rate)
2232}
2233
2234fn sub2api_usage_windows(value: &serde_json::Value) -> Vec<ProviderUsageWindow> {
2235 ["daily", "weekly", "monthly"]
2236 .into_iter()
2237 .filter_map(|period| {
2238 let used = if period == "daily" && sub2api_daily_subscription_usage_is_lazy_stale(value)
2239 {
2240 Some(UsdAmount::ZERO)
2241 } else {
2242 first_amount_from_paths(
2243 value,
2244 &[],
2245 &[
2246 &format!("subscription.{period}_usage_usd"),
2247 &format!("data.subscription.{period}_usage_usd"),
2248 ],
2249 None,
2250 )
2251 };
2252 let limit = first_amount_from_paths(
2253 value,
2254 &[],
2255 &[
2256 &format!("subscription.{period}_limit_usd"),
2257 &format!("data.subscription.{period}_limit_usd"),
2258 ],
2259 None,
2260 );
2261 if used.is_none() && limit.is_none() {
2262 return None;
2263 }
2264 let unlimited = limit.map(|limit| limit.is_zero());
2265 let remaining = match (limit, used) {
2266 (Some(limit), Some(used)) if !limit.is_zero() => Some(limit.saturating_sub(used)),
2267 _ => None,
2268 };
2269 Some(ProviderUsageWindow {
2270 period: period.to_string(),
2271 used_usd: used.map(amount_to_string),
2272 limit_usd: limit.map(amount_to_string),
2273 remaining_usd: remaining.map(amount_to_string),
2274 unlimited,
2275 })
2276 })
2277 .collect()
2278}
2279
2280fn sub2api_rate_limit_window_from_json(
2281 value: &serde_json::Value,
2282) -> Option<RateLimitWindowSnapshot> {
2283 let period = first_string_from_paths(value, &["window", "period", "name"])?;
2284 let limit = first_u64_from_paths(value, &["limit"]);
2285 if limit == Some(0) {
2286 return None;
2287 }
2288 let remaining = first_u64_from_paths(value, &["remaining"])?;
2289 if remaining > 0 {
2290 return None;
2291 }
2292 let reset_at_ms = first_secs_from_paths(value, &["reset_at", "resets_at", "resetAt"])
2293 .map(|secs| secs.saturating_mul(1000));
2294 Some(RateLimitWindowSnapshot {
2295 period: format!("rate_limit:{period}"),
2296 reset_at_ms,
2297 })
2298}
2299
2300fn sub2api_limiting_rate_limit_window(
2301 value: &serde_json::Value,
2302) -> Option<RateLimitWindowSnapshot> {
2303 ["rate_limits", "data.rate_limits"]
2304 .into_iter()
2305 .find_map(|path| array_from_json_path(value, path))
2306 .and_then(|items| {
2307 items
2308 .iter()
2309 .filter_map(sub2api_rate_limit_window_from_json)
2310 .max_by_key(|window| window.reset_at_ms.unwrap_or(0))
2311 })
2312}
2313
2314fn sub2api_model_stats(value: &serde_json::Value) -> Vec<ProviderUsageModelStat> {
2315 [
2316 "model_stats",
2317 "data.model_stats",
2318 "modelStats",
2319 "data.modelStats",
2320 ]
2321 .into_iter()
2322 .find_map(|path| array_from_json_path(value, path))
2323 .map(|items| {
2324 items
2325 .iter()
2326 .filter_map(sub2api_model_stat_from_json)
2327 .collect::<Vec<_>>()
2328 })
2329 .unwrap_or_default()
2330}
2331
2332fn sub2api_model_stat_from_json(value: &serde_json::Value) -> Option<ProviderUsageModelStat> {
2333 let model = first_string_from_paths(value, &["model", "model_name", "name"])?;
2334 let input_cost = first_amount_from_paths(value, &[], &["input_cost_usd", "input_cost"], None);
2335 let output_cost =
2336 first_amount_from_paths(value, &[], &["output_cost_usd", "output_cost"], None);
2337 let total_cost =
2338 first_amount_from_paths(value, &[], &["total_cost_usd", "total_cost", "cost"], None)
2339 .or_else(|| match (input_cost, output_cost) {
2340 (Some(input), Some(output)) => Some(input.saturating_add(output)),
2341 _ => None,
2342 });
2343 let input_tokens = first_u64_from_paths(value, &["input_tokens", "prompt_tokens"]);
2344 let output_tokens = first_u64_from_paths(value, &["output_tokens", "completion_tokens"]);
2345 let total_tokens =
2346 first_u64_from_paths(value, &["total_tokens", "tokens"]).or_else(|| {
2347 match (input_tokens, output_tokens) {
2348 (Some(input), Some(output)) => input.checked_add(output),
2349 _ => None,
2350 }
2351 });
2352 Some(ProviderUsageModelStat {
2353 model,
2354 request_count: first_u64_from_paths(value, &["request_count", "requests", "count"]),
2355 input_tokens,
2356 output_tokens,
2357 total_tokens,
2358 input_cost_usd: input_cost.map(amount_to_string),
2359 output_cost_usd: output_cost.map(amount_to_string),
2360 total_cost_usd: total_cost.map(amount_to_string),
2361 })
2362}
2363
2364fn sub2api_usage_alerts(value: &serde_json::Value) -> Vec<ProviderUsageAlert> {
2365 let mut alerts = Vec::new();
2366 if let (Some(used), Some(limit)) = (
2367 first_amount_from_paths(
2368 value,
2369 &[],
2370 &[
2371 "subscription.daily_usage_usd",
2372 "data.subscription.daily_usage_usd",
2373 ],
2374 None,
2375 ),
2376 first_amount_from_paths(
2377 value,
2378 &[],
2379 &[
2380 "subscription.daily_limit_usd",
2381 "data.subscription.daily_limit_usd",
2382 ],
2383 None,
2384 ),
2385 ) && !limit.is_zero()
2386 {
2387 let used = if sub2api_daily_subscription_usage_is_lazy_stale(value) {
2388 UsdAmount::ZERO
2389 } else {
2390 used
2391 };
2392 let used_femto = used.femto_usd();
2393 let limit_femto = limit.femto_usd();
2394 if used_femto.saturating_mul(100) >= limit_femto.saturating_mul(95) {
2395 alerts.push(ProviderUsageAlert {
2396 kind: ProviderUsageAlertKind::DailyUsage95,
2397 message: "daily usage is at or above 95%".to_string(),
2398 });
2399 } else if used_femto.saturating_mul(100) >= limit_femto.saturating_mul(80) {
2400 alerts.push(ProviderUsageAlert {
2401 kind: ProviderUsageAlertKind::DailyUsage80,
2402 message: "daily usage is at or above 80%".to_string(),
2403 });
2404 }
2405 }
2406
2407 if let Some(remaining) = sub2api_remaining_balance(value)
2408 && let Some(threshold) = UsdAmount::from_decimal_str(LOW_BALANCE_ALERT_THRESHOLD_USD)
2409 && remaining <= threshold
2410 {
2411 alerts.push(ProviderUsageAlert {
2412 kind: ProviderUsageAlertKind::LowBalance,
2413 message: "remaining balance is low".to_string(),
2414 });
2415 }
2416
2417 if let Some(expires_at_secs) = first_secs_from_paths(
2418 value,
2419 &[
2420 "subscription.expires_at",
2421 "data.subscription.expires_at",
2422 "subscription.expiresAt",
2423 "data.subscription.expiresAt",
2424 ],
2425 ) {
2426 let now = unix_now_secs();
2427 if expires_at_secs <= now {
2428 alerts.push(ProviderUsageAlert {
2429 kind: ProviderUsageAlertKind::SubscriptionExpired,
2430 message: "subscription has expired".to_string(),
2431 });
2432 } else if expires_at_secs <= now.saturating_add(EXPIRING_SOON_WINDOW_SECS) {
2433 alerts.push(ProviderUsageAlert {
2434 kind: ProviderUsageAlertKind::SubscriptionExpiringSoon,
2435 message: "subscription expires within 7 days".to_string(),
2436 });
2437 }
2438 }
2439
2440 alerts.sort_by_key(|alert| alert.kind);
2441 alerts.dedup_by_key(|alert| alert.kind);
2442 alerts
2443}
2444
2445fn sub2api_subscription_limit_snapshot(
2446 value: &serde_json::Value,
2447 period: &'static str,
2448 limit_paths: &[&str],
2449 usage_paths: &[&str],
2450) -> Option<QuotaWindowSnapshot> {
2451 let budget = first_amount_from_paths(value, &[], limit_paths, None)?;
2452 if budget.is_zero() {
2453 return None;
2454 }
2455 let spent = if period == "daily" && sub2api_daily_subscription_usage_is_lazy_stale(value) {
2456 UsdAmount::ZERO
2457 } else {
2458 first_amount_from_paths(value, &[], usage_paths, None).unwrap_or(UsdAmount::ZERO)
2459 };
2460 let remaining = budget.saturating_sub(spent);
2461 Some(QuotaWindowSnapshot {
2462 period,
2463 remaining,
2464 used: spent,
2465 limit: budget,
2466 })
2467}
2468
2469fn sub2api_limiting_subscription_window(value: &serde_json::Value) -> Option<QuotaWindowSnapshot> {
2470 let windows = [
2471 sub2api_subscription_limit_snapshot(
2472 value,
2473 "daily",
2474 &[
2475 "subscription.daily_limit_usd",
2476 "data.subscription.daily_limit_usd",
2477 ],
2478 &[
2479 "subscription.daily_usage_usd",
2480 "data.subscription.daily_usage_usd",
2481 ],
2482 ),
2483 sub2api_subscription_limit_snapshot(
2484 value,
2485 "weekly",
2486 &[
2487 "subscription.weekly_limit_usd",
2488 "data.subscription.weekly_limit_usd",
2489 ],
2490 &[
2491 "subscription.weekly_usage_usd",
2492 "data.subscription.weekly_usage_usd",
2493 ],
2494 ),
2495 sub2api_subscription_limit_snapshot(
2496 value,
2497 "monthly",
2498 &[
2499 "subscription.monthly_limit_usd",
2500 "data.subscription.monthly_limit_usd",
2501 ],
2502 &[
2503 "subscription.monthly_usage_usd",
2504 "data.subscription.monthly_usage_usd",
2505 ],
2506 ),
2507 ];
2508
2509 windows
2510 .into_iter()
2511 .flatten()
2512 .min_by_key(|window| window.remaining)
2513}
2514
2515fn sub2api_usage_snapshot_from_json(
2516 provider: &UsageProviderConfig,
2517 upstream: &UpstreamRef,
2518 value: &serde_json::Value,
2519 fetched_at_ms: u64,
2520 stale_after_ms: Option<u64>,
2521) -> ProviderBalanceSnapshot {
2522 if json_value_at_path(value, "isValid").and_then(bool_from_json) == Some(false) {
2523 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2524 .with_error("sub2api usage response reported invalid API key");
2525 }
2526
2527 let mode = first_string_from_paths(value, &["mode", "data.mode"]);
2528 let has_subscription = has_any_json_path(value, &["subscription", "data.subscription"]);
2529
2530 if mode.as_deref() == Some("quota_limited") {
2531 let rate_limit_window = sub2api_limiting_rate_limit_window(value);
2532 let quota_remaining = first_amount_from_paths(
2533 value,
2534 &provider.extract.remaining_balance_paths,
2535 &[
2536 "quota.remaining",
2537 "data.quota.remaining",
2538 "remaining",
2539 "data.remaining",
2540 ],
2541 provider.extract.remaining_divisor,
2542 );
2543 let quota_limit = first_amount_from_paths(
2544 value,
2545 &provider.extract.monthly_budget_paths,
2546 &["quota.limit", "data.quota.limit"],
2547 provider.extract.monthly_budget_divisor,
2548 );
2549 let quota_used = first_amount_from_paths(
2550 value,
2551 &provider.extract.monthly_spent_paths,
2552 &["quota.used", "data.quota.used"],
2553 provider.extract.monthly_spent_divisor,
2554 );
2555 let quota_exhausted = first_bool_from_paths(
2556 value,
2557 &provider.extract.exhausted_paths,
2558 &[
2559 "exhausted",
2560 "data.exhausted",
2561 "quota_exhausted",
2562 "data.quota_exhausted",
2563 ],
2564 )
2565 .or_else(|| quota_remaining.map(UsdAmount::is_zero));
2566 let exhausted = Some(quota_exhausted.unwrap_or(false) || rate_limit_window.is_some());
2567
2568 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2569 if let Some(rate_limit_window) = rate_limit_window.clone()
2570 && quota_exhausted != Some(true)
2571 {
2572 snapshot.quota_period = Some(rate_limit_window.period);
2573 snapshot.quota_resets_at_ms = rate_limit_window.reset_at_ms;
2574 } else {
2575 snapshot.quota_period = Some("quota".to_string());
2576 snapshot.quota_remaining_usd = quota_remaining.map(amount_to_string);
2577 snapshot.quota_limit_usd = quota_limit.map(amount_to_string);
2578 snapshot.quota_used_usd = quota_used.map(amount_to_string);
2579 snapshot.monthly_budget_usd = quota_limit.map(amount_to_string);
2580 snapshot.monthly_spent_usd = quota_used.map(amount_to_string);
2581 }
2582 snapshot.exhausted = exhausted;
2583 populate_sub2api_usage_fields(&mut snapshot, value);
2584 snapshot.refresh_status(fetched_at_ms);
2585 return snapshot;
2586 }
2587
2588 if mode.as_deref() == Some("unrestricted") && has_subscription {
2589 let limiting_window = sub2api_limiting_subscription_window(value);
2590 let exhausted = first_bool_from_paths(
2591 value,
2592 &provider.extract.exhausted_paths,
2593 &[
2594 "exhausted",
2595 "data.exhausted",
2596 "quota_exhausted",
2597 "data.quota_exhausted",
2598 ],
2599 )
2600 .or_else(|| limiting_window.map(|window| window.remaining.is_zero()))
2601 .or(Some(false));
2602
2603 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2604 if let Some(window) = limiting_window {
2605 snapshot.quota_period = Some(window.period.to_string());
2606 snapshot.quota_remaining_usd = Some(amount_to_string(window.remaining));
2607 snapshot.quota_limit_usd = Some(amount_to_string(window.limit));
2608 snapshot.quota_used_usd = Some(amount_to_string(window.used));
2609 snapshot.monthly_budget_usd = Some(amount_to_string(window.limit));
2610 snapshot.monthly_spent_usd = Some(amount_to_string(window.used));
2611 }
2612 snapshot.exhaustion_affects_routing = false;
2613 snapshot.exhausted = exhausted;
2614 populate_sub2api_usage_fields(&mut snapshot, value);
2615 snapshot.refresh_status(fetched_at_ms);
2616 return snapshot;
2617 }
2618
2619 let mut snapshot =
2620 balance_http_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms);
2621 populate_sub2api_usage_fields(&mut snapshot, value);
2622 snapshot.refresh_status(fetched_at_ms);
2623 snapshot
2624}
2625
2626fn sub2api_auth_me_snapshot_from_json(
2627 provider: &UsageProviderConfig,
2628 upstream: &UpstreamRef,
2629 value: &serde_json::Value,
2630 fetched_at_ms: u64,
2631 stale_after_ms: Option<u64>,
2632) -> ProviderBalanceSnapshot {
2633 if json_value_at_path(value, "code")
2634 .and_then(|value| value.as_i64())
2635 .is_some_and(|code| code != 0)
2636 {
2637 let message = json_value_at_path(value, "message")
2638 .and_then(|value| value.as_str())
2639 .unwrap_or("sub2api auth/me response reported failure");
2640 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2641 .with_error(message.to_string());
2642 }
2643
2644 balance_http_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
2645}
2646
2647fn rightcode_available_prefixes(value: &serde_json::Value) -> Vec<String> {
2648 array_from_json_path(value, "available_prefixes")
2649 .into_iter()
2650 .flatten()
2651 .filter_map(|value| value.as_str())
2652 .map(str::trim)
2653 .filter(|value| !value.is_empty())
2654 .map(ToOwned::to_owned)
2655 .collect()
2656}
2657
2658fn rightcode_subscription_window(value: &serde_json::Value) -> Option<QuotaWindowSnapshot> {
2659 let limit = json_value_at_path(value, "total_quota").and_then(amount_from_json)?;
2660 if limit < UsdAmount::from_decimal_str("10").unwrap_or(UsdAmount::ZERO) {
2661 return None;
2662 }
2663 let raw_remaining = json_value_at_path(value, "remaining_quota").and_then(amount_from_json)?;
2664 let reset_today = json_value_at_path(value, "reset_today").and_then(bool_from_json);
2665 let remaining = if reset_today == Some(true) {
2666 raw_remaining
2667 } else {
2668 raw_remaining.saturating_add(limit)
2669 };
2670 let used = limit.saturating_sub(remaining);
2671 Some(QuotaWindowSnapshot {
2672 period: "daily",
2673 remaining,
2674 used,
2675 limit,
2676 })
2677}
2678
2679fn rightcode_account_summary_snapshot_from_json(
2680 provider: &UsageProviderConfig,
2681 upstream: &UpstreamRef,
2682 value: &serde_json::Value,
2683 upstream_base_url: &str,
2684 fetched_at_ms: u64,
2685 stale_after_ms: Option<u64>,
2686) -> ProviderBalanceSnapshot {
2687 let balance = json_value_at_path(value, "balance").and_then(amount_from_json);
2688 let provider_prefixes = base_path_prefixes(upstream_base_url);
2689 let mut matched_windows = Vec::new();
2690 let mut matched_plan_names = Vec::new();
2691
2692 if let Some(subscriptions) = array_from_json_path(value, "subscriptions") {
2693 for subscription in subscriptions {
2694 let available_prefixes = rightcode_available_prefixes(subscription);
2695 if !path_prefixes_match(&provider_prefixes, &available_prefixes) {
2696 continue;
2697 }
2698 let Some(window) = rightcode_subscription_window(subscription) else {
2699 continue;
2700 };
2701 matched_windows.push(window);
2702 if let Some(name) = json_value_at_path(subscription, "name").and_then(string_from_json)
2703 {
2704 matched_plan_names.push(name);
2705 }
2706 }
2707 }
2708
2709 if balance.is_none() && matched_windows.is_empty() {
2710 return snapshot_error(
2711 provider,
2712 upstream,
2713 fetched_at_ms,
2714 stale_after_ms,
2715 "rightcode account summary missing balance and matching subscription quota fields",
2716 );
2717 }
2718
2719 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2720 snapshot.total_balance_usd = balance.map(amount_to_string);
2721
2722 if !matched_windows.is_empty() {
2723 let mut remaining = UsdAmount::ZERO;
2724 let mut used = UsdAmount::ZERO;
2725 let mut limit = UsdAmount::ZERO;
2726 for window in matched_windows {
2727 remaining = remaining.saturating_add(window.remaining);
2728 used = used.saturating_add(window.used);
2729 limit = limit.saturating_add(window.limit);
2730 }
2731 snapshot.quota_period = Some("daily".to_string());
2732 snapshot.quota_remaining_usd = Some(amount_to_string(remaining));
2733 snapshot.quota_used_usd = Some(amount_to_string(used));
2734 snapshot.quota_limit_usd = Some(amount_to_string(limit));
2735 if !matched_plan_names.is_empty() {
2736 matched_plan_names.sort();
2737 matched_plan_names.dedup();
2738 snapshot.plan_name = Some(matched_plan_names.join(", "));
2739 }
2740 snapshot.exhausted = Some(remaining.is_zero() && balance.is_none_or(UsdAmount::is_zero));
2741 } else {
2742 snapshot.exhausted = balance.map(UsdAmount::is_zero);
2743 }
2744
2745 snapshot.refresh_status(fetched_at_ms);
2746 snapshot
2747}
2748
2749fn new_api_token_usage_snapshot_from_json(
2750 provider: &UsageProviderConfig,
2751 upstream: &UpstreamRef,
2752 value: &serde_json::Value,
2753 fetched_at_ms: u64,
2754 stale_after_ms: Option<u64>,
2755) -> ProviderBalanceSnapshot {
2756 if json_value_at_path(value, "success").and_then(bool_from_json) == Some(false)
2757 || json_value_at_path(value, "code").and_then(bool_from_json) == Some(false)
2758 {
2759 let message = json_value_at_path(value, "message")
2760 .and_then(|value| value.as_str())
2761 .unwrap_or("new api token usage response reported failure");
2762 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2763 .with_error(message.to_string());
2764 }
2765
2766 let mut effective = provider.extract.clone();
2767 if effective.remaining_balance_paths.is_empty() {
2768 effective.remaining_balance_paths = vec![
2769 "data.total_available".to_string(),
2770 "data.remain_quota".to_string(),
2771 "total_available".to_string(),
2772 "remain_quota".to_string(),
2773 ];
2774 }
2775 if effective.monthly_spent_paths.is_empty() {
2776 effective.monthly_spent_paths = vec![
2777 "data.total_used".to_string(),
2778 "data.used_quota".to_string(),
2779 "total_used".to_string(),
2780 "used_quota".to_string(),
2781 ];
2782 }
2783 if effective.monthly_budget_paths.is_empty() {
2784 effective.monthly_budget_paths = vec![
2785 "data.total_granted".to_string(),
2786 "total_granted".to_string(),
2787 ];
2788 }
2789 effective.remaining_divisor = effective.remaining_divisor.or(Some(500_000));
2790 effective.monthly_spent_divisor = effective.monthly_spent_divisor.or(Some(500_000));
2791 effective.monthly_budget_divisor = effective.monthly_budget_divisor.or(Some(500_000));
2792
2793 let unlimited_quota =
2794 first_bool_from_paths(value, &[], &["data.unlimited_quota", "unlimited_quota"])
2795 == Some(true);
2796 let remaining_balance = first_amount_from_paths(
2797 value,
2798 &effective.remaining_balance_paths,
2799 &[],
2800 effective.remaining_divisor,
2801 );
2802 let monthly_spent = first_amount_from_paths(
2803 value,
2804 &effective.monthly_spent_paths,
2805 &[],
2806 effective.monthly_spent_divisor,
2807 );
2808 let monthly_budget = first_amount_from_paths(
2809 value,
2810 &effective.monthly_budget_paths,
2811 &[],
2812 effective.monthly_budget_divisor,
2813 )
2814 .or_else(|| match (remaining_balance, monthly_spent) {
2815 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
2816 _ => None,
2817 });
2818 let exhausted = if unlimited_quota {
2819 Some(false)
2820 } else {
2821 first_bool_from_paths(
2822 value,
2823 &effective.exhausted_paths,
2824 &["data.exhausted", "exhausted"],
2825 )
2826 .or_else(|| remaining_balance.map(UsdAmount::is_zero))
2827 };
2828
2829 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2830 snapshot.plan_name = first_string_from_paths(value, &["data.name", "name"]);
2831 snapshot.unlimited_quota = Some(unlimited_quota);
2832 if !unlimited_quota {
2833 snapshot.quota_period = Some("token".to_string());
2834 snapshot.quota_remaining_usd = remaining_balance.map(amount_to_string);
2835 snapshot.quota_limit_usd = monthly_budget.map(amount_to_string);
2836 snapshot.quota_used_usd = monthly_spent.map(amount_to_string);
2837 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
2838 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2839 } else {
2840 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2841 }
2842 snapshot.exhausted = exhausted;
2843 snapshot.refresh_status(fetched_at_ms);
2844 snapshot
2845}
2846
2847fn new_api_snapshot_from_json(
2848 provider: &UsageProviderConfig,
2849 upstream: &UpstreamRef,
2850 value: &serde_json::Value,
2851 fetched_at_ms: u64,
2852 stale_after_ms: Option<u64>,
2853) -> ProviderBalanceSnapshot {
2854 if json_value_at_path(value, "success").and_then(bool_from_json) == Some(false) {
2855 let message = json_value_at_path(value, "message")
2856 .and_then(|value| value.as_str())
2857 .unwrap_or("new api balance response reported failure");
2858 return base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms)
2859 .with_error(message.to_string());
2860 }
2861
2862 let mut effective = provider.extract.clone();
2863 if effective.remaining_balance_paths.is_empty() {
2864 effective.remaining_balance_paths = vec!["data.quota".to_string(), "quota".to_string()];
2865 }
2866 if effective.monthly_spent_paths.is_empty() {
2867 effective.monthly_spent_paths =
2868 vec!["data.used_quota".to_string(), "used_quota".to_string()];
2869 }
2870 effective.remaining_divisor = effective.remaining_divisor.or(Some(500_000));
2871 effective.monthly_spent_divisor = effective.monthly_spent_divisor.or(Some(500_000));
2872 effective.monthly_budget_divisor = effective.monthly_budget_divisor.or(Some(500_000));
2873
2874 let remaining_balance = first_amount_from_paths(
2875 value,
2876 &effective.remaining_balance_paths,
2877 &[],
2878 effective.remaining_divisor,
2879 );
2880 let monthly_spent = first_amount_from_paths(
2881 value,
2882 &effective.monthly_spent_paths,
2883 &[],
2884 effective.monthly_spent_divisor,
2885 );
2886 let monthly_budget = first_amount_from_paths(
2887 value,
2888 &effective.monthly_budget_paths,
2889 &["data.total_quota", "total_quota"],
2890 effective.monthly_budget_divisor,
2891 )
2892 .or_else(|| match (remaining_balance, monthly_spent) {
2893 (Some(remaining), Some(spent)) => Some(remaining.saturating_add(spent)),
2894 _ => None,
2895 });
2896 let unlimited_quota =
2897 first_bool_from_paths(value, &[], &["data.unlimited_quota", "unlimited_quota"])
2898 == Some(true);
2899 let exhausted = if unlimited_quota {
2900 Some(false)
2901 } else {
2902 first_bool_from_paths(
2903 value,
2904 &effective.exhausted_paths,
2905 &["data.exhausted", "exhausted"],
2906 )
2907 .or_else(|| remaining_balance.map(UsdAmount::is_zero))
2908 };
2909
2910 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2911 snapshot.unlimited_quota = Some(unlimited_quota);
2912 if !unlimited_quota {
2913 snapshot.quota_period = Some("quota".to_string());
2914 snapshot.quota_remaining_usd = remaining_balance.map(amount_to_string);
2915 snapshot.quota_limit_usd = monthly_budget.map(amount_to_string);
2916 snapshot.quota_used_usd = monthly_spent.map(amount_to_string);
2917 snapshot.monthly_budget_usd = monthly_budget.map(amount_to_string);
2918 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2919 } else {
2920 snapshot.monthly_spent_usd = monthly_spent.map(amount_to_string);
2921 }
2922 snapshot.exhausted = exhausted;
2923 snapshot.refresh_status(fetched_at_ms);
2924 snapshot
2925}
2926
2927fn openai_cost_result_usd_amount(result: &serde_json::Value) -> Option<UsdAmount> {
2928 let amount = json_value_at_path(result, "amount.value").and_then(amount_from_json)?;
2929 let currency = json_value_at_path(result, "amount.currency").and_then(|value| value.as_str());
2930 match currency {
2931 Some(currency) if currency.eq_ignore_ascii_case("usd") => Some(amount),
2932 None => Some(amount),
2933 _ => None,
2934 }
2935}
2936
2937fn openai_organization_costs_total(value: &serde_json::Value) -> Option<UsdAmount> {
2938 let buckets = json_value_at_path(value, "data")?.as_array()?;
2939 let mut total = UsdAmount::ZERO;
2940
2941 for bucket in buckets {
2942 let Some(results) =
2943 json_value_at_path(bucket, "results").and_then(|value| value.as_array())
2944 else {
2945 continue;
2946 };
2947 for result in results {
2948 if let Some(amount) = openai_cost_result_usd_amount(result) {
2949 total = total.saturating_add(amount);
2950 }
2951 }
2952 }
2953
2954 Some(total)
2955}
2956
2957fn openai_organization_costs_snapshot_from_json(
2958 provider: &UsageProviderConfig,
2959 upstream: &UpstreamRef,
2960 value: &serde_json::Value,
2961 fetched_at_ms: u64,
2962 stale_after_ms: Option<u64>,
2963) -> ProviderBalanceSnapshot {
2964 let spent = first_amount_from_paths(
2965 value,
2966 &provider.extract.monthly_spent_paths,
2967 &[],
2968 provider.extract.monthly_spent_divisor,
2969 )
2970 .or_else(|| openai_organization_costs_total(value));
2971
2972 let mut snapshot = base_snapshot(provider, upstream, fetched_at_ms, stale_after_ms);
2973 snapshot.monthly_spent_usd = spent.map(amount_to_string);
2974 snapshot.exhausted = None;
2975 snapshot.exhaustion_affects_routing = false;
2976 snapshot.refresh_status(fetched_at_ms);
2977 snapshot
2978}
2979
2980fn balance_exhaustion_policy_action(
2981 endpoint_key: ProviderEndpointKey,
2982 observed_at_ms: u64,
2983) -> Option<PolicyAction> {
2984 const BALANCE_EXHAUSTION_COOLDOWN_SECS: u64 = 24 * 60 * 60;
2985 let mut signal = ProviderSignal::high_confidence_route_facing(
2986 ProviderSignalKind::Balance,
2987 ProviderSignalSource::BalanceSnapshot,
2988 ProviderSignalTarget::ProviderEndpoint {
2989 provider_endpoint_key: endpoint_key,
2990 },
2991 observed_at_ms,
2992 );
2993 signal.reset_after_secs = Some(BALANCE_EXHAUSTION_COOLDOWN_SECS);
2994 signal.reason = Some("balance_exhausted".to_string());
2995 PolicyAction::cooldown_from_signal(signal, observed_at_ms, 0, observed_at_ms)
2996}
2997
2998async fn sync_balance_policy_action_for_endpoint(
2999 state: &Arc<ProxyState>,
3000 service_name: &str,
3001 endpoint_key: ProviderEndpointKey,
3002 exhausted: bool,
3003) {
3004 if exhausted {
3005 if let Some(action) =
3006 balance_exhaustion_policy_action(endpoint_key, crate::logging::now_ms())
3007 {
3008 state.upsert_owned_policy_action(service_name, action).await;
3009 }
3010 } else {
3011 state
3012 .clear_owned_policy_action(
3013 service_name,
3014 &endpoint_key,
3015 PolicyActionKind::Cooldown,
3016 ProviderSignalKind::Balance,
3017 ProviderSignalSource::BalanceSnapshot,
3018 )
3019 .await;
3020 }
3021}
3022
3023async fn update_usage_exhausted(
3024 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
3025 state: &Arc<ProxyState>,
3026 cfg: &ProxyConfig,
3027 service_name: &str,
3028 upstreams: &[UpstreamRef],
3029 exhausted: bool,
3030) {
3031 if let Ok(mut map) = lb_states.lock() {
3032 for uref in upstreams {
3033 let service = match service_manager(cfg, service_name).station(&uref.station_name) {
3034 Some(s) => s,
3035 None => continue,
3036 };
3037
3038 let entry = map
3039 .entry(uref.station_name.clone())
3040 .or_insert_with(LbState::default);
3041 entry.ensure_layout(service.name.as_str(), &service.upstreams);
3042 if uref.index < entry.usage_exhausted.len() {
3043 entry.usage_exhausted[uref.index] = exhausted;
3044 }
3045 }
3046 }
3047
3048 for uref in upstreams {
3049 if let Some(endpoint_key) = uref.provider_endpoint.clone() {
3050 sync_balance_policy_action_for_endpoint(
3051 state,
3052 service_name,
3053 endpoint_key.clone(),
3054 exhausted,
3055 )
3056 .await;
3057 state
3058 .set_provider_endpoint_usage_exhausted(service_name, endpoint_key, exhausted)
3059 .await;
3060 }
3061 }
3062}
3063
3064fn provider_hosts_for_diagnostics(
3065 cfg: &ProxyConfig,
3066 service_name: &str,
3067 provider: &UsageProviderConfig,
3068) -> Vec<String> {
3069 let mut hosts: Vec<String> = Vec::new();
3070 for service in service_manager(cfg, service_name).stations().values() {
3071 for upstream in &service.upstreams {
3072 if domain_matches(&upstream.base_url, &provider.domains)
3073 && let Ok(url) = reqwest::Url::parse(&upstream.base_url)
3074 && let Some(host) = url.host_str()
3075 {
3076 hosts.push(host.to_string());
3077 }
3078 }
3079 }
3080 hosts.sort();
3081 hosts.dedup();
3082 hosts
3083}
3084
3085fn warn_if_provider_spans_hosts(
3086 cfg: &ProxyConfig,
3087 service_name: &str,
3088 provider: &UsageProviderConfig,
3089) {
3090 let hosts = provider_hosts_for_diagnostics(cfg, service_name, provider);
3091 if hosts.len() > 1 {
3092 warn!(
3093 "usage provider '{}' is associated with multiple hosts: {:?}; \
3094将按统一额度处理这些 upstream,如需区分配额请拆分为多个 provider 配置",
3095 provider.id, hosts
3096 );
3097 }
3098}
3099
3100fn snapshot_from_provider_json(
3101 provider: &UsageProviderConfig,
3102 upstream: &UpstreamRef,
3103 value: &serde_json::Value,
3104 upstream_base_url: &str,
3105 fetched_at_ms: u64,
3106 stale_after_ms: Option<u64>,
3107) -> ProviderBalanceSnapshot {
3108 match provider.kind {
3109 ProviderKind::BudgetHttpJson => {
3110 budget_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
3111 }
3112 ProviderKind::YescodeProfile => {
3113 yescode_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
3114 }
3115 ProviderKind::OpenAiBalanceHttpJson => balance_http_snapshot_from_json(
3116 provider,
3117 upstream,
3118 value,
3119 fetched_at_ms,
3120 stale_after_ms,
3121 ),
3122 ProviderKind::Sub2ApiUsage => sub2api_usage_snapshot_from_json(
3123 provider,
3124 upstream,
3125 value,
3126 fetched_at_ms,
3127 stale_after_ms,
3128 ),
3129 ProviderKind::Sub2ApiAuthMe => sub2api_auth_me_snapshot_from_json(
3130 provider,
3131 upstream,
3132 value,
3133 fetched_at_ms,
3134 stale_after_ms,
3135 ),
3136 ProviderKind::NewApiTokenUsage => new_api_token_usage_snapshot_from_json(
3137 provider,
3138 upstream,
3139 value,
3140 fetched_at_ms,
3141 stale_after_ms,
3142 ),
3143 ProviderKind::NewApiUserSelf => {
3144 new_api_snapshot_from_json(provider, upstream, value, fetched_at_ms, stale_after_ms)
3145 }
3146 ProviderKind::RightCodeAccountSummary => rightcode_account_summary_snapshot_from_json(
3147 provider,
3148 upstream,
3149 value,
3150 upstream_base_url,
3151 fetched_at_ms,
3152 stale_after_ms,
3153 ),
3154 ProviderKind::OpenAiOrganizationCosts => openai_organization_costs_snapshot_from_json(
3155 provider,
3156 upstream,
3157 value,
3158 fetched_at_ms,
3159 stale_after_ms,
3160 ),
3161 }
3162}
3163
3164async fn refresh_provider_target(
3165 params: RefreshProviderTargetParams<'_>,
3166) -> UsageProviderRefreshOutcome {
3167 let RefreshProviderTargetParams {
3168 client,
3169 provider,
3170 target,
3171 cfg,
3172 lb_states,
3173 state,
3174 service_name,
3175 interval_secs,
3176 } = params;
3177
3178 let upstreams = vec![target.upstream.clone()];
3179 let fetched_at_ms = unix_now_ms();
3180 let stale_after_ms = stale_after_ms(fetched_at_ms, interval_secs);
3181 if let Some(suppression) =
3182 usage_provider_target_suppression_active(&provider.id, target, Instant::now())
3183 {
3184 update_usage_exhausted(
3185 lb_states,
3186 state,
3187 cfg,
3188 service_name,
3189 &upstreams,
3190 suppression.routing_exhausted,
3191 )
3192 .await;
3193 warn!(
3194 "usage provider '{}' skipped {}[{}]: balance refresh suppressed: {}",
3195 provider.id, target.upstream.station_name, target.upstream.index, suppression.reason
3196 );
3197 return UsageProviderRefreshOutcome::Failed;
3198 }
3199 if let Some(decision) = existing_usage_provider_target_suppression_decision(
3200 state,
3201 cfg,
3202 service_name,
3203 &provider.id,
3204 target,
3205 fetched_at_ms,
3206 )
3207 .await
3208 {
3209 remember_usage_provider_target_suppression(
3210 &provider.id,
3211 target,
3212 decision.ttl,
3213 decision.reason.clone(),
3214 decision.routing_exhausted,
3215 Instant::now(),
3216 );
3217 update_usage_exhausted(
3218 lb_states,
3219 state,
3220 cfg,
3221 service_name,
3222 &upstreams,
3223 decision.routing_exhausted,
3224 )
3225 .await;
3226 warn!(
3227 "usage provider '{}' skipped {}[{}]: existing balance snapshot suppresses refresh: {}",
3228 provider.id, target.upstream.station_name, target.upstream.index, decision.reason
3229 );
3230 return UsageProviderRefreshOutcome::Failed;
3231 }
3232
3233 let Some(token) = resolve_token(provider, &upstreams, cfg, service_name) else {
3234 let snapshot = if provider.kind == ProviderKind::OpenAiOrganizationCosts {
3235 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
3236 } else {
3237 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
3238 .with_error("no usable token; checked provider token_env and upstream auth")
3239 };
3240 state
3241 .record_provider_balance_snapshot(service_name, snapshot)
3242 .await;
3243 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3244 if provider.kind == ProviderKind::OpenAiOrganizationCosts {
3245 warn!(
3246 "usage provider '{}' is missing OPENAI_ADMIN_KEY; OpenAI official costs stay unknown",
3247 provider.id
3248 );
3249 } else {
3250 warn!(
3251 "usage provider '{}' has no usable token (checked token_env and associated upstream auth_token); \
3252跳过本次用量查询,请检查 usage_providers.json 和 ~/.codex-helper/config.json",
3253 provider.id
3254 );
3255 }
3256 return UsageProviderRefreshOutcome::MissingToken;
3257 };
3258
3259 match poll_provider_http_json(client, provider, &target.base_url, &token).await {
3260 Ok(value) => {
3261 let snapshot = snapshot_from_provider_json(
3262 provider,
3263 &upstreams[0],
3264 &value,
3265 &target.base_url,
3266 fetched_at_ms,
3267 stale_after_ms,
3268 );
3269 let snapshot_error = usage_provider_snapshot_error(&snapshot).map(str::to_string);
3270 let suppression_decision =
3271 usage_provider_suppression_decision_from_snapshot(&snapshot, cfg, fetched_at_ms);
3272 let exhausted_for_lb = suppression_decision
3273 .as_ref()
3274 .is_some_and(|decision| decision.routing_exhausted);
3275 if let Some(decision) = suppression_decision.as_ref() {
3276 remember_usage_provider_target_suppression(
3277 &provider.id,
3278 target,
3279 decision.ttl,
3280 decision.reason.as_str(),
3281 decision.routing_exhausted,
3282 Instant::now(),
3283 );
3284 } else {
3285 clear_usage_provider_target_suppression(&provider.id, target);
3286 }
3287 update_usage_exhausted(
3288 lb_states,
3289 state,
3290 cfg,
3291 service_name,
3292 &upstreams,
3293 exhausted_for_lb,
3294 )
3295 .await;
3296 state
3297 .record_provider_balance_snapshot(service_name, snapshot)
3298 .await;
3299 if let Some(error) = snapshot_error {
3300 warn!(
3301 "usage provider '{}' returned error snapshot for {}[{}]: {}",
3302 provider.id, target.upstream.station_name, target.upstream.index, error
3303 );
3304 return UsageProviderRefreshOutcome::Failed;
3305 }
3306 info!(
3307 "usage provider '{}' refreshed {}[{}], exhausted = {}, routing_trusted = {}",
3308 provider.id,
3309 target.upstream.station_name,
3310 target.upstream.index,
3311 exhausted_for_lb,
3312 provider.trust_exhaustion_for_routing
3313 );
3314 UsageProviderRefreshOutcome::Refreshed
3315 }
3316 Err(err) => {
3317 let error = err.to_string();
3318 let terminal_failure = usage_provider_error_is_terminal(&error);
3319 if terminal_failure {
3320 remember_usage_provider_target_suppression(
3321 &provider.id,
3322 target,
3323 USAGE_PROVIDER_TERMINAL_FAILURE_TTL,
3324 error.clone(),
3325 true,
3326 Instant::now(),
3327 );
3328 }
3329 state
3330 .record_provider_balance_snapshot(
3331 service_name,
3332 base_snapshot(provider, &upstreams[0], fetched_at_ms, stale_after_ms)
3333 .with_error(error.clone()),
3334 )
3335 .await;
3336 update_usage_exhausted(
3337 lb_states,
3338 state,
3339 cfg,
3340 service_name,
3341 &upstreams,
3342 terminal_failure,
3343 )
3344 .await;
3345 warn!(
3346 "usage provider '{}' poll failed for {}[{}]: {}",
3347 provider.id, target.upstream.station_name, target.upstream.index, error
3348 );
3349 UsageProviderRefreshOutcome::Failed
3350 }
3351 }
3352}
3353
3354struct ConfiguredRefreshJob<'a> {
3355 provider: &'a UsageProviderConfig,
3356 target: UsageProviderTarget,
3357 interval_secs: u64,
3358}
3359
3360struct AutoRefreshJob {
3361 target: UsageProviderTarget,
3362}
3363
3364async fn run_configured_refresh_job<'a>(
3365 client: &'a Client,
3366 job: ConfiguredRefreshJob<'a>,
3367 cfg: &'a ProxyConfig,
3368 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
3369 state: &'a Arc<ProxyState>,
3370 service_name: &'a str,
3371) -> (String, UsageProviderRefreshOutcome) {
3372 let provider_id = job.provider.id.clone();
3373 let outcome = refresh_provider_target(RefreshProviderTargetParams {
3374 client,
3375 provider: job.provider,
3376 target: &job.target,
3377 cfg,
3378 lb_states,
3379 state,
3380 service_name,
3381 interval_secs: job.interval_secs,
3382 })
3383 .await;
3384 (provider_id, outcome)
3385}
3386
3387async fn run_auto_refresh_job(
3388 client: &Client,
3389 job: AutoRefreshJob,
3390 cfg: &ProxyConfig,
3391 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
3392 state: &Arc<ProxyState>,
3393 service_name: &str,
3394) -> UsageProviderRefreshOutcome {
3395 auto_probe_provider_target(client, &job.target, cfg, lb_states, state, service_name).await
3396}
3397
3398async fn run_configured_refresh_jobs<'a>(
3399 client: &'a Client,
3400 jobs: Vec<ConfiguredRefreshJob<'a>>,
3401 cfg: &'a ProxyConfig,
3402 lb_states: &'a Arc<Mutex<HashMap<String, LbState>>>,
3403 state: &'a Arc<ProxyState>,
3404 service_name: &'a str,
3405) -> Vec<(String, UsageProviderRefreshOutcome)> {
3406 let mut pending = jobs.into_iter();
3407 let mut running = FuturesUnordered::new();
3408 let mut results = Vec::new();
3409 let concurrency = BALANCE_REFRESH_CONCURRENCY.max(1);
3410
3411 for _ in 0..concurrency {
3412 let Some(job) = pending.next() else {
3413 break;
3414 };
3415 running.push(run_configured_refresh_job(
3416 client,
3417 job,
3418 cfg,
3419 lb_states,
3420 state,
3421 service_name,
3422 ));
3423 }
3424
3425 while let Some(result) = running.next().await {
3426 results.push(result);
3427 if let Some(job) = pending.next() {
3428 running.push(run_configured_refresh_job(
3429 client,
3430 job,
3431 cfg,
3432 lb_states,
3433 state,
3434 service_name,
3435 ));
3436 }
3437 }
3438
3439 results
3440}
3441
3442async fn run_auto_refresh_jobs(
3443 client: &Client,
3444 jobs: Vec<AutoRefreshJob>,
3445 cfg: &ProxyConfig,
3446 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
3447 state: &Arc<ProxyState>,
3448 service_name: &str,
3449) -> Vec<UsageProviderRefreshOutcome> {
3450 let mut pending = jobs.into_iter();
3451 let mut running = FuturesUnordered::new();
3452 let mut results = Vec::new();
3453 let concurrency = BALANCE_REFRESH_CONCURRENCY.max(1);
3454
3455 for _ in 0..concurrency {
3456 let Some(job) = pending.next() else {
3457 break;
3458 };
3459 running.push(run_auto_refresh_job(
3460 client,
3461 job,
3462 cfg,
3463 lb_states,
3464 state,
3465 service_name,
3466 ));
3467 }
3468
3469 while let Some(result) = running.next().await {
3470 results.push(result);
3471 if let Some(job) = pending.next() {
3472 running.push(run_auto_refresh_job(
3473 client,
3474 job,
3475 cfg,
3476 lb_states,
3477 state,
3478 service_name,
3479 ));
3480 }
3481 }
3482
3483 results
3484}
3485
3486fn auto_snapshot_is_usable(snapshot: &ProviderBalanceSnapshot) -> bool {
3487 snapshot.error.is_none()
3488 && matches!(
3489 snapshot.status,
3490 BalanceSnapshotStatus::Ok | BalanceSnapshotStatus::Exhausted
3491 )
3492}
3493
3494fn normalized_error_text(value: &str) -> String {
3495 value
3496 .chars()
3497 .map(|ch| match ch {
3498 '_' | '-' | '.' | '/' | ':' | '[' | ']' | '(' | ')' | ',' | ';' => ' ',
3499 _ => ch.to_ascii_lowercase(),
3500 })
3501 .collect::<String>()
3502}
3503
3504fn usage_provider_error_is_terminal(error: &str) -> bool {
3505 let normalized = normalized_error_text(error);
3506 let terminal_markers = [
3507 "user inactive",
3508 "user account is not active",
3509 "account is not active",
3510 "account inactive",
3511 "account disabled",
3512 "user disabled",
3513 "api key disabled",
3514 "api key is disabled",
3515 "api key inactive",
3516 "api key is not active",
3517 "key inactive",
3518 "key disabled",
3519 "invalid api key",
3520 "invalid token",
3521 "invalid bearer token",
3522 "token invalid",
3523 "unauthorized api key",
3524 "insufficient balance",
3525 "balance insufficient",
3526 "insufficient quota",
3527 "quota exhausted",
3528 "quota exceeded",
3529 "no balance",
3530 "余额不足",
3531 "额度不足",
3532 "配额不足",
3533 "账户未激活",
3534 "账号未激活",
3535 "用户未激活",
3536 "账户已禁用",
3537 "账号已禁用",
3538 "用户已禁用",
3539 "密钥无效",
3540 "令牌无效",
3541 ];
3542 terminal_markers
3543 .iter()
3544 .any(|marker| normalized.contains(marker))
3545}
3546
3547fn quota_period_is_current_day(period: &str) -> bool {
3548 let normalized = period.trim().to_ascii_lowercase();
3549 matches!(
3550 normalized.as_str(),
3551 "daily" | "day" | "today" | "current_day" | "current-day" | "1d" | "24h" | "今日" | "今天"
3552 )
3553}
3554
3555fn quota_period_is_refreshable_window(period: &str) -> bool {
3556 let normalized = period.trim().to_ascii_lowercase();
3557 quota_period_is_current_day(&normalized)
3558 || matches!(
3559 normalized.as_str(),
3560 "weekly" | "week" | "7d" | "monthly" | "month"
3561 )
3562 || normalized.starts_with("rate_limit:")
3563}
3564
3565fn snapshot_has_current_day_quota_exhaustion(snapshot: &ProviderBalanceSnapshot) -> bool {
3566 snapshot.exhausted == Some(true)
3567 && snapshot
3568 .quota_period
3569 .as_deref()
3570 .is_some_and(quota_period_is_current_day)
3571}
3572
3573fn snapshot_has_refreshable_window_exhaustion(snapshot: &ProviderBalanceSnapshot) -> bool {
3574 snapshot.exhausted == Some(true)
3575 && snapshot
3576 .quota_period
3577 .as_deref()
3578 .is_some_and(quota_period_is_refreshable_window)
3579}
3580
3581fn duration_millis_u64(duration: Duration) -> u64 {
3582 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
3583}
3584
3585fn duration_until_ms(deadline_ms: u64, now_ms: u64) -> Option<Duration> {
3586 deadline_ms
3587 .checked_sub(now_ms)
3588 .filter(|remaining_ms| *remaining_ms > 0)
3589 .map(Duration::from_millis)
3590}
3591
3592fn snapshot_freshness_ttl(snapshot: &ProviderBalanceSnapshot, now_ms: u64) -> Option<Duration> {
3593 snapshot
3594 .stale_after_ms
3595 .and_then(|stale_after_ms| duration_until_ms(stale_after_ms, now_ms))
3596}
3597
3598fn current_day_quota_suppression_ttl(
3599 snapshot: &ProviderBalanceSnapshot,
3600 cfg: &ProxyConfig,
3601 now_ms: u64,
3602) -> Option<Duration> {
3603 let reset_at_ms = next_reset_at_ms(
3604 snapshot.fetched_at_ms,
3605 cfg.ui.usage_forecast.reset_utc_offset.as_str(),
3606 cfg.ui.usage_forecast.reset_time.as_str(),
3607 )?;
3608 let suppress_until_ms = reset_at_ms.saturating_add(duration_millis_u64(
3609 USAGE_PROVIDER_DAILY_RESET_SUPPRESSION_GRACE,
3610 ));
3611 duration_until_ms(suppress_until_ms, now_ms)
3612}
3613
3614fn usage_provider_snapshot_suppression_ttl(
3615 snapshot: &ProviderBalanceSnapshot,
3616 cfg: &ProxyConfig,
3617 now_ms: u64,
3618) -> Option<Duration> {
3619 if let Some(reset_at_ms) = snapshot.quota_resets_at_ms {
3620 let suppress_until_ms = reset_at_ms.saturating_add(duration_millis_u64(
3621 USAGE_PROVIDER_DAILY_RESET_SUPPRESSION_GRACE,
3622 ));
3623 return duration_until_ms(suppress_until_ms, now_ms)
3624 .or_else(|| snapshot_freshness_ttl(snapshot, now_ms));
3625 }
3626
3627 if snapshot_has_current_day_quota_exhaustion(snapshot) {
3628 return current_day_quota_suppression_ttl(snapshot, cfg, now_ms)
3629 .or_else(|| snapshot_freshness_ttl(snapshot, now_ms));
3630 }
3631
3632 if snapshot_has_refreshable_window_exhaustion(snapshot) {
3633 return snapshot_freshness_ttl(snapshot, now_ms);
3634 }
3635
3636 if snapshot.stale_at(now_ms) {
3637 None
3638 } else {
3639 Some(USAGE_PROVIDER_EXHAUSTED_SUPPRESSION_TTL)
3640 }
3641}
3642
3643fn usage_provider_snapshot_suppression_reason(
3644 snapshot: &ProviderBalanceSnapshot,
3645) -> Option<String> {
3646 if snapshot.status_at(snapshot.fetched_at_ms) != BalanceSnapshotStatus::Exhausted {
3647 return None;
3648 }
3649
3650 if snapshot_has_current_day_quota_exhaustion(snapshot) {
3651 let period = snapshot.quota_period.as_deref().unwrap_or("daily");
3652 return Some(format!(
3653 "{period} package quota exhausted for current period"
3654 ));
3655 }
3656
3657 if snapshot_has_refreshable_window_exhaustion(snapshot) {
3658 let period = snapshot.quota_period.as_deref().unwrap_or("usage");
3659 return Some(format!(
3660 "{period} usage window exhausted for current period"
3661 ));
3662 }
3663
3664 if snapshot.routing_exhausted() {
3665 return Some("balance exhausted".to_string());
3666 }
3667
3668 None
3669}
3670
3671fn usage_provider_snapshot_error(snapshot: &ProviderBalanceSnapshot) -> Option<&str> {
3672 snapshot
3673 .error
3674 .as_deref()
3675 .map(str::trim)
3676 .filter(|value| !value.is_empty())
3677}
3678
3679fn usage_provider_snapshot_terminal_error(snapshot: &ProviderBalanceSnapshot) -> Option<&str> {
3680 usage_provider_snapshot_error(snapshot).filter(|error| usage_provider_error_is_terminal(error))
3681}
3682
3683fn usage_provider_suppression_decision_from_snapshot(
3684 snapshot: &ProviderBalanceSnapshot,
3685 cfg: &ProxyConfig,
3686 now_ms: u64,
3687) -> Option<ProviderTargetSuppressionDecision> {
3688 if let Some(error) = usage_provider_snapshot_terminal_error(snapshot) {
3689 if snapshot.stale_at(now_ms) {
3690 return None;
3691 }
3692 return Some(ProviderTargetSuppressionDecision {
3693 reason: error.to_string(),
3694 routing_exhausted: true,
3695 ttl: USAGE_PROVIDER_TERMINAL_FAILURE_TTL,
3696 });
3697 }
3698
3699 let reason = usage_provider_snapshot_suppression_reason(snapshot)?;
3700 usage_provider_snapshot_suppression_ttl(snapshot, cfg, now_ms).map(|ttl| {
3701 ProviderTargetSuppressionDecision {
3702 reason,
3703 routing_exhausted: true,
3704 ttl,
3705 }
3706 })
3707}
3708
3709fn provider_balance_snapshot_matches_target(
3710 snapshot: &ProviderBalanceSnapshot,
3711 provider_id: &str,
3712 target: &UsageProviderTarget,
3713) -> bool {
3714 if snapshot.provider_id != provider_id {
3715 return false;
3716 }
3717 if snapshot.station_name.as_deref() != Some(target.upstream.station_name.as_str()) {
3718 return false;
3719 }
3720 if snapshot.upstream_index != Some(target.upstream.index) {
3721 return false;
3722 }
3723
3724 match (
3725 snapshot.provider_endpoint_key.as_deref(),
3726 target
3727 .upstream
3728 .provider_endpoint
3729 .as_ref()
3730 .map(ProviderEndpointKey::stable_key),
3731 ) {
3732 (Some(snapshot_key), Some(target_key)) => snapshot_key == target_key,
3733 _ => true,
3734 }
3735}
3736
3737async fn existing_usage_provider_target_suppression_decision(
3738 state: &Arc<ProxyState>,
3739 cfg: &ProxyConfig,
3740 service_name: &str,
3741 provider_id: &str,
3742 target: &UsageProviderTarget,
3743 now_ms: u64,
3744) -> Option<ProviderTargetSuppressionDecision> {
3745 let view = state.get_provider_balance_view(service_name).await;
3746 view.get(&target.upstream.station_name)
3747 .and_then(|snapshots| {
3748 snapshots
3749 .iter()
3750 .filter(|snapshot| {
3751 provider_balance_snapshot_matches_target(snapshot, provider_id, target)
3752 })
3753 .max_by_key(|snapshot| snapshot.fetched_at_ms)
3754 .and_then(|snapshot| {
3755 usage_provider_suppression_decision_from_snapshot(snapshot, cfg, now_ms)
3756 })
3757 })
3758}
3759
3760fn auto_probe_error_summary(probe_errors: &[String]) -> Option<String> {
3761 (!probe_errors.is_empty()).then(|| format!("attempts failed: {}", probe_errors.join("; ")))
3762}
3763
3764async fn auto_probe_provider_target(
3765 client: &Client,
3766 target: &UsageProviderTarget,
3767 cfg: &ProxyConfig,
3768 lb_states: &Arc<Mutex<HashMap<String, LbState>>>,
3769 state: &Arc<ProxyState>,
3770 service_name: &str,
3771) -> UsageProviderRefreshOutcome {
3772 let upstreams = vec![target.upstream.clone()];
3773 let fetched_at_ms = unix_now_ms();
3774 let interval_secs = DEFAULT_POLL_INTERVAL_SECS;
3775 let stale_after_ms = stale_after_ms(fetched_at_ms, interval_secs);
3776
3777 if is_official_openai_base_url(&target.base_url) {
3778 let provider = auto_openai_official_provider(target);
3779 let Some(token) = resolve_token(&provider, &upstreams, cfg, service_name) else {
3780 state
3781 .record_provider_balance_snapshot(
3782 service_name,
3783 base_snapshot(&provider, &target.upstream, fetched_at_ms, stale_after_ms),
3784 )
3785 .await;
3786 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3787 warn!(
3788 "OpenAI organization costs require OPENAI_ADMIN_KEY; balance stays unknown for {}[{}]",
3789 target.upstream.station_name, target.upstream.index
3790 );
3791 return UsageProviderRefreshOutcome::MissingToken;
3792 };
3793
3794 return match poll_provider_http_json(client, &provider, &target.base_url, &token).await {
3795 Ok(value) => {
3796 let snapshot = snapshot_from_provider_json(
3797 &provider,
3798 &upstreams[0],
3799 &value,
3800 &target.base_url,
3801 fetched_at_ms,
3802 stale_after_ms,
3803 );
3804 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false)
3805 .await;
3806 state
3807 .record_provider_balance_snapshot(service_name, snapshot)
3808 .await;
3809 UsageProviderRefreshOutcome::Refreshed
3810 }
3811 Err(err) => {
3812 state
3813 .record_provider_balance_snapshot(
3814 service_name,
3815 base_snapshot(&provider, &upstreams[0], fetched_at_ms, stale_after_ms)
3816 .with_error(err.to_string()),
3817 )
3818 .await;
3819 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false)
3820 .await;
3821 warn!(
3822 "OpenAI organization costs poll failed for {}[{}]: {}",
3823 target.upstream.station_name, target.upstream.index, err
3824 );
3825 UsageProviderRefreshOutcome::Failed
3826 }
3827 };
3828 }
3829
3830 let first_provider = auto_usage_provider(target, first_auto_probe_kind(target));
3831 let provider_id = first_provider.id.clone();
3832 if let Some(suppression) =
3833 usage_provider_target_suppression_active(&provider_id, target, Instant::now())
3834 {
3835 update_usage_exhausted(
3836 lb_states,
3837 state,
3838 cfg,
3839 service_name,
3840 &upstreams,
3841 suppression.routing_exhausted,
3842 )
3843 .await;
3844 warn!(
3845 "auto usage provider '{}' skipped {}[{}]: balance refresh suppressed: {}",
3846 first_provider.id,
3847 target.upstream.station_name,
3848 target.upstream.index,
3849 suppression.reason
3850 );
3851 return UsageProviderRefreshOutcome::Failed;
3852 }
3853 if let Some(decision) = existing_usage_provider_target_suppression_decision(
3854 state,
3855 cfg,
3856 service_name,
3857 &provider_id,
3858 target,
3859 fetched_at_ms,
3860 )
3861 .await
3862 {
3863 remember_usage_provider_target_suppression(
3864 &provider_id,
3865 target,
3866 decision.ttl,
3867 decision.reason.clone(),
3868 decision.routing_exhausted,
3869 Instant::now(),
3870 );
3871 update_usage_exhausted(
3872 lb_states,
3873 state,
3874 cfg,
3875 service_name,
3876 &upstreams,
3877 decision.routing_exhausted,
3878 )
3879 .await;
3880 warn!(
3881 "auto usage provider '{}' skipped {}[{}]: existing balance snapshot suppresses refresh: {}",
3882 first_provider.id, target.upstream.station_name, target.upstream.index, decision.reason
3883 );
3884 return UsageProviderRefreshOutcome::Failed;
3885 }
3886
3887 let Some(token) = resolve_token(&first_provider, &upstreams, cfg, service_name) else {
3888 state
3889 .record_provider_balance_snapshot(
3890 service_name,
3891 base_snapshot(
3892 &first_provider,
3893 &target.upstream,
3894 fetched_at_ms,
3895 stale_after_ms,
3896 )
3897 .with_error("no usable token; checked upstream auth"),
3898 )
3899 .await;
3900 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3901 return UsageProviderRefreshOutcome::MissingToken;
3902 };
3903
3904 let probe_order = auto_probe_kind_order(&provider_id, target);
3905 if probe_order.is_empty() {
3906 let error = "all balance probe kinds are temporarily suppressed";
3907 warn!(
3908 "auto usage provider '{}' skipped {}[{}]: {}",
3909 first_provider.id, target.upstream.station_name, target.upstream.index, error
3910 );
3911 state
3912 .record_provider_balance_snapshot(
3913 service_name,
3914 base_snapshot(
3915 &first_provider,
3916 &target.upstream,
3917 fetched_at_ms,
3918 stale_after_ms,
3919 )
3920 .with_error(error),
3921 )
3922 .await;
3923 update_usage_exhausted(lb_states, state, cfg, service_name, &upstreams, false).await;
3924 return UsageProviderRefreshOutcome::Failed;
3925 }
3926
3927 let mut probe_errors = Vec::new();
3928 for kind in probe_order {
3929 let provider = auto_usage_provider(target, kind);
3930 match poll_provider_http_json(client, &provider, &target.base_url, &token).await {
3931 Ok(value) => {
3932 let snapshot = snapshot_from_provider_json(
3933 &provider,
3934 &upstreams[0],
3935 &value,
3936 &target.base_url,
3937 fetched_at_ms,
3938 stale_after_ms,
3939 );
3940 if auto_snapshot_is_usable(&snapshot) {
3941 remember_auto_probe_kind_success(&provider_id, target, kind);
3942 let suppression_decision = usage_provider_suppression_decision_from_snapshot(
3943 &snapshot,
3944 cfg,
3945 fetched_at_ms,
3946 );
3947 let exhausted_for_lb = suppression_decision
3948 .as_ref()
3949 .is_some_and(|decision| decision.routing_exhausted);
3950 if let Some(decision) = suppression_decision.as_ref() {
3951 remember_usage_provider_target_suppression(
3952 &provider_id,
3953 target,
3954 decision.ttl,
3955 decision.reason.as_str(),
3956 decision.routing_exhausted,
3957 Instant::now(),
3958 );
3959 } else {
3960 clear_usage_provider_target_suppression(&provider_id, target);
3961 }
3962 update_usage_exhausted(
3963 lb_states,
3964 state,
3965 cfg,
3966 service_name,
3967 &upstreams,
3968 exhausted_for_lb,
3969 )
3970 .await;
3971 state
3972 .record_provider_balance_snapshot(service_name, snapshot)
3973 .await;
3974 info!(
3975 "auto usage provider '{}' refreshed {}[{}] via {:?}, exhausted = {}",
3976 provider.id,
3977 target.upstream.station_name,
3978 target.upstream.index,
3979 kind,
3980 exhausted_for_lb
3981 );
3982 return UsageProviderRefreshOutcome::Refreshed;
3983 }
3984 remember_auto_probe_kind_failure(&provider_id, target, kind, Instant::now());
3985 let error = snapshot.error.unwrap_or_else(|| {
3986 format!("auto probe {:?} returned no usable balance fields", kind)
3987 });
3988 probe_errors.push(format!("{:?}: {}", kind, error));
3989 }
3990 Err(err) => {
3991 remember_auto_probe_kind_failure(&provider_id, target, kind, Instant::now());
3992 probe_errors.push(format!("{:?}: {}", kind, err));
3993 }
3994 }
3995 }
3996
3997 if let Some(error) = auto_probe_error_summary(&probe_errors) {
3998 warn!(
3999 "auto usage provider '{}' found no usable balance endpoint for {}[{}]: {}",
4000 first_provider.id, target.upstream.station_name, target.upstream.index, error
4001 );
4002 state
4003 .record_provider_balance_snapshot(
4004 service_name,
4005 base_snapshot(
4006 &first_provider,
4007 &target.upstream,
4008 fetched_at_ms,
4009 stale_after_ms,
4010 )
4011 .with_error(error.clone()),
4012 )
4013 .await;
4014 let terminal_failure = usage_provider_error_is_terminal(error.as_str());
4015 if terminal_failure {
4016 remember_usage_provider_target_suppression(
4017 &provider_id,
4018 target,
4019 USAGE_PROVIDER_TERMINAL_FAILURE_TTL,
4020 error.clone(),
4021 true,
4022 Instant::now(),
4023 );
4024 }
4025 update_usage_exhausted(
4026 lb_states,
4027 state,
4028 cfg,
4029 service_name,
4030 &upstreams,
4031 terminal_failure,
4032 )
4033 .await;
4034 }
4035 UsageProviderRefreshOutcome::Failed
4036}
4037
4038pub async fn refresh_balances_for_service(
4039 client: &Client,
4040 cfg: Arc<ProxyConfig>,
4041 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
4042 state: Arc<ProxyState>,
4043 service_name: &str,
4044 station_name_filter: Option<&str>,
4045 provider_id_filter: Option<&str>,
4046) -> UsageProviderRefreshSummary {
4047 if cfg!(test) {
4049 return UsageProviderRefreshSummary::default();
4050 }
4051
4052 let station_name_filter = station_name_filter
4053 .map(str::trim)
4054 .filter(|value| !value.is_empty());
4055 let provider_id_filter = provider_id_filter
4056 .map(str::trim)
4057 .filter(|value| !value.is_empty());
4058 let providers_file = load_providers();
4059 let mut summary = UsageProviderRefreshSummary {
4060 providers_configured: providers_file.providers.len(),
4061 ..UsageProviderRefreshSummary::default()
4062 };
4063
4064 let poll_map = LAST_USAGE_POLL.get_or_init(|| Mutex::new(HashMap::new()));
4065 let mut configured_jobs = Vec::new();
4066 let mut configured_job_keys = HashSet::new();
4067 for provider in &providers_file.providers {
4068 if provider_id_filter.is_some_and(|filter| filter != provider.id.as_str()) {
4069 continue;
4070 }
4071
4072 let targets = matching_provider_targets(&cfg, service_name, provider, station_name_filter);
4073 if targets.is_empty() {
4074 continue;
4075 }
4076
4077 summary.providers_matched += 1;
4078 summary.upstreams_matched += targets.len();
4079 warn_if_provider_spans_hosts(&cfg, service_name, provider);
4080
4081 let interval_secs = snapshot_refresh_interval_secs(provider);
4082 for target in targets {
4083 summary.attempted += 1;
4084 configured_job_keys.insert(target_key(&target));
4085 configured_jobs.push(ConfiguredRefreshJob {
4086 provider,
4087 target,
4088 interval_secs,
4089 });
4090 }
4091 }
4092
4093 let mut refreshed_provider_ids = HashSet::new();
4094 if !configured_jobs.is_empty() {
4095 for (provider_id, outcome) in run_configured_refresh_jobs(
4096 client,
4097 configured_jobs,
4098 &cfg,
4099 &lb_states,
4100 &state,
4101 service_name,
4102 )
4103 .await
4104 {
4105 match outcome {
4106 UsageProviderRefreshOutcome::Refreshed => {
4107 summary.refreshed += 1;
4108 refreshed_provider_ids.insert(provider_id);
4109 }
4110 UsageProviderRefreshOutcome::Failed => summary.failed += 1,
4111 UsageProviderRefreshOutcome::MissingToken => summary.missing_token += 1,
4112 }
4113 }
4114 }
4115
4116 let mut auto_jobs = Vec::new();
4117 for target in usage_provider_targets(&cfg, service_name, station_name_filter) {
4118 if configured_job_keys.contains(&target_key(&target)) {
4119 continue;
4120 }
4121 if !auto_target_matches_provider_id_filter(&target, provider_id_filter) {
4122 continue;
4123 }
4124
4125 summary.attempted += 1;
4126 summary.auto_attempted += 1;
4127 auto_jobs.push(AutoRefreshJob { target });
4128 }
4129
4130 if !auto_jobs.is_empty() {
4131 for outcome in
4132 run_auto_refresh_jobs(client, auto_jobs, &cfg, &lb_states, &state, service_name).await
4133 {
4134 match outcome {
4135 UsageProviderRefreshOutcome::Refreshed => {
4136 summary.refreshed += 1;
4137 summary.auto_refreshed += 1;
4138 }
4139 UsageProviderRefreshOutcome::Failed => {
4140 summary.failed += 1;
4141 summary.auto_failed += 1;
4142 }
4143 UsageProviderRefreshOutcome::MissingToken => {
4144 summary.missing_token += 1;
4145 }
4146 }
4147 }
4148 }
4149
4150 if !refreshed_provider_ids.is_empty()
4151 && let Ok(mut map) = poll_map.lock()
4152 {
4153 let now = Instant::now();
4154 for provider_id in refreshed_provider_ids {
4155 map.insert(provider_id, now);
4156 }
4157 }
4158
4159 summary
4160}
4161
4162pub fn enqueue_poll_for_codex_provider_endpoint(
4166 client: Client,
4167 cfg: Arc<ProxyConfig>,
4168 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
4169 state: Arc<ProxyState>,
4170 service_name: &str,
4171 provider_endpoint: ProviderEndpointKey,
4172) {
4173 let key = provider_endpoint.clone();
4174 let Some(initial_sleep_for) = enqueue_request_balance_refresh(key.clone()) else {
4175 return;
4176 };
4177
4178 let service_name = service_name.to_string();
4179 tokio::spawn(async move {
4180 let mut sleep_for = initial_sleep_for;
4181 loop {
4182 tokio::time::sleep(sleep_for).await;
4183 match take_request_balance_refresh_if_due(&key) {
4184 RequestBalanceQueueDue::Due => {}
4185 RequestBalanceQueueDue::NotDue(delay) => {
4186 sleep_for = delay;
4187 continue;
4188 }
4189 RequestBalanceQueueDue::Missing => return,
4190 }
4191
4192 let current_target = usage_provider_target_for_provider_endpoint(
4193 &cfg,
4194 service_name.as_str(),
4195 &provider_endpoint,
4196 );
4197 match poll_for_codex_target(
4198 client.clone(),
4199 cfg.clone(),
4200 lb_states.clone(),
4201 state.clone(),
4202 service_name.as_str(),
4203 current_target,
4204 )
4205 .await
4206 {
4207 RequestBalancePollOutcome::Attempted | RequestBalancePollOutcome::Skipped => {
4208 return;
4209 }
4210 RequestBalancePollOutcome::Deferred(delay) => {
4211 schedule_request_balance_refresh_at(key.clone(), Instant::now() + delay);
4212 sleep_for = delay;
4213 }
4214 }
4215 }
4216 });
4217}
4218
4219async fn poll_for_codex_target(
4220 client: Client,
4221 cfg: Arc<ProxyConfig>,
4222 lb_states: Arc<Mutex<HashMap<String, LbState>>>,
4223 state: Arc<ProxyState>,
4224 service_name: &str,
4225 current_target: Option<UsageProviderTarget>,
4226) -> RequestBalancePollOutcome {
4227 if cfg!(test) {
4230 return RequestBalancePollOutcome::Skipped;
4231 }
4232
4233 let providers_file = load_providers();
4234 let Some(current_target) = current_target else {
4235 return RequestBalancePollOutcome::Skipped;
4236 };
4237
4238 let now = Instant::now();
4239 let poll_map = LAST_USAGE_POLL.get_or_init(|| Mutex::new(HashMap::new()));
4240 let mut matched_configured_provider = false;
4241 let mut configured_jobs = Vec::new();
4242 let mut next_cooldown = None::<Duration>;
4243
4244 for provider in &providers_file.providers {
4245 if !domain_matches(¤t_target.base_url, &provider.domains) {
4246 continue;
4247 }
4248 matched_configured_provider = true;
4249
4250 let Some(interval_secs) = effective_poll_interval_secs(provider) else {
4251 continue;
4252 };
4253
4254 {
4255 let mut map = match poll_map.lock() {
4256 Ok(m) => m,
4257 Err(_) => continue,
4258 };
4259 if let Some(last) = map.get(&provider.id)
4260 && let Some(cooldown) = remaining_poll_cooldown(*last, interval_secs, now)
4261 {
4262 next_cooldown =
4263 Some(next_cooldown.map_or(cooldown, |existing| existing.min(cooldown)));
4264 continue;
4265 }
4266 map.insert(provider.id.clone(), now);
4267 }
4268
4269 warn_if_provider_spans_hosts(&cfg, service_name, provider);
4270 configured_jobs.push(ConfiguredRefreshJob {
4271 provider,
4272 target: current_target.clone(),
4273 interval_secs,
4274 });
4275 }
4276
4277 if !configured_jobs.is_empty() {
4278 let _ = run_configured_refresh_jobs(
4279 &client,
4280 configured_jobs,
4281 &cfg,
4282 &lb_states,
4283 &state,
4284 service_name,
4285 )
4286 .await;
4287 return RequestBalancePollOutcome::Attempted;
4288 }
4289
4290 if matched_configured_provider {
4291 return next_cooldown
4292 .map(RequestBalancePollOutcome::Deferred)
4293 .unwrap_or(RequestBalancePollOutcome::Skipped);
4294 }
4295
4296 let auto_provider = if is_official_openai_base_url(¤t_target.base_url) {
4297 auto_openai_official_provider(¤t_target)
4298 } else {
4299 auto_usage_provider(¤t_target, first_auto_probe_kind(¤t_target))
4300 };
4301 let Some(interval_secs) = effective_poll_interval_secs(&auto_provider) else {
4302 return RequestBalancePollOutcome::Skipped;
4303 };
4304
4305 {
4306 let mut map = match poll_map.lock() {
4307 Ok(m) => m,
4308 Err(_) => return RequestBalancePollOutcome::Skipped,
4309 };
4310 if let Some(last) = map.get(&auto_provider.id)
4311 && let Some(cooldown) = remaining_poll_cooldown(*last, interval_secs, now)
4312 {
4313 return RequestBalancePollOutcome::Deferred(cooldown);
4314 }
4315 map.insert(auto_provider.id.clone(), now);
4316 }
4317
4318 let _ = auto_probe_provider_target(
4319 &client,
4320 ¤t_target,
4321 &cfg,
4322 &lb_states,
4323 &state,
4324 service_name,
4325 )
4326 .await;
4327 RequestBalancePollOutcome::Attempted
4328}
4329
4330#[cfg(test)]
4331mod tests {
4332 use super::*;
4333
4334 use crate::balance::BalanceSnapshotStatus;
4335 use crate::config::{ServiceConfig, UpstreamAuth, UpstreamConfig};
4336 use axum::routing::get;
4337 use std::net::SocketAddr;
4338 use std::sync::atomic::{AtomicUsize, Ordering};
4339 use tokio::net::TcpListener;
4340
4341 async fn spawn_axum_server(app: axum::Router) -> (SocketAddr, tokio::task::JoinHandle<()>) {
4342 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
4343 let addr = listener.local_addr().expect("local_addr");
4344 let handle = tokio::spawn(async move {
4345 axum::serve(listener, app).await.expect("serve");
4346 });
4347 (addr, handle)
4348 }
4349
4350 fn provider(id: &str, kind: ProviderKind) -> UsageProviderConfig {
4351 UsageProviderConfig {
4352 id: id.to_string(),
4353 kind,
4354 domains: vec!["example.com".to_string()],
4355 endpoint: "https://example.com/usage".to_string(),
4356 token_env: None,
4357 require_token_env: false,
4358 poll_interval_secs: Some(60),
4359 refresh_on_request: true,
4360 trust_exhaustion_for_routing: true,
4361 headers: BTreeMap::new(),
4362 variables: BTreeMap::new(),
4363 extract: UsageProviderExtractConfig::default(),
4364 }
4365 }
4366
4367 fn upstream() -> UpstreamRef {
4368 UpstreamRef {
4369 station_name: "right".to_string(),
4370 index: 1,
4371 provider_endpoint: None,
4372 }
4373 }
4374
4375 fn endpoint_upstream() -> UpstreamRef {
4376 UpstreamRef {
4377 station_name: "right".to_string(),
4378 index: 1,
4379 provider_endpoint: Some(ProviderEndpointKey::new("codex", "right", "default")),
4380 }
4381 }
4382
4383 fn upstream_config(base_url: &str) -> UpstreamConfig {
4384 UpstreamConfig {
4385 base_url: base_url.to_string(),
4386 auth: UpstreamAuth::default(),
4387 tags: HashMap::new(),
4388 supported_models: HashMap::new(),
4389 model_mapping: HashMap::new(),
4390 }
4391 }
4392
4393 fn endpoint_upstream_config(
4394 base_url: &str,
4395 provider_id: &str,
4396 endpoint_id: &str,
4397 ) -> UpstreamConfig {
4398 let mut upstream = upstream_config(base_url);
4399 upstream
4400 .tags
4401 .insert("provider_id".to_string(), provider_id.to_string());
4402 upstream
4403 .tags
4404 .insert("endpoint_id".to_string(), endpoint_id.to_string());
4405 upstream
4406 }
4407
4408 fn service_config(name: &str, upstreams: Vec<UpstreamConfig>) -> ServiceConfig {
4409 ServiceConfig {
4410 name: name.to_string(),
4411 alias: None,
4412 enabled: true,
4413 level: 1,
4414 upstreams,
4415 }
4416 }
4417
4418 fn proxy_config(stations: Vec<ServiceConfig>) -> ProxyConfig {
4419 let mut cfg = ProxyConfig::default();
4420 cfg.codex.configs = stations
4421 .into_iter()
4422 .map(|station| (station.name.clone(), station))
4423 .collect();
4424 cfg
4425 }
4426
4427 fn usage_provider_target(base_url: &str, provider_id: &str) -> UsageProviderTarget {
4428 UsageProviderTarget {
4429 upstream: UpstreamRef {
4430 station_name: "routing".to_string(),
4431 index: 0,
4432 provider_endpoint: Some(ProviderEndpointKey::new("codex", provider_id, "default")),
4433 },
4434 base_url: base_url.to_string(),
4435 provider_id: Some(provider_id.to_string()),
4436 }
4437 }
4438
4439 fn usage_provider_target_at(
4440 station_name: &str,
4441 upstream_index: usize,
4442 base_url: &str,
4443 provider_id: &str,
4444 ) -> UsageProviderTarget {
4445 UsageProviderTarget {
4446 upstream: UpstreamRef {
4447 station_name: station_name.to_string(),
4448 index: upstream_index,
4449 provider_endpoint: Some(ProviderEndpointKey::new("codex", provider_id, "default")),
4450 },
4451 base_url: base_url.to_string(),
4452 provider_id: Some(provider_id.to_string()),
4453 }
4454 }
4455
4456 fn clear_auto_probe_kind_state(provider_id: &str) {
4457 if let Some(hints) = AUTO_PROBE_KIND_HINTS.get()
4458 && let Ok(mut hints) = hints.lock()
4459 {
4460 hints.remove(provider_id);
4461 }
4462 if let Some(failures) = AUTO_PROBE_KIND_FAILURES.get()
4463 && let Ok(mut failures) = failures.lock()
4464 {
4465 failures.retain(|key, _| key.provider_id != provider_id);
4466 }
4467 clear_usage_provider_target_suppressions_for_provider(provider_id);
4468 }
4469
4470 #[test]
4471 fn budget_snapshot_reports_monthly_budget_and_exhaustion() {
4472 let snapshot = budget_snapshot_from_json(
4473 &provider("packycode", ProviderKind::BudgetHttpJson),
4474 &upstream(),
4475 &serde_json::json!({
4476 "monthly_budget_usd": "10.50",
4477 "monthly_spent_usd": 10.5
4478 }),
4479 100,
4480 Some(1_000),
4481 );
4482
4483 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
4484 assert_eq!(snapshot.exhausted, Some(true));
4485 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10.5"));
4486 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("10.5"));
4487 }
4488
4489 #[test]
4490 fn budget_snapshot_keeps_missing_amounts_unknown() {
4491 let snapshot = budget_snapshot_from_json(
4492 &provider("packycode", ProviderKind::BudgetHttpJson),
4493 &upstream(),
4494 &serde_json::json!({}),
4495 100,
4496 Some(1_000),
4497 );
4498
4499 assert_eq!(snapshot.status, BalanceSnapshotStatus::Unknown);
4500 assert_eq!(snapshot.exhausted, None);
4501 }
4502
4503 #[test]
4504 fn yescode_snapshot_sums_subscription_and_paygo_balances() {
4505 let snapshot = yescode_snapshot_from_json(
4506 &provider("yescode", ProviderKind::YescodeProfile),
4507 &upstream(),
4508 &serde_json::json!({
4509 "subscription_balance": "1.25",
4510 "pay_as_you_go_balance": 2.5
4511 }),
4512 100,
4513 Some(1_000),
4514 );
4515
4516 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
4517 assert_eq!(snapshot.exhausted, Some(false));
4518 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.75"));
4519 assert_eq!(snapshot.subscription_balance_usd.as_deref(), Some("1.25"));
4520 assert_eq!(snapshot.paygo_balance_usd.as_deref(), Some("2.5"));
4521 }
4522
4523 #[test]
4524 fn openai_balance_endpoint_defaults_to_base_user_balance_without_v1() {
4525 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
4526 provider.endpoint.clear();
4527
4528 let endpoint =
4529 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
4530
4531 assert_eq!(endpoint, "https://relay.example.com/user/balance");
4532 }
4533
4534 #[test]
4535 fn sub2api_usage_endpoint_defaults_to_upstream_usage_under_v1() {
4536 let mut provider = provider("sub2api", ProviderKind::Sub2ApiUsage);
4537 provider.endpoint.clear();
4538
4539 let endpoint =
4540 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
4541
4542 assert_eq!(endpoint, "https://relay.example.com/v1/usage");
4543 }
4544
4545 #[test]
4546 fn sub2api_auth_me_endpoint_defaults_to_dashboard_path_without_v1() {
4547 let mut provider = provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe);
4548 provider.endpoint.clear();
4549
4550 let endpoint =
4551 resolve_endpoint(&provider, "https://relay.example.com/v1", "token").expect("endpoint");
4552
4553 assert_eq!(endpoint, "https://relay.example.com/api/v1/auth/me");
4554 }
4555
4556 #[test]
4557 fn provider_templates_support_variables_for_custom_headers_or_queries() {
4558 let mut provider = provider("newapi", ProviderKind::NewApiUserSelf);
4559 provider.endpoint = "{{base_url}}/api/user/self?user={{userId}}".to_string();
4560 provider
4561 .variables
4562 .insert("userId".to_string(), "42".to_string());
4563
4564 let endpoint = resolve_endpoint(&provider, "https://newapi.example.com/v1", "token")
4565 .expect("endpoint");
4566
4567 assert_eq!(endpoint, "https://newapi.example.com/api/user/self?user=42");
4568 }
4569
4570 #[test]
4571 fn new_api_token_usage_endpoint_defaults_to_model_key_usage_path() {
4572 let mut provider = provider("newapi-token", ProviderKind::NewApiTokenUsage);
4573 provider.endpoint.clear();
4574
4575 let endpoint = resolve_endpoint(&provider, "https://newapi.example.com/v1", "token")
4576 .expect("endpoint");
4577
4578 assert_eq!(endpoint, "https://newapi.example.com/api/usage/token/");
4579 }
4580
4581 #[test]
4582 fn openai_organization_costs_endpoint_defaults_to_official_v1_costs_window() {
4583 let mut provider = provider("openai", ProviderKind::OpenAiOrganizationCosts);
4584 provider.endpoint.clear();
4585
4586 let endpoint =
4587 resolve_endpoint(&provider, "https://api.openai.com/v1", "token").expect("endpoint");
4588
4589 assert!(endpoint.starts_with("https://api.openai.com/v1/organization/costs?start_time="));
4590 assert!(endpoint.ends_with("&limit=30"));
4591 let start_time = endpoint
4592 .split("start_time=")
4593 .nth(1)
4594 .and_then(|value| value.split('&').next())
4595 .and_then(|value| value.parse::<u64>().ok())
4596 .expect("numeric start_time");
4597 assert!(start_time > 0);
4598 }
4599
4600 #[test]
4601 fn require_token_env_prevents_upstream_model_key_fallback() {
4602 let mut cfg = proxy_config(vec![service_config(
4603 "right",
4604 vec![upstream_config("https://api.openai.com/v1")],
4605 )]);
4606 cfg.codex
4607 .configs
4608 .get_mut("right")
4609 .expect("station")
4610 .upstreams[0]
4611 .auth
4612 .auth_token = Some("model-key".to_string());
4613
4614 let mut provider = provider("openai", ProviderKind::OpenAiOrganizationCosts);
4615 provider.token_env = Some("__CODEX_HELPER_TEST_MISSING_TOKEN_ENV__".to_string());
4616 provider.require_token_env = true;
4617 let upstreams = [UpstreamRef {
4618 station_name: "right".to_string(),
4619 index: 0,
4620 provider_endpoint: None,
4621 }];
4622
4623 assert_eq!(resolve_token(&provider, &upstreams, &cfg, "codex"), None);
4624
4625 provider.require_token_env = false;
4626 assert_eq!(
4627 resolve_token(&provider, &upstreams, &cfg, "codex").as_deref(),
4628 Some("model-key")
4629 );
4630 }
4631
4632 #[test]
4633 fn effective_poll_interval_respects_disable_flag_zero_and_minimum() {
4634 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
4635
4636 provider.poll_interval_secs = Some(0);
4637 assert_eq!(effective_poll_interval_secs(&provider), None);
4638
4639 provider.poll_interval_secs = Some(10);
4640 assert_eq!(
4641 effective_poll_interval_secs(&provider),
4642 Some(MIN_POLL_INTERVAL_SECS)
4643 );
4644
4645 provider.poll_interval_secs = None;
4646 assert_eq!(
4647 effective_poll_interval_secs(&provider),
4648 Some(DEFAULT_POLL_INTERVAL_SECS)
4649 );
4650
4651 provider.refresh_on_request = false;
4652 assert_eq!(effective_poll_interval_secs(&provider), None);
4653 }
4654
4655 #[test]
4656 fn usage_provider_http_error_detail_extracts_json_code_and_message() {
4657 let detail = usage_provider_http_error_detail(
4658 r#"{"code":"USER_INACTIVE","message":"User account is not active"}"#,
4659 )
4660 .expect("error detail");
4661
4662 assert_eq!(detail, "USER_INACTIVE: User account is not active");
4663 assert!(usage_provider_error_is_terminal(&detail));
4664 }
4665
4666 #[test]
4667 fn current_day_quota_exhaustion_blocks_followup_usage_even_when_display_only() {
4668 let mut snapshot = ProviderBalanceSnapshot {
4669 status: BalanceSnapshotStatus::Exhausted,
4670 exhausted: Some(true),
4671 exhaustion_affects_routing: false,
4672 quota_period: Some("daily".to_string()),
4673 quota_remaining_usd: Some("0".to_string()),
4674 ..ProviderBalanceSnapshot::default()
4675 };
4676 snapshot.refresh_status(0);
4677
4678 assert!(!snapshot.routing_exhausted());
4679 assert!(
4680 usage_provider_suppression_decision_from_snapshot(
4681 &snapshot,
4682 &ProxyConfig::default(),
4683 snapshot.fetched_at_ms,
4684 )
4685 .is_some()
4686 );
4687 assert_eq!(
4688 usage_provider_snapshot_suppression_reason(&snapshot).as_deref(),
4689 Some("daily package quota exhausted for current period")
4690 );
4691 }
4692
4693 #[test]
4694 fn current_day_quota_suppression_expires_after_configured_reset() {
4695 let cfg = ProxyConfig::default();
4696 let fetched_at_ms = 1_700_000_000_000;
4697 let reset_at_ms = next_reset_at_ms(
4698 fetched_at_ms,
4699 cfg.ui.usage_forecast.reset_utc_offset.as_str(),
4700 cfg.ui.usage_forecast.reset_time.as_str(),
4701 )
4702 .expect("default reset config is valid");
4703 let grace_ms = duration_millis_u64(USAGE_PROVIDER_DAILY_RESET_SUPPRESSION_GRACE);
4704 let suppress_until_ms = reset_at_ms + grace_ms;
4705
4706 let mut snapshot = ProviderBalanceSnapshot {
4707 fetched_at_ms,
4708 stale_after_ms: Some(fetched_at_ms + 60_000),
4709 status: BalanceSnapshotStatus::Exhausted,
4710 exhausted: Some(true),
4711 exhaustion_affects_routing: false,
4712 quota_period: Some("daily".to_string()),
4713 quota_remaining_usd: Some("0".to_string()),
4714 ..ProviderBalanceSnapshot::default()
4715 };
4716 snapshot.refresh_status(fetched_at_ms);
4717
4718 let decision = usage_provider_suppression_decision_from_snapshot(
4719 &snapshot,
4720 &cfg,
4721 suppress_until_ms - 1,
4722 )
4723 .expect("daily exhaustion should suppress until reset grace expires");
4724 assert_eq!(decision.ttl, Duration::from_millis(1));
4725 assert!(decision.routing_exhausted);
4726
4727 assert!(
4728 usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, suppress_until_ms)
4729 .is_none(),
4730 "a stale daily exhaustion snapshot must not suppress refresh after the reset boundary"
4731 );
4732 }
4733
4734 #[test]
4735 fn stale_non_daily_exhaustion_snapshot_does_not_renew_suppression() {
4736 let cfg = ProxyConfig::default();
4737 let mut snapshot = ProviderBalanceSnapshot {
4738 fetched_at_ms: 1_000,
4739 stale_after_ms: Some(2_000),
4740 status: BalanceSnapshotStatus::Exhausted,
4741 exhausted: Some(true),
4742 quota_period: Some("quota".to_string()),
4743 quota_remaining_usd: Some("0".to_string()),
4744 ..ProviderBalanceSnapshot::default()
4745 };
4746 snapshot.refresh_status(1_000);
4747
4748 let fresh_decision =
4749 usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, 1_500)
4750 .expect("fresh exhausted quota should suppress follow-up polling");
4751 assert_eq!(fresh_decision.ttl, USAGE_PROVIDER_EXHAUSTED_SUPPRESSION_TTL);
4752
4753 assert!(
4754 usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, 2_001).is_none(),
4755 "stale exhausted quota snapshots must not be used to renew suppression forever"
4756 );
4757 }
4758
4759 #[test]
4760 fn refreshable_weekly_window_exhaustion_suppresses_only_while_snapshot_is_fresh() {
4761 let cfg = ProxyConfig::default();
4762 let mut snapshot = ProviderBalanceSnapshot {
4763 fetched_at_ms: 1_000,
4764 stale_after_ms: Some(10_000),
4765 status: BalanceSnapshotStatus::Exhausted,
4766 exhausted: Some(true),
4767 exhaustion_affects_routing: false,
4768 quota_period: Some("weekly".to_string()),
4769 quota_remaining_usd: Some("0".to_string()),
4770 ..ProviderBalanceSnapshot::default()
4771 };
4772 snapshot.refresh_status(1_000);
4773
4774 let decision = usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, 4_000)
4775 .expect("fresh weekly window exhaustion should block follow-up usage");
4776 assert_eq!(decision.ttl, Duration::from_millis(6_000));
4777 assert!(decision.routing_exhausted);
4778
4779 assert!(
4780 usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, 10_001).is_none(),
4781 "weekly/monthly windows without explicit reset_at should be re-queried after staleness"
4782 );
4783 }
4784
4785 #[test]
4786 fn rate_limit_reset_at_drives_suppression_ttl() {
4787 let cfg = ProxyConfig::default();
4788 let reset_at_ms = 120_000;
4789 let suppress_until_ms =
4790 reset_at_ms + duration_millis_u64(USAGE_PROVIDER_DAILY_RESET_SUPPRESSION_GRACE);
4791 let mut snapshot = ProviderBalanceSnapshot {
4792 fetched_at_ms: 1_000,
4793 stale_after_ms: Some(10_000),
4794 status: BalanceSnapshotStatus::Exhausted,
4795 exhausted: Some(true),
4796 quota_period: Some("rate_limit:5h".to_string()),
4797 quota_resets_at_ms: Some(reset_at_ms),
4798 ..ProviderBalanceSnapshot::default()
4799 };
4800 snapshot.refresh_status(1_000);
4801
4802 let decision = usage_provider_suppression_decision_from_snapshot(
4803 &snapshot,
4804 &cfg,
4805 suppress_until_ms - 1,
4806 )
4807 .expect("rate limit should suppress until reset grace expires");
4808 assert_eq!(decision.ttl, Duration::from_millis(1));
4809
4810 assert!(
4811 usage_provider_suppression_decision_from_snapshot(&snapshot, &cfg, suppress_until_ms)
4812 .is_none()
4813 );
4814 }
4815
4816 #[test]
4817 fn auto_provider_uses_stable_target_id_across_probe_kinds() {
4818 let target = UsageProviderTarget {
4819 upstream: UpstreamRef {
4820 station_name: "input/sub".to_string(),
4821 index: 2,
4822 provider_endpoint: None,
4823 },
4824 base_url: "https://ai.input.im/v1".to_string(),
4825 provider_id: None,
4826 };
4827
4828 let sub2api = auto_usage_provider(&target, ProviderKind::Sub2ApiUsage);
4829 let newapi_token = auto_usage_provider(&target, ProviderKind::NewApiTokenUsage);
4830 let newapi = auto_usage_provider(&target, ProviderKind::NewApiUserSelf);
4831
4832 assert_eq!(sub2api.id, "auto:balance:input-sub:2");
4833 assert_eq!(sub2api.id, newapi_token.id);
4834 assert_eq!(sub2api.id, newapi.id);
4835 assert_eq!(sub2api.domains, vec!["ai.input.im".to_string()]);
4836 assert_eq!(
4837 resolve_endpoint(&sub2api, &target.base_url, "token").unwrap(),
4838 "https://ai.input.im/v1/usage"
4839 );
4840 assert_eq!(
4841 resolve_endpoint(&newapi_token, &target.base_url, "token").unwrap(),
4842 "https://ai.input.im/api/usage/token/"
4843 );
4844 }
4845
4846 #[test]
4847 fn auto_probe_prefers_rightcode_adapter_for_rightcode_hosts() {
4848 let target = UsageProviderTarget {
4849 upstream: UpstreamRef {
4850 station_name: "right".to_string(),
4851 index: 0,
4852 provider_endpoint: None,
4853 },
4854 base_url: "https://www.right.codes/codex/v1".to_string(),
4855 provider_id: Some("right".to_string()),
4856 };
4857
4858 assert_eq!(
4859 first_auto_probe_kind(&target),
4860 ProviderKind::RightCodeAccountSummary
4861 );
4862 assert_eq!(
4863 resolve_endpoint(
4864 &auto_usage_provider(&target, ProviderKind::RightCodeAccountSummary),
4865 &target.base_url,
4866 "token"
4867 )
4868 .unwrap(),
4869 "https://www.right.codes/account/summary"
4870 );
4871 assert_eq!(
4872 auto_usage_provider(&target, ProviderKind::RightCodeAccountSummary)
4873 .token_env
4874 .as_deref(),
4875 None
4876 );
4877 }
4878
4879 #[test]
4880 fn auto_provider_id_prefers_runtime_provider_tag() {
4881 let target = UsageProviderTarget {
4882 upstream: UpstreamRef {
4883 station_name: "routing".to_string(),
4884 index: 0,
4885 provider_endpoint: Some(ProviderEndpointKey::new("codex", "input", "default")),
4886 },
4887 base_url: "https://ai.input.im/v1".to_string(),
4888 provider_id: Some("input".to_string()),
4889 };
4890
4891 let provider = auto_usage_provider(&target, ProviderKind::Sub2ApiUsage);
4892
4893 assert_eq!(provider.id, "input");
4894 }
4895
4896 #[test]
4897 fn auto_target_provider_id_filter_matches_runtime_provider_tag() {
4898 let target = UsageProviderTarget {
4899 upstream: UpstreamRef {
4900 station_name: "routing".to_string(),
4901 index: 6,
4902 provider_endpoint: Some(ProviderEndpointKey::new("codex", "input6", "default")),
4903 },
4904 base_url: "https://input.9z1.me/v1".to_string(),
4905 provider_id: Some("input6".to_string()),
4906 };
4907
4908 assert!(auto_target_matches_provider_id_filter(&target, None));
4909 assert!(auto_target_matches_provider_id_filter(
4910 &target,
4911 Some("input6")
4912 ));
4913 assert!(!auto_target_matches_provider_id_filter(
4914 &target,
4915 Some("input5")
4916 ));
4917 }
4918
4919 #[test]
4920 fn auto_target_provider_id_filter_matches_generated_auto_id() {
4921 let target = UsageProviderTarget {
4922 upstream: UpstreamRef {
4923 station_name: "routing".to_string(),
4924 index: 6,
4925 provider_endpoint: None,
4926 },
4927 base_url: "https://input.9z1.me/v1".to_string(),
4928 provider_id: None,
4929 };
4930
4931 assert!(auto_target_matches_provider_id_filter(
4932 &target,
4933 Some("auto:balance:routing:6")
4934 ));
4935 assert!(!auto_target_matches_provider_id_filter(
4936 &target,
4937 Some("input6")
4938 ));
4939 }
4940
4941 #[test]
4942 fn provider_endpoint_target_lookup_uses_endpoint_identity() {
4943 let cfg = proxy_config(vec![service_config(
4944 "routing",
4945 vec![
4946 endpoint_upstream_config("https://input.example/v1", "input", "default"),
4947 endpoint_upstream_config("https://right.example/v1", "right", "default"),
4948 ],
4949 )]);
4950
4951 let target = usage_provider_target_for_provider_endpoint(
4952 &cfg,
4953 "codex",
4954 &ProviderEndpointKey::new("codex", "right", "default"),
4955 )
4956 .expect("provider endpoint target");
4957
4958 assert_eq!(target.upstream.station_name, "routing");
4959 assert_eq!(target.upstream.index, 1);
4960 assert_eq!(
4961 target
4962 .upstream
4963 .provider_endpoint
4964 .as_ref()
4965 .map(ProviderEndpointKey::stable_key)
4966 .as_deref(),
4967 Some("codex/right/default")
4968 );
4969 assert_eq!(target.base_url, "https://right.example/v1");
4970 assert_eq!(target.provider_id.as_deref(), Some("right"));
4971 }
4972
4973 #[test]
4974 fn request_balance_queue_deduplicates_until_due() {
4975 let key = ProviderEndpointKey::new("codex", "input", "default");
4976 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
4977 {
4978 let mut queue = queue.lock().expect("queue");
4979 queue.remove(&key);
4980 }
4981
4982 assert_eq!(
4983 enqueue_request_balance_refresh(key.clone()),
4984 Some(REQUEST_BALANCE_REFRESH_DELAY)
4985 );
4986 assert_eq!(enqueue_request_balance_refresh(key.clone()), None);
4987 assert!(matches!(
4988 take_request_balance_refresh_if_due(&key),
4989 RequestBalanceQueueDue::NotDue(_)
4990 ));
4991
4992 {
4993 let mut queue = queue.lock().expect("queue");
4994 queue.insert(key.clone(), Instant::now() - Duration::from_secs(1));
4995 }
4996
4997 assert_eq!(
4998 take_request_balance_refresh_if_due(&key),
4999 RequestBalanceQueueDue::Due
5000 );
5001 assert_eq!(
5002 take_request_balance_refresh_if_due(&key),
5003 RequestBalanceQueueDue::Missing
5004 );
5005 assert_eq!(
5006 enqueue_request_balance_refresh(key.clone()),
5007 Some(REQUEST_BALANCE_REFRESH_DELAY)
5008 );
5009
5010 queue.lock().expect("queue").remove(&key);
5011 }
5012
5013 #[test]
5014 fn request_balance_queue_does_not_extend_due_refresh() {
5015 let key = ProviderEndpointKey::new("codex", "input", "default");
5016 let queue = REQUEST_BALANCE_QUEUE.get_or_init(|| Mutex::new(HashMap::new()));
5017 {
5018 let mut queue = queue.lock().expect("queue");
5019 queue.insert(key.clone(), Instant::now() - Duration::from_secs(1));
5020 }
5021
5022 assert_eq!(
5023 enqueue_request_balance_refresh(key.clone()),
5024 Some(Duration::ZERO)
5025 );
5026 assert_eq!(
5027 take_request_balance_refresh_if_due(&key),
5028 RequestBalanceQueueDue::Due
5029 );
5030 assert_eq!(
5031 take_request_balance_refresh_if_due(&key),
5032 RequestBalanceQueueDue::Missing
5033 );
5034
5035 queue.lock().expect("queue").remove(&key);
5036 }
5037
5038 #[test]
5039 fn remaining_poll_cooldown_returns_only_unexpired_window() {
5040 let now = Instant::now();
5041 assert_eq!(
5042 remaining_poll_cooldown(now - Duration::from_secs(30), 60, now),
5043 Some(Duration::from_secs(30))
5044 );
5045 assert_eq!(
5046 remaining_poll_cooldown(now - Duration::from_secs(60), 60, now),
5047 None
5048 );
5049 assert_eq!(
5050 remaining_poll_cooldown(now - Duration::from_secs(61), 60, now),
5051 None
5052 );
5053 }
5054
5055 #[test]
5056 fn configured_target_keys_prevent_auto_probe_for_explicit_balance_domains() {
5057 let cfg = proxy_config(vec![
5058 service_config("explicit", vec![upstream_config("https://example.com/v1")]),
5059 service_config("auto", vec![upstream_config("https://ai.input.im/v1")]),
5060 ]);
5061 let configured = configured_target_keys(
5062 &cfg,
5063 "codex",
5064 &[provider("relay", ProviderKind::OpenAiBalanceHttpJson)],
5065 None,
5066 );
5067 let auto_targets = usage_provider_targets(&cfg, "codex", None)
5068 .into_iter()
5069 .filter(|target| !configured.contains(&target_key(target)))
5070 .map(|target| target.upstream.station_name)
5071 .collect::<Vec<_>>();
5072
5073 assert_eq!(auto_targets, vec!["auto".to_string()]);
5074 }
5075
5076 #[test]
5077 fn auto_probe_accepts_only_usable_balance_snapshots() {
5078 let usable = sub2api_usage_snapshot_from_json(
5079 &provider("auto", ProviderKind::Sub2ApiUsage),
5080 &upstream(),
5081 &serde_json::json!({
5082 "isValid": true,
5083 "remaining": 1
5084 }),
5085 100,
5086 Some(1_000),
5087 );
5088 let unusable = balance_http_snapshot_from_json(
5089 &provider("auto", ProviderKind::OpenAiBalanceHttpJson),
5090 &upstream(),
5091 &serde_json::json!({ "ok": true }),
5092 100,
5093 Some(1_000),
5094 );
5095
5096 assert!(auto_snapshot_is_usable(&usable));
5097 assert!(!auto_snapshot_is_usable(&unusable));
5098 }
5099
5100 #[test]
5101 fn auto_probe_error_summary_keeps_all_attempt_failures() {
5102 let errors = vec![
5103 "Sub2ApiUsage: HTTP 404".to_string(),
5104 "NewApiTokenUsage: missing quota fields".to_string(),
5105 "OpenAiBalanceHttpJson: non-JSON response".to_string(),
5106 ];
5107
5108 let summary = auto_probe_error_summary(&errors).expect("summary");
5109
5110 assert!(summary.contains("Sub2ApiUsage: HTTP 404"));
5111 assert!(summary.contains("NewApiTokenUsage: missing quota fields"));
5112 assert!(summary.contains("OpenAiBalanceHttpJson: non-JSON response"));
5113 }
5114
5115 #[test]
5116 fn auto_probe_kind_order_prioritizes_remembered_success() {
5117 let provider_id = "input-order-success";
5118 clear_auto_probe_kind_state(provider_id);
5119 let target = usage_provider_target("https://relay.example.com/v1", provider_id);
5120
5121 remember_auto_probe_kind_success(provider_id, &target, ProviderKind::NewApiUserSelf);
5122
5123 let order = auto_probe_kind_order(provider_id, &target);
5124
5125 assert_eq!(order.first(), Some(&ProviderKind::NewApiUserSelf));
5126 assert_eq!(
5127 order
5128 .iter()
5129 .filter(|kind| **kind == ProviderKind::Sub2ApiUsage)
5130 .count(),
5131 1
5132 );
5133 clear_auto_probe_kind_state(provider_id);
5134 }
5135
5136 #[test]
5137 fn auto_probe_kind_order_temporarily_skips_recent_failures() {
5138 let provider_id = "input-order-failure";
5139 clear_auto_probe_kind_state(provider_id);
5140 let target = usage_provider_target("https://relay.example.com/v1", provider_id);
5141 let now = Instant::now();
5142
5143 remember_auto_probe_kind_failure(provider_id, &target, ProviderKind::Sub2ApiUsage, now);
5144
5145 let order = auto_probe_kind_order(provider_id, &target);
5146
5147 assert!(!order.contains(&ProviderKind::Sub2ApiUsage));
5148 assert!(order.contains(&ProviderKind::NewApiTokenUsage));
5149
5150 if let Some(failures) = AUTO_PROBE_KIND_FAILURES.get()
5151 && let Ok(mut failures) = failures.lock()
5152 {
5153 failures.insert(
5154 AutoProbeKindFailureKey {
5155 provider_id: provider_id.to_string(),
5156 target: auto_probe_target_key(&target),
5157 kind: ProviderKind::Sub2ApiUsage,
5158 },
5159 now - AUTO_PROBE_KIND_FAILURE_TTL - Duration::from_secs(1),
5160 );
5161 }
5162
5163 let order = auto_probe_kind_order(provider_id, &target);
5164 assert!(order.contains(&ProviderKind::Sub2ApiUsage));
5165 clear_auto_probe_kind_state(provider_id);
5166 }
5167
5168 #[test]
5169 fn auto_probe_kind_order_can_be_empty_when_all_kinds_are_suppressed() {
5170 let provider_id = "input-order-all-suppressed";
5171 clear_auto_probe_kind_state(provider_id);
5172 let target = usage_provider_target("https://relay.example.com/v1", provider_id);
5173 let now = Instant::now();
5174
5175 for kind in AUTO_PROBE_KINDS {
5176 if kind == ProviderKind::RightCodeAccountSummary {
5177 continue;
5178 }
5179 remember_auto_probe_kind_failure(provider_id, &target, kind, now);
5180 }
5181
5182 assert!(auto_probe_kind_order(provider_id, &target).is_empty());
5183 clear_auto_probe_kind_state(provider_id);
5184 }
5185
5186 #[test]
5187 fn auto_probe_kind_failures_do_not_suppress_distinct_targets_with_same_provider_id() {
5188 let provider_id = "input-shared-provider";
5189 clear_auto_probe_kind_state(provider_id);
5190 let routing_target =
5191 usage_provider_target_at("routing", 0, "https://relay.example.com/v1", provider_id);
5192 let catalog_target =
5193 usage_provider_target_at("input", 0, "https://relay.example.com/v1", provider_id);
5194 let now = Instant::now();
5195
5196 for kind in AUTO_PROBE_KINDS {
5197 if kind == ProviderKind::RightCodeAccountSummary {
5198 continue;
5199 }
5200 remember_auto_probe_kind_failure(provider_id, &routing_target, kind, now);
5201 }
5202
5203 assert!(auto_probe_kind_order(provider_id, &routing_target).is_empty());
5204 assert!(
5205 !auto_probe_kind_order(provider_id, &catalog_target).is_empty(),
5206 "a routing target's temporary failures must not hide catalog balance probes for the same provider"
5207 );
5208 clear_auto_probe_kind_state(provider_id);
5209 }
5210
5211 #[tokio::test]
5212 async fn auto_probe_suppressed_order_records_error_snapshot() {
5213 let provider_id = "input-suppressed-snapshot";
5214 clear_auto_probe_kind_state(provider_id);
5215 let mut upstream =
5216 endpoint_upstream_config("https://relay.example.com/v1", provider_id, "default");
5217 upstream.auth.auth_token = Some("model-key".to_string());
5218 let cfg = proxy_config(vec![service_config("routing", vec![upstream])]);
5219 let lb_states = Arc::new(Mutex::new(HashMap::new()));
5220 let state = ProxyState::new();
5221 let target = usage_provider_target("https://relay.example.com/v1", provider_id);
5222 let now = Instant::now();
5223
5224 for kind in AUTO_PROBE_KINDS {
5225 if kind == ProviderKind::RightCodeAccountSummary {
5226 continue;
5227 }
5228 remember_auto_probe_kind_failure(provider_id, &target, kind, now);
5229 }
5230
5231 let outcome =
5232 auto_probe_provider_target(&Client::new(), &target, &cfg, &lb_states, &state, "codex")
5233 .await;
5234
5235 assert_eq!(outcome, UsageProviderRefreshOutcome::Failed);
5236 let view = state.get_provider_balance_view("codex").await;
5237 let snapshot = view
5238 .get("routing")
5239 .and_then(|snapshots| {
5240 snapshots
5241 .iter()
5242 .find(|snapshot| snapshot.provider_id == provider_id)
5243 })
5244 .expect("suppressed auto probe snapshot");
5245 assert_eq!(snapshot.status, BalanceSnapshotStatus::Error);
5246 assert_eq!(
5247 snapshot.error.as_deref(),
5248 Some("all balance probe kinds are temporarily suppressed")
5249 );
5250 assert_eq!(snapshot.upstream_index, Some(0));
5251 let guard = lb_states.lock().expect("lb states");
5252 assert!(
5253 !guard
5254 .get("routing")
5255 .and_then(|entry| entry.usage_exhausted.first())
5256 .copied()
5257 .unwrap_or(true)
5258 );
5259 clear_auto_probe_kind_state(provider_id);
5260 }
5261
5262 #[tokio::test]
5263 async fn auto_probe_terminal_auth_failure_keeps_route_exhausted_during_suppression() {
5264 let provider_id = "input-terminal-auth";
5265 clear_auto_probe_kind_state(provider_id);
5266 let request_count = Arc::new(AtomicUsize::new(0));
5267 let counter = request_count.clone();
5268 let app = axum::Router::new().fallback(get(move || {
5269 let counter = counter.clone();
5270 async move {
5271 counter.fetch_add(1, Ordering::SeqCst);
5272 (
5273 axum::http::StatusCode::UNAUTHORIZED,
5274 axum::Json(serde_json::json!({
5275 "code": "USER_INACTIVE",
5276 "message": "User account is not active"
5277 })),
5278 )
5279 }
5280 }));
5281 let (addr, handle) = spawn_axum_server(app).await;
5282 let base_url = format!("http://{addr}/v1");
5283 let mut upstream = endpoint_upstream_config(&base_url, provider_id, "default");
5284 upstream.auth.auth_token = Some("model-key".to_string());
5285 let cfg = proxy_config(vec![service_config("routing", vec![upstream])]);
5286 let lb_states = Arc::new(Mutex::new(HashMap::new()));
5287 let state = ProxyState::new();
5288 let target = usage_provider_target(&base_url, provider_id);
5289
5290 let outcome =
5291 auto_probe_provider_target(&Client::new(), &target, &cfg, &lb_states, &state, "codex")
5292 .await;
5293
5294 assert_eq!(outcome, UsageProviderRefreshOutcome::Failed);
5295 assert!(request_count.load(Ordering::SeqCst) > 0);
5296 {
5297 let guard = lb_states.lock().expect("lb states");
5298 assert!(
5299 guard
5300 .get("routing")
5301 .and_then(|entry| entry.usage_exhausted.first())
5302 .copied()
5303 .unwrap_or(false)
5304 );
5305 }
5306 let view = state.get_provider_balance_view("codex").await;
5307 let snapshot = view
5308 .get("routing")
5309 .and_then(|snapshots| {
5310 snapshots
5311 .iter()
5312 .find(|snapshot| snapshot.provider_id == provider_id)
5313 })
5314 .expect("terminal auth failure snapshot");
5315 assert!(
5316 snapshot
5317 .error
5318 .as_deref()
5319 .unwrap_or("")
5320 .contains("User account is not active")
5321 );
5322
5323 let requests_after_first_probe = request_count.load(Ordering::SeqCst);
5324 let suppressed_outcome =
5325 auto_probe_provider_target(&Client::new(), &target, &cfg, &lb_states, &state, "codex")
5326 .await;
5327
5328 assert_eq!(suppressed_outcome, UsageProviderRefreshOutcome::Failed);
5329 assert_eq!(
5330 request_count.load(Ordering::SeqCst),
5331 requests_after_first_probe
5332 );
5333 let guard = lb_states.lock().expect("lb states");
5334 assert!(
5335 guard
5336 .get("routing")
5337 .and_then(|entry| entry.usage_exhausted.first())
5338 .copied()
5339 .unwrap_or(false)
5340 );
5341 clear_auto_probe_kind_state(provider_id);
5342 handle.abort();
5343 }
5344
5345 #[tokio::test]
5346 async fn auto_probe_daily_package_exhaustion_suppresses_followup_refresh() {
5347 let provider_id = "input-daily-exhausted";
5348 clear_auto_probe_kind_state(provider_id);
5349 let request_count = Arc::new(AtomicUsize::new(0));
5350 let counter = request_count.clone();
5351 let app = axum::Router::new().fallback(get(move || {
5352 let counter = counter.clone();
5353 async move {
5354 counter.fetch_add(1, Ordering::SeqCst);
5355 axum::Json(serde_json::json!({
5356 "isValid": true,
5357 "mode": "unrestricted",
5358 "planName": "CodeX Lite",
5359 "subscription": {
5360 "daily_usage_usd": 100,
5361 "daily_limit_usd": 100,
5362 "weekly_usage_usd": 100,
5363 "weekly_limit_usd": 0,
5364 "monthly_usage_usd": 100,
5365 "monthly_limit_usd": 0
5366 }
5367 }))
5368 }
5369 }));
5370 let (addr, handle) = spawn_axum_server(app).await;
5371 let base_url = format!("http://{addr}/v1");
5372 let mut upstream = endpoint_upstream_config(&base_url, provider_id, "default");
5373 upstream.auth.auth_token = Some("model-key".to_string());
5374 let cfg = proxy_config(vec![service_config("routing", vec![upstream])]);
5375 let lb_states = Arc::new(Mutex::new(HashMap::new()));
5376 let state = ProxyState::new();
5377 let target = usage_provider_target(&base_url, provider_id);
5378
5379 let outcome =
5380 auto_probe_provider_target(&Client::new(), &target, &cfg, &lb_states, &state, "codex")
5381 .await;
5382
5383 assert_eq!(outcome, UsageProviderRefreshOutcome::Refreshed);
5384 assert_eq!(request_count.load(Ordering::SeqCst), 1);
5385 {
5386 let guard = lb_states.lock().expect("lb states");
5387 assert!(
5388 guard
5389 .get("routing")
5390 .and_then(|entry| entry.usage_exhausted.first())
5391 .copied()
5392 .unwrap_or(false)
5393 );
5394 }
5395 let view = state.get_provider_balance_view("codex").await;
5396 let snapshot = view
5397 .get("routing")
5398 .and_then(|snapshots| {
5399 snapshots
5400 .iter()
5401 .find(|snapshot| snapshot.provider_id == provider_id)
5402 })
5403 .expect("daily exhausted snapshot");
5404 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
5405 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
5406 assert!(snapshot.routing_ignored_exhaustion());
5407
5408 let suppressed_outcome =
5409 auto_probe_provider_target(&Client::new(), &target, &cfg, &lb_states, &state, "codex")
5410 .await;
5411
5412 assert_eq!(suppressed_outcome, UsageProviderRefreshOutcome::Failed);
5413 assert_eq!(request_count.load(Ordering::SeqCst), 1);
5414 let guard = lb_states.lock().expect("lb states");
5415 assert!(
5416 guard
5417 .get("routing")
5418 .and_then(|entry| entry.usage_exhausted.first())
5419 .copied()
5420 .unwrap_or(false)
5421 );
5422 clear_auto_probe_kind_state(provider_id);
5423 handle.abort();
5424 }
5425
5426 #[test]
5427 fn openai_balance_snapshot_reads_common_sub2api_balance_shape() {
5428 let snapshot = balance_http_snapshot_from_json(
5429 &provider("sub2api", ProviderKind::OpenAiBalanceHttpJson),
5430 &upstream(),
5431 &serde_json::json!({
5432 "balance": "1.25"
5433 }),
5434 100,
5435 Some(1_000),
5436 );
5437
5438 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5439 assert_eq!(snapshot.exhausted, Some(false));
5440 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("1.25"));
5441 }
5442
5443 #[test]
5444 fn json_path_supports_array_indices_for_official_balance_shapes() {
5445 let value = serde_json::json!({
5446 "balance_infos": [
5447 { "currency": "CNY", "total_balance": "3.25" }
5448 ]
5449 });
5450
5451 assert_eq!(
5452 json_value_at_path(&value, "balance_infos.0.total_balance")
5453 .and_then(|value| value.as_str()),
5454 Some("3.25")
5455 );
5456 }
5457
5458 #[test]
5459 fn openai_balance_snapshot_reads_cc_switch_official_balance_shapes() {
5460 let snapshot = balance_http_snapshot_from_json(
5461 &provider("deepseek", ProviderKind::OpenAiBalanceHttpJson),
5462 &upstream(),
5463 &serde_json::json!({
5464 "balance_infos": [
5465 { "currency": "CNY", "total_balance": "3.25" }
5466 ],
5467 "is_available": true
5468 }),
5469 100,
5470 Some(1_000),
5471 );
5472
5473 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5474 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.25"));
5475
5476 let snapshot = balance_http_snapshot_from_json(
5477 &provider("siliconflow", ProviderKind::OpenAiBalanceHttpJson),
5478 &upstream(),
5479 &serde_json::json!({
5480 "code": 20000,
5481 "data": {
5482 "totalBalance": "8.5",
5483 "chargeBalance": "2.5"
5484 }
5485 }),
5486 100,
5487 Some(1_000),
5488 );
5489
5490 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5491 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("8.5"));
5492 assert_eq!(snapshot.paygo_balance_usd.as_deref(), Some("2.5"));
5493 }
5494
5495 #[test]
5496 fn openai_balance_snapshot_can_derive_remaining_from_total_and_used() {
5497 let mut provider = provider("openrouter", ProviderKind::OpenAiBalanceHttpJson);
5498 provider.extract.monthly_budget_paths = vec!["data.total_credits".to_string()];
5499 provider.extract.monthly_spent_paths = vec!["data.total_usage".to_string()];
5500 provider.extract.derive_remaining_from_budget_and_spent = true;
5501
5502 let snapshot = balance_http_snapshot_from_json(
5503 &provider,
5504 &upstream(),
5505 &serde_json::json!({
5506 "data": {
5507 "total_credits": "10",
5508 "total_usage": "4"
5509 }
5510 }),
5511 100,
5512 Some(1_000),
5513 );
5514
5515 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5516 assert_eq!(snapshot.exhausted, Some(false));
5517 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("6"));
5518 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10"));
5519 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("4"));
5520 }
5521
5522 #[test]
5523 fn openai_balance_snapshot_supports_divisor_for_minor_units() {
5524 let mut provider = provider("novita", ProviderKind::OpenAiBalanceHttpJson);
5525 provider.extract.remaining_balance_paths = vec!["availableBalance".to_string()];
5526 provider.extract.remaining_divisor = Some(10_000);
5527
5528 let snapshot = balance_http_snapshot_from_json(
5529 &provider,
5530 &upstream(),
5531 &serde_json::json!({
5532 "availableBalance": 12345
5533 }),
5534 100,
5535 Some(1_000),
5536 );
5537
5538 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5539 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("1.2345"));
5540 }
5541
5542 #[test]
5543 fn sub2api_usage_snapshot_reads_all_api_hub_usage_shape() {
5544 let snapshot = sub2api_usage_snapshot_from_json(
5545 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5546 &upstream(),
5547 &serde_json::json!({
5548 "isValid": true,
5549 "mode": "unrestricted",
5550 "planName": "CodeX Air",
5551 "remaining": 165.0877165,
5552 "usage": {
5553 "today": {
5554 "cost": 0,
5555 "requests": 0,
5556 "total_tokens": 0
5557 },
5558 "total": {
5559 "cost": 354.194748,
5560 "requests": 2691,
5561 "total_tokens": 384084697
5562 }
5563 }
5564 }),
5565 100,
5566 Some(1_000),
5567 );
5568
5569 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5570 assert_eq!(snapshot.exhausted, Some(false));
5571 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Air"));
5572 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("165.0877165"));
5573 assert_eq!(snapshot.total_used_usd.as_deref(), Some("354.194748"));
5574 assert_eq!(snapshot.today_used_usd.as_deref(), Some("0"));
5575 assert_eq!(snapshot.total_requests, Some(2691));
5576 assert_eq!(snapshot.today_requests, Some(0));
5577 assert_eq!(snapshot.total_tokens, Some(384084697));
5578 assert_eq!(snapshot.today_tokens, Some(0));
5579 }
5580
5581 #[test]
5582 fn sub2api_usage_snapshot_reads_rates_model_stats_windows_and_alerts() {
5583 let snapshot = sub2api_usage_snapshot_from_json(
5584 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5585 &upstream(),
5586 &serde_json::json!({
5587 "isValid": true,
5588 "mode": "unrestricted",
5589 "plan_name": "CodeX Pro",
5590 "remaining": 9,
5591 "subscription": {
5592 "daily_usage_usd": 95,
5593 "daily_limit_usd": 100,
5594 "weekly_usage_usd": "120.5",
5595 "weekly_limit_usd": 0,
5596 "monthly_usage_usd": 300.25,
5597 "monthly_limit_usd": 1000,
5598 "expires_at": "2026-05-09T12:00:00.000Z"
5599 },
5600 "usage": {
5601 "today": {
5602 "request_count": "7",
5603 "input_tokens": 100,
5604 "output_tokens": 25,
5605 "total_cost_usd": "1.5"
5606 },
5607 "total": {
5608 "requests": 42,
5609 "tokens": 1234,
5610 "cost": 9.25
5611 },
5612 "average_duration_ms": "842.7",
5613 "rpm": "0.7",
5614 "tpm": 85.3
5615 },
5616 "model_stats": [
5617 {
5618 "model": "gpt-4o-mini",
5619 "request_count": "7",
5620 "prompt_tokens": 100,
5621 "completion_tokens": 25,
5622 "input_cost": "0.12",
5623 "output_cost": "0.34"
5624 }
5625 ]
5626 }),
5627 100,
5628 Some(1_000),
5629 );
5630
5631 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5632 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Pro"));
5633 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("9"));
5634 assert_eq!(snapshot.today_requests, Some(7));
5635 assert_eq!(snapshot.today_tokens, Some(100));
5636 assert_eq!(snapshot.today_used_usd.as_deref(), Some("1.5"));
5637 assert_eq!(snapshot.total_requests, Some(42));
5638 assert_eq!(snapshot.total_tokens, Some(1234));
5639 assert_eq!(snapshot.total_used_usd.as_deref(), Some("9.25"));
5640 let rate = snapshot.usage_rate.expect("rate");
5641 assert_eq!(rate.average_duration_ms.as_deref(), Some("842.7"));
5642 assert_eq!(rate.rpm.as_deref(), Some("0.7"));
5643 assert_eq!(rate.tpm.as_deref(), Some("85.3"));
5644 assert_eq!(snapshot.usage_windows.len(), 3);
5645 assert_eq!(snapshot.usage_windows[0].period, "daily");
5646 assert_eq!(
5647 snapshot.usage_windows[0].remaining_usd.as_deref(),
5648 Some("5")
5649 );
5650 assert_eq!(snapshot.usage_windows[1].unlimited, Some(true));
5651 assert_eq!(snapshot.usage_model_stats.len(), 1);
5652 assert_eq!(snapshot.usage_model_stats[0].model, "gpt-4o-mini");
5653 assert_eq!(snapshot.usage_model_stats[0].request_count, Some(7));
5654 assert_eq!(snapshot.usage_model_stats[0].total_tokens, Some(125));
5655 assert_eq!(
5656 snapshot.usage_model_stats[0].total_cost_usd.as_deref(),
5657 Some("0.46")
5658 );
5659 assert_eq!(
5660 snapshot
5661 .usage_alerts
5662 .iter()
5663 .map(|alert| alert.kind)
5664 .collect::<Vec<_>>(),
5665 vec![
5666 ProviderUsageAlertKind::DailyUsage95,
5667 ProviderUsageAlertKind::LowBalance,
5668 ProviderUsageAlertKind::SubscriptionExpired,
5669 ]
5670 );
5671 }
5672
5673 #[test]
5674 fn sub2api_subscription_lazy_daily_reset_projects_today_capacity() {
5675 let snapshot = sub2api_usage_snapshot_from_json(
5676 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5677 &upstream(),
5678 &serde_json::json!({
5679 "isValid": true,
5680 "mode": "unrestricted",
5681 "planName": "CodeX Lite 年度",
5682 "remaining": 0,
5683 "subscription": {
5684 "daily_usage_usd": 100.468025,
5685 "daily_limit_usd": 100,
5686 "weekly_usage_usd": 401.441684,
5687 "weekly_limit_usd": 0,
5688 "monthly_usage_usd": 401.441684,
5689 "monthly_limit_usd": 0
5690 },
5691 "usage": {
5692 "today": { "cost": 0, "requests": 0, "total_tokens": 0 },
5693 "total": { "cost": 702.492098, "requests": 42, "total_tokens": 1234 }
5694 }
5695 }),
5696 100,
5697 Some(1_000),
5698 );
5699
5700 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5701 assert_eq!(snapshot.exhausted, Some(false));
5702 assert_eq!(snapshot.plan_name.as_deref(), Some("CodeX Lite 年度"));
5703 assert_eq!(snapshot.total_balance_usd, None);
5704 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
5705 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("100"));
5706 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("100"));
5707 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("0"));
5708 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("100"));
5709 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0"));
5710 assert_eq!(snapshot.total_used_usd.as_deref(), Some("702.492098"));
5711 assert_eq!(snapshot.today_used_usd.as_deref(), Some("0"));
5712 assert_eq!(snapshot.today_requests, Some(0));
5713 assert_eq!(snapshot.today_tokens, Some(0));
5714 assert_eq!(snapshot.usage_windows[0].period, "daily");
5715 assert_eq!(snapshot.usage_windows[0].used_usd.as_deref(), Some("0"));
5716 assert_eq!(
5717 snapshot.usage_windows[0].remaining_usd.as_deref(),
5718 Some("100")
5719 );
5720 assert!(
5721 !snapshot
5722 .usage_alerts
5723 .iter()
5724 .any(|alert| alert.kind == ProviderUsageAlertKind::DailyUsage95)
5725 );
5726 assert!(
5727 !snapshot.routing_exhausted(),
5728 "sub2api /v1/usage skips billing checks; subscription windows are reset lazily on real requests"
5729 );
5730 }
5731
5732 #[test]
5733 fn sub2api_subscription_same_day_daily_exhaustion_remains_exhausted() {
5734 let snapshot = sub2api_usage_snapshot_from_json(
5735 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5736 &upstream(),
5737 &serde_json::json!({
5738 "isValid": true,
5739 "mode": "unrestricted",
5740 "planName": "CodeX Lite 年度",
5741 "remaining": 0,
5742 "subscription": {
5743 "daily_usage_usd": 100.468025,
5744 "daily_limit_usd": 100,
5745 "weekly_usage_usd": 401.441684,
5746 "weekly_limit_usd": 0,
5747 "monthly_usage_usd": 401.441684,
5748 "monthly_limit_usd": 0
5749 },
5750 "usage": {
5751 "today": { "cost": 100.468025, "requests": 8, "total_tokens": 1234 },
5752 "total": { "cost": 702.492098, "requests": 42, "total_tokens": 1234 }
5753 }
5754 }),
5755 100,
5756 Some(1_000),
5757 );
5758
5759 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
5760 assert_eq!(snapshot.exhausted, Some(true));
5761 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
5762 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
5763 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("100.468025"));
5764 assert_eq!(snapshot.today_used_usd.as_deref(), Some("100.468025"));
5765 assert!(!snapshot.routing_exhausted());
5766 }
5767
5768 #[test]
5769 fn sub2api_quota_limited_zero_remaining_still_marks_exhausted() {
5770 let snapshot = sub2api_usage_snapshot_from_json(
5771 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5772 &upstream(),
5773 &serde_json::json!({
5774 "isValid": true,
5775 "mode": "quota_limited",
5776 "quota": {
5777 "limit": 10,
5778 "used": 10,
5779 "remaining": 0,
5780 "unit": "USD"
5781 }
5782 }),
5783 100,
5784 Some(1_000),
5785 );
5786
5787 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
5788 assert_eq!(snapshot.exhausted, Some(true));
5789 assert_eq!(snapshot.total_balance_usd, None);
5790 assert_eq!(snapshot.quota_period.as_deref(), Some("quota"));
5791 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
5792 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("10"));
5793 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("10"));
5794 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("10"));
5795 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("10"));
5796 assert!(snapshot.routing_exhausted());
5797 }
5798
5799 #[test]
5800 fn sub2api_quota_limited_rate_limit_exhaustion_marks_temporary_window() {
5801 let reset_at = "2026-01-02T03:04:05Z";
5802 let reset_at_ms = parse_timestamp_secs(reset_at).expect("timestamp") * 1000;
5803 let snapshot = sub2api_usage_snapshot_from_json(
5804 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5805 &upstream(),
5806 &serde_json::json!({
5807 "isValid": true,
5808 "mode": "quota_limited",
5809 "rate_limits": [
5810 {
5811 "window": "5h",
5812 "limit": 100,
5813 "used": 100,
5814 "remaining": 0,
5815 "reset_at": reset_at
5816 }
5817 ]
5818 }),
5819 100,
5820 Some(1_000),
5821 );
5822
5823 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
5824 assert_eq!(snapshot.exhausted, Some(true));
5825 assert_eq!(snapshot.quota_period.as_deref(), Some("rate_limit:5h"));
5826 assert_eq!(snapshot.quota_resets_at_ms, Some(reset_at_ms));
5827 assert_eq!(snapshot.quota_remaining_usd, None);
5828 assert!(snapshot.routing_exhausted());
5829 }
5830
5831 #[test]
5832 fn sub2api_quota_limited_total_quota_exhaustion_wins_over_rate_limit_window() {
5833 let snapshot = sub2api_usage_snapshot_from_json(
5834 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5835 &upstream(),
5836 &serde_json::json!({
5837 "isValid": true,
5838 "mode": "quota_limited",
5839 "quota": {
5840 "limit": 10,
5841 "used": 10,
5842 "remaining": 0,
5843 "unit": "USD"
5844 },
5845 "rate_limits": [
5846 {
5847 "window": "5h",
5848 "limit": 100,
5849 "used": 100,
5850 "remaining": 0,
5851 "reset_at": "2026-01-02T03:04:05Z"
5852 }
5853 ]
5854 }),
5855 100,
5856 Some(1_000),
5857 );
5858
5859 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
5860 assert_eq!(snapshot.quota_period.as_deref(), Some("quota"));
5861 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
5862 assert_eq!(snapshot.quota_resets_at_ms, None);
5863 }
5864
5865 #[test]
5866 fn sub2api_usage_snapshot_marks_invalid_key_as_error() {
5867 let snapshot = sub2api_usage_snapshot_from_json(
5868 &provider("sub2api", ProviderKind::Sub2ApiUsage),
5869 &upstream(),
5870 &serde_json::json!({
5871 "isValid": false
5872 }),
5873 100,
5874 Some(1_000),
5875 );
5876
5877 assert_eq!(snapshot.status, BalanceSnapshotStatus::Error);
5878 assert_eq!(
5879 snapshot.error.as_deref(),
5880 Some("sub2api usage response reported invalid API key")
5881 );
5882 }
5883
5884 #[test]
5885 fn sub2api_auth_me_snapshot_reads_dashboard_balance_envelope() {
5886 let snapshot = sub2api_auth_me_snapshot_from_json(
5887 &provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe),
5888 &upstream(),
5889 &serde_json::json!({
5890 "code": 0,
5891 "message": "ok",
5892 "data": {
5893 "id": 42,
5894 "username": "demo",
5895 "balance": "12.5"
5896 }
5897 }),
5898 100,
5899 Some(1_000),
5900 );
5901
5902 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5903 assert_eq!(snapshot.exhausted, Some(false));
5904 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("12.5"));
5905 }
5906
5907 #[test]
5908 fn rightcode_endpoint_defaults_to_account_summary() {
5909 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
5910 provider.endpoint.clear();
5911
5912 let endpoint = resolve_endpoint(&provider, "https://www.right.codes/codex/v1", "token")
5913 .expect("endpoint");
5914
5915 assert_eq!(endpoint, "https://www.right.codes/account/summary");
5916 }
5917
5918 #[test]
5919 fn rightcode_account_summary_reads_matching_subscription_and_balance() {
5920 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
5921 provider.trust_exhaustion_for_routing = false;
5922
5923 let snapshot = rightcode_account_summary_snapshot_from_json(
5924 &provider,
5925 &upstream(),
5926 &serde_json::json!({
5927 "balance": 3.25,
5928 "subscriptions": [
5929 {
5930 "name": "Daily",
5931 "total_quota": 20,
5932 "remaining_quota": 7.5,
5933 "reset_today": true,
5934 "available_prefixes": ["/codex"]
5935 },
5936 {
5937 "name": "Other",
5938 "total_quota": 99,
5939 "remaining_quota": 99,
5940 "reset_today": true,
5941 "available_prefixes": ["/claude"]
5942 },
5943 {
5944 "name": "Badge",
5945 "total_quota": 5,
5946 "remaining_quota": 5,
5947 "reset_today": true,
5948 "available_prefixes": ["/codex"]
5949 }
5950 ]
5951 }),
5952 "https://www.right.codes/codex/v1",
5953 100,
5954 Some(1_000),
5955 );
5956
5957 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5958 assert_eq!(snapshot.exhausted, Some(false));
5959 assert!(!snapshot.routing_exhausted());
5960 assert_eq!(snapshot.plan_name.as_deref(), Some("Daily"));
5961 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("3.25"));
5962 assert_eq!(snapshot.paygo_balance_usd, None);
5963 assert_eq!(snapshot.quota_period.as_deref(), Some("daily"));
5964 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("7.5"));
5965 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("20"));
5966 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("12.5"));
5967 }
5968
5969 #[test]
5970 fn rightcode_account_summary_accounts_for_not_reset_today() {
5971 let provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
5972
5973 let snapshot = rightcode_account_summary_snapshot_from_json(
5974 &provider,
5975 &upstream(),
5976 &serde_json::json!({
5977 "subscriptions": [
5978 {
5979 "name": "Daily",
5980 "total_quota": 20,
5981 "remaining_quota": 7.5,
5982 "reset_today": false,
5983 "available_prefixes": ["/codex"]
5984 }
5985 ]
5986 }),
5987 "https://www.right.codes/codex/v1",
5988 100,
5989 Some(1_000),
5990 );
5991
5992 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
5993 assert_eq!(snapshot.exhausted, Some(false));
5994 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("27.5"));
5995 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("20"));
5996 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("0"));
5997 }
5998
5999 #[test]
6000 fn rightcode_zero_daily_quota_without_balance_is_display_only_exhaustion_by_default() {
6001 let mut provider = provider("rightcode", ProviderKind::RightCodeAccountSummary);
6002 provider.trust_exhaustion_for_routing = false;
6003
6004 let snapshot = rightcode_account_summary_snapshot_from_json(
6005 &provider,
6006 &upstream(),
6007 &serde_json::json!({
6008 "balance": 0,
6009 "subscriptions": [
6010 {
6011 "name": "Daily",
6012 "total_quota": 20,
6013 "remaining_quota": 0,
6014 "reset_today": true,
6015 "available_prefixes": ["/codex"]
6016 }
6017 ]
6018 }),
6019 "https://www.right.codes/codex/v1",
6020 100,
6021 Some(1_000),
6022 );
6023
6024 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
6025 assert_eq!(snapshot.exhausted, Some(true));
6026 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("0"));
6027 assert!(!snapshot.routing_exhausted());
6028 assert!(snapshot.routing_ignored_exhaustion());
6029 }
6030
6031 #[test]
6032 fn sub2api_auth_me_snapshot_marks_business_error() {
6033 let snapshot = sub2api_auth_me_snapshot_from_json(
6034 &provider("sub2api-auth", ProviderKind::Sub2ApiAuthMe),
6035 &upstream(),
6036 &serde_json::json!({
6037 "code": 401,
6038 "message": "login required"
6039 }),
6040 100,
6041 Some(1_000),
6042 );
6043
6044 assert_eq!(snapshot.status, BalanceSnapshotStatus::Error);
6045 assert_eq!(snapshot.error.as_deref(), Some("login required"));
6046 }
6047
6048 #[test]
6049 fn provider_can_disable_routing_trust_for_exhausted_balance() {
6050 let mut provider = provider("sub2api", ProviderKind::OpenAiBalanceHttpJson);
6051 provider.trust_exhaustion_for_routing = false;
6052
6053 let snapshot = balance_http_snapshot_from_json(
6054 &provider,
6055 &upstream(),
6056 &serde_json::json!({
6057 "balance": "0"
6058 }),
6059 100,
6060 Some(1_000),
6061 );
6062
6063 assert_eq!(snapshot.status, BalanceSnapshotStatus::Exhausted);
6064 assert_eq!(snapshot.exhausted, Some(true));
6065 assert!(!snapshot.exhaustion_affects_routing);
6066 assert!(!snapshot.routing_exhausted());
6067 }
6068
6069 #[test]
6070 fn provider_exhaustion_trust_defaults_to_enabled_when_omitted() {
6071 let provider: UsageProviderConfig = serde_json::from_value(serde_json::json!({
6072 "id": "sub2api",
6073 "kind": "openai_balance_http_json",
6074 "domains": ["example.com"]
6075 }))
6076 .expect("provider config");
6077
6078 assert!(provider.trust_exhaustion_for_routing);
6079 }
6080
6081 #[tokio::test]
6082 async fn provider_missing_token_clears_stale_lb_exhaustion_marker() {
6083 let cfg = proxy_config(vec![service_config(
6084 "right",
6085 vec![
6086 upstream_config("https://primary.example/v1"),
6087 upstream_config("https://backup.example/v1"),
6088 ],
6089 )]);
6090 let lb_states = Arc::new(Mutex::new(HashMap::new()));
6091 let target = UsageProviderTarget {
6092 upstream: upstream(),
6093 base_url: "https://backup.example/v1".to_string(),
6094 provider_id: Some("right".to_string()),
6095 };
6096 let upstreams = vec![target.upstream.clone()];
6097 let state = ProxyState::new();
6098 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, true).await;
6099 {
6100 let guard = lb_states.lock().expect("lb states");
6101 assert!(
6102 guard
6103 .get("right")
6104 .and_then(|entry| entry.usage_exhausted.get(1))
6105 .copied()
6106 .unwrap_or(false)
6107 );
6108 }
6109
6110 let outcome = refresh_provider_target(RefreshProviderTargetParams {
6111 client: &Client::new(),
6112 provider: &provider("sub2api", ProviderKind::OpenAiBalanceHttpJson),
6113 target: &target,
6114 cfg: &cfg,
6115 lb_states: &lb_states,
6116 state: &state,
6117 service_name: "codex",
6118 interval_secs: 60,
6119 })
6120 .await;
6121
6122 assert_eq!(outcome, UsageProviderRefreshOutcome::MissingToken);
6123 let guard = lb_states.lock().expect("lb states");
6124 assert!(
6125 !guard
6126 .get("right")
6127 .and_then(|entry| entry.usage_exhausted.get(1))
6128 .copied()
6129 .unwrap_or(true)
6130 );
6131 }
6132
6133 #[tokio::test]
6134 async fn usage_exhaustion_syncs_owned_balance_policy_action() {
6135 let cfg = proxy_config(vec![service_config(
6136 "right",
6137 vec![
6138 upstream_config("https://primary.example/v1"),
6139 upstream_config("https://backup.example/v1"),
6140 ],
6141 )]);
6142 let lb_states = Arc::new(Mutex::new(HashMap::new()));
6143 let upstreams = vec![endpoint_upstream()];
6144 let state = ProxyState::new();
6145
6146 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, true).await;
6147 let actions = state.list_policy_actions("codex").await;
6148 assert_eq!(actions.len(), 1);
6149 assert_eq!(actions[0].source_signal.kind, ProviderSignalKind::Balance);
6150 assert_eq!(actions[0].reason, "balance_exhausted");
6151
6152 update_usage_exhausted(&lb_states, &state, &cfg, "codex", &upstreams, false).await;
6153 assert!(state.list_policy_actions("codex").await.is_empty());
6154 }
6155
6156 #[test]
6157 fn openai_balance_snapshot_supports_custom_paths_and_derived_budget() {
6158 let mut provider = provider("custom", ProviderKind::OpenAiBalanceHttpJson);
6159 provider.extract.remaining_balance_paths = vec!["payload.remaining_usd".to_string()];
6160 provider.extract.monthly_spent_paths = vec!["payload.used_usd".to_string()];
6161 provider.extract.derive_budget_from_remaining_and_spent = true;
6162
6163 let snapshot = balance_http_snapshot_from_json(
6164 &provider,
6165 &upstream(),
6166 &serde_json::json!({
6167 "payload": {
6168 "remaining_usd": "2",
6169 "used_usd": "0.5"
6170 }
6171 }),
6172 100,
6173 Some(1_000),
6174 );
6175
6176 assert_eq!(snapshot.total_balance_usd.as_deref(), Some("2"));
6177 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
6178 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("2.5"));
6179 assert_eq!(snapshot.exhausted, Some(false));
6180 }
6181
6182 #[test]
6183 fn new_api_snapshot_converts_quota_units_like_cc_switch_template() {
6184 let snapshot = new_api_snapshot_from_json(
6185 &provider("newapi", ProviderKind::NewApiUserSelf),
6186 &upstream(),
6187 &serde_json::json!({
6188 "success": true,
6189 "data": {
6190 "quota": 500000,
6191 "used_quota": 250000
6192 }
6193 }),
6194 100,
6195 Some(1_000),
6196 );
6197
6198 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
6199 assert_eq!(snapshot.exhausted, Some(false));
6200 assert_eq!(snapshot.total_balance_usd, None);
6201 assert_eq!(snapshot.quota_period.as_deref(), Some("quota"));
6202 assert_eq!(snapshot.quota_remaining_usd.as_deref(), Some("1"));
6203 assert_eq!(snapshot.quota_limit_usd.as_deref(), Some("1.5"));
6204 assert_eq!(snapshot.quota_used_usd.as_deref(), Some("0.5"));
6205 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
6206 assert_eq!(snapshot.monthly_budget_usd.as_deref(), Some("1.5"));
6207 }
6208
6209 #[test]
6210 fn new_api_user_self_honors_unlimited_quota_flag() {
6211 let snapshot = new_api_snapshot_from_json(
6212 &provider("newapi", ProviderKind::NewApiUserSelf),
6213 &upstream(),
6214 &serde_json::json!({
6215 "success": true,
6216 "data": {
6217 "quota": 0,
6218 "used_quota": 250000,
6219 "unlimited_quota": true
6220 }
6221 }),
6222 100,
6223 Some(1_000),
6224 );
6225
6226 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
6227 assert_eq!(snapshot.exhausted, Some(false));
6228 assert_eq!(snapshot.total_balance_usd, None);
6229 assert_eq!(snapshot.monthly_budget_usd, None);
6230 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
6231 assert_eq!(snapshot.unlimited_quota, Some(true));
6232 }
6233
6234 #[test]
6235 fn new_api_token_usage_honors_unlimited_quota_flag() {
6236 let snapshot = new_api_token_usage_snapshot_from_json(
6237 &provider("newapi-token", ProviderKind::NewApiTokenUsage),
6238 &upstream(),
6239 &serde_json::json!({
6240 "code": true,
6241 "message": "ok",
6242 "data": {
6243 "object": "token_usage",
6244 "name": "demo-token",
6245 "total_granted": 0,
6246 "total_used": 250000,
6247 "total_available": 0,
6248 "unlimited_quota": true
6249 }
6250 }),
6251 100,
6252 Some(1_000),
6253 );
6254
6255 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
6256 assert_eq!(snapshot.exhausted, Some(false));
6257 assert_eq!(snapshot.plan_name.as_deref(), Some("demo-token"));
6258 assert_eq!(snapshot.total_balance_usd, None);
6259 assert_eq!(snapshot.monthly_budget_usd, None);
6260 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("0.5"));
6261 assert_eq!(snapshot.unlimited_quota, Some(true));
6262 }
6263
6264 #[test]
6265 fn openai_organization_costs_sums_official_cost_buckets_without_exhaustion() {
6266 let snapshot = openai_organization_costs_snapshot_from_json(
6267 &provider("openai", ProviderKind::OpenAiOrganizationCosts),
6268 &upstream(),
6269 &serde_json::json!({
6270 "object": "page",
6271 "data": [
6272 {
6273 "object": "bucket",
6274 "start_time": 1710000000,
6275 "end_time": 1710086400,
6276 "results": [
6277 {
6278 "object": "organization.costs.result",
6279 "amount": { "value": 1.25, "currency": "usd" }
6280 },
6281 {
6282 "object": "organization.costs.result",
6283 "amount": { "value": "2.5", "currency": "usd" }
6284 },
6285 {
6286 "object": "organization.costs.result",
6287 "amount": { "value": 99, "currency": "eur" }
6288 }
6289 ]
6290 },
6291 {
6292 "object": "bucket",
6293 "results": [
6294 {
6295 "object": "organization.costs.result",
6296 "amount": { "value": "0.25", "currency": "USD" }
6297 }
6298 ]
6299 }
6300 ],
6301 "has_more": false
6302 }),
6303 100,
6304 Some(1_000),
6305 );
6306
6307 assert_eq!(snapshot.status, BalanceSnapshotStatus::Ok);
6308 assert_eq!(snapshot.exhausted, None);
6309 assert!(!snapshot.exhaustion_affects_routing);
6310 assert!(!snapshot.routing_exhausted());
6311 assert_eq!(snapshot.monthly_spent_usd.as_deref(), Some("4"));
6312 assert_eq!(snapshot.total_balance_usd, None);
6313 }
6314}