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