Skip to main content

codex_helper_core/
usage_providers.rs

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