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