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