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