Skip to main content

ironflow_core/account/
mod.rs

1//! Provider Accounts: credentials and usage limits of AI provider accounts.
2//!
3//! A Provider Account is an account at an AI provider (v1: a Claude Pro/Max
4//! subscription). This module holds the provider-agnostic pieces:
5//!
6//! * [`AccountWindow`] - one usage-limit window (e.g. the 5 hour window) as
7//!   last observed.
8//! * [`AccountCredential`] / [`AccountSession`] - the credential a provider
9//!   injects into its process, and the recorder it reports observed windows to.
10//! * [`AccountKind`] - a kind of account (form fields, credential validation,
11//!   live credential check). [`ClaudeSubscriptionKind`] is the v1 kind.
12//!
13//! Selection among accounts lives in [`crate::account_strategy`].
14//!
15//! # Examples
16//!
17//! ```
18//! use ironflow_core::account::{AccountKind, ClaudeSubscriptionKind};
19//!
20//! let kind = ClaudeSubscriptionKind::new();
21//! assert!(kind.validate_credential("not-a-token").is_err());
22//! ```
23
24use std::fmt;
25use std::future::Future;
26use std::mem;
27use std::pin::Pin;
28use std::sync::{Arc, Mutex};
29use std::time::Duration;
30
31use chrono::{DateTime, TimeDelta, TimeZone, Utc};
32use reqwest::header::{HeaderMap, RETRY_AFTER};
33use reqwest::{Client, StatusCode};
34use serde::{Deserialize, Serialize};
35use serde_json::json;
36use strum::{Display, EnumString};
37use thiserror::Error;
38
39/// Status of a usage-limit window as reported by the provider.
40///
41/// # Examples
42///
43/// ```
44/// use ironflow_core::account::WindowStatus;
45///
46/// assert_eq!(WindowStatus::AllowedWarning.to_string(), "allowed_warning");
47/// ```
48#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display, EnumString)]
49#[serde(rename_all = "snake_case")]
50#[strum(serialize_all = "snake_case")]
51pub enum WindowStatus {
52    /// Requests are allowed.
53    Allowed,
54    /// Requests are allowed, but the window is close to its limit.
55    AllowedWarning,
56    /// Requests are rejected until the window resets.
57    Rejected,
58}
59
60/// One usage-limit window of a Provider Account, as last observed.
61///
62/// # Examples
63///
64/// ```
65/// use chrono::Utc;
66/// use ironflow_core::account::{AccountWindow, WindowStatus};
67///
68/// let window = AccountWindow {
69///     window: "seven_day".to_string(),
70///     utilization: 1.0,
71///     resets_at: None,
72///     status: WindowStatus::Rejected,
73///     model_scope: Some("opus".to_string()),
74///     observed_at: Utc::now(),
75/// };
76/// assert!(window.applies_to("claude-opus-4-1"));
77/// assert!(!window.applies_to("claude-sonnet-4-5"));
78/// assert!(window.is_exhausted(Utc::now()));
79/// ```
80#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
81pub struct AccountWindow {
82    /// Window name, e.g. `"five_hour"` or `"seven_day"`.
83    pub window: String,
84    /// Fraction of the window used, between `0.0` and `1.0`.
85    pub utilization: f64,
86    /// When the window resets, if known.
87    pub resets_at: Option<DateTime<Utc>>,
88    /// Provider-reported status.
89    pub status: WindowStatus,
90    /// Model family the window applies to (`"opus"`), `None` for every model.
91    pub model_scope: Option<String>,
92    /// When the window was observed.
93    pub observed_at: DateTime<Utc>,
94}
95
96impl AccountWindow {
97    /// Whether this window constrains requests for `model`.
98    ///
99    /// # Examples
100    ///
101    /// See [`AccountWindow`].
102    pub fn applies_to(&self, model: &str) -> bool {
103        match &self.model_scope {
104            None => true,
105            Some(scope) => model
106                .to_ascii_lowercase()
107                .contains(&scope.to_ascii_lowercase()),
108        }
109    }
110
111    /// Whether the window rejects requests at `now`.
112    ///
113    /// A rejected window stops blocking once its `resets_at` has passed.
114    ///
115    /// # Examples
116    ///
117    /// See [`AccountWindow`].
118    pub fn is_exhausted(&self, now: DateTime<Utc>) -> bool {
119        self.status == WindowStatus::Rejected && self.resets_at.is_none_or(|reset| reset > now)
120    }
121
122    /// Utilization at `now`: `0.0` once the window has reset.
123    ///
124    /// # Examples
125    ///
126    /// ```
127    /// use chrono::{TimeDelta, Utc};
128    /// use ironflow_core::account::{AccountWindow, WindowStatus};
129    ///
130    /// let now = Utc::now();
131    /// let window = AccountWindow {
132    ///     window: "five_hour".to_string(),
133    ///     utilization: 0.9,
134    ///     resets_at: Some(now - TimeDelta::minutes(1)),
135    ///     status: WindowStatus::Allowed,
136    ///     model_scope: None,
137    ///     observed_at: now,
138    /// };
139    /// assert_eq!(window.effective_utilization(now), 0.0);
140    /// ```
141    pub fn effective_utilization(&self, now: DateTime<Utc>) -> f64 {
142        match self.resets_at {
143            Some(reset) if reset <= now => 0.0,
144            _ => self.utilization,
145        }
146    }
147}
148
149/// The credential of a Provider Account, as injected into a provider process.
150///
151/// Its [`Debug`] output never shows the value, and it is not serializable.
152///
153/// # Examples
154///
155/// ```
156/// use ironflow_core::account::AccountCredential;
157///
158/// let credential = AccountCredential::new("CLAUDE_CODE_OAUTH_TOKEN", "secret".to_string());
159/// assert_eq!(credential.env_var(), "CLAUDE_CODE_OAUTH_TOKEN");
160/// assert!(!format!("{credential:?}").contains("secret"));
161/// ```
162#[derive(Clone)]
163pub struct AccountCredential {
164    env_var: &'static str,
165    value: String,
166}
167
168impl AccountCredential {
169    /// Wrap a credential value with the environment variable carrying it.
170    ///
171    /// # Examples
172    ///
173    /// See [`AccountCredential`].
174    pub fn new(env_var: &'static str, value: String) -> Self {
175        Self { env_var, value }
176    }
177
178    /// Environment variable the credential is exposed as.
179    ///
180    /// # Examples
181    ///
182    /// See [`AccountCredential`].
183    pub fn env_var(&self) -> &'static str {
184        self.env_var
185    }
186
187    /// The raw credential value. Never log it.
188    ///
189    /// # Examples
190    ///
191    /// ```
192    /// use ironflow_core::account::AccountCredential;
193    ///
194    /// let credential = AccountCredential::new("TOKEN", "value".to_string());
195    /// assert_eq!(credential.expose(), "value");
196    /// ```
197    pub fn expose(&self) -> &str {
198        &self.value
199    }
200}
201
202impl fmt::Debug for AccountCredential {
203    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
204        f.debug_struct("AccountCredential")
205            .field("env_var", &self.env_var)
206            .field("value", &"[REDACTED]")
207            .finish()
208    }
209}
210
211/// Collects the usage windows a provider observes during an invocation.
212///
213/// Only the last observation per `(window, model_scope)` is kept.
214///
215/// # Examples
216///
217/// ```
218/// use chrono::Utc;
219/// use ironflow_core::account::{AccountWindow, RateLimitRecorder, WindowStatus};
220///
221/// let recorder = RateLimitRecorder::default();
222/// let window = AccountWindow {
223///     window: "five_hour".to_string(),
224///     utilization: 0.2,
225///     resets_at: None,
226///     status: WindowStatus::Allowed,
227///     model_scope: None,
228///     observed_at: Utc::now(),
229/// };
230/// recorder.record(window.clone());
231/// recorder.record(AccountWindow { utilization: 0.3, ..window });
232/// let windows = recorder.take();
233/// assert_eq!(windows.len(), 1);
234/// assert_eq!(windows[0].utilization, 0.3);
235/// assert!(recorder.take().is_empty());
236/// ```
237#[derive(Debug, Clone, Default)]
238pub struct RateLimitRecorder(Arc<Mutex<Vec<AccountWindow>>>);
239
240impl RateLimitRecorder {
241    /// Record an observed window, replacing a previous one with the same
242    /// `(window, model_scope)`.
243    ///
244    /// # Examples
245    ///
246    /// See [`RateLimitRecorder`].
247    pub fn record(&self, window: AccountWindow) {
248        let mut windows = self.0.lock().unwrap_or_else(|e| e.into_inner());
249        windows.retain(|w| !(w.window == window.window && w.model_scope == window.model_scope));
250        windows.push(window);
251    }
252
253    /// Drain every recorded window.
254    ///
255    /// # Examples
256    ///
257    /// See [`RateLimitRecorder`].
258    pub fn take(&self) -> Vec<AccountWindow> {
259        let mut windows = self.0.lock().unwrap_or_else(|e| e.into_inner());
260        mem::take(&mut *windows)
261    }
262}
263
264/// The Provider Account a single invocation runs under.
265///
266/// # Examples
267///
268/// ```
269/// use ironflow_core::account::{AccountCredential, AccountSession, RateLimitRecorder};
270///
271/// let session = AccountSession::new(
272///     AccountCredential::new("CLAUDE_CODE_OAUTH_TOKEN", "secret".to_string()),
273///     RateLimitRecorder::default(),
274/// );
275/// assert_eq!(session.credential().env_var(), "CLAUDE_CODE_OAUTH_TOKEN");
276/// assert!(!format!("{session:?}").contains("secret"));
277/// ```
278#[derive(Clone)]
279pub struct AccountSession {
280    credential: AccountCredential,
281    recorder: RateLimitRecorder,
282}
283
284impl AccountSession {
285    /// Build a session from a credential and the recorder observations go to.
286    ///
287    /// # Examples
288    ///
289    /// See [`AccountSession`].
290    pub fn new(credential: AccountCredential, recorder: RateLimitRecorder) -> Self {
291        Self {
292            credential,
293            recorder,
294        }
295    }
296
297    /// The credential to inject.
298    ///
299    /// # Examples
300    ///
301    /// See [`AccountSession`].
302    pub fn credential(&self) -> &AccountCredential {
303        &self.credential
304    }
305
306    /// The recorder observed windows are pushed to.
307    ///
308    /// # Examples
309    ///
310    /// ```
311    /// use ironflow_core::account::{AccountCredential, AccountSession, RateLimitRecorder};
312    ///
313    /// let recorder = RateLimitRecorder::default();
314    /// let session = AccountSession::new(AccountCredential::new("T", "v".to_string()), recorder);
315    /// assert!(session.recorder().take().is_empty());
316    /// ```
317    pub fn recorder(&self) -> &RateLimitRecorder {
318        &self.recorder
319    }
320}
321
322impl fmt::Debug for AccountSession {
323    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
324        f.debug_struct("AccountSession")
325            .field("credential", &self.credential)
326            .finish_non_exhaustive()
327    }
328}
329
330/// Outcome of a successful live credential check.
331///
332/// # Examples
333///
334/// ```
335/// use ironflow_core::account::CredentialCheck;
336///
337/// let check = CredentialCheck::Valid { windows: Vec::new() };
338/// assert!(check.windows().is_empty());
339/// ```
340#[derive(Debug, Clone, PartialEq)]
341pub enum CredentialCheck {
342    /// The provider accepted the credential.
343    Valid {
344        /// Windows reported with the answer.
345        windows: Vec<AccountWindow>,
346    },
347    /// The credential is valid but currently rate limited.
348    Limited {
349        /// Windows reported with the answer, at least one rejected.
350        windows: Vec<AccountWindow>,
351    },
352}
353
354impl CredentialCheck {
355    /// Windows reported by the check.
356    ///
357    /// # Examples
358    ///
359    /// See [`CredentialCheck`].
360    pub fn windows(&self) -> &[AccountWindow] {
361        match self {
362            Self::Valid { windows } | Self::Limited { windows } => windows,
363        }
364    }
365}
366
367/// Errors raised while validating or checking an account credential.
368///
369/// # Examples
370///
371/// ```
372/// use ironflow_core::account::AccountError;
373///
374/// let err = AccountError::Unauthorized { status: 401 };
375/// assert!(err.to_string().contains("401"));
376/// ```
377#[derive(Debug, Error)]
378pub enum AccountError {
379    /// The credential format was rejected before any call.
380    #[error("invalid credential: {0}")]
381    InvalidCredential(String),
382    /// The provider rejected the credential.
383    #[error("credential rejected by provider (HTTP {status})")]
384    Unauthorized {
385        /// HTTP status returned (401 or 403).
386        status: u16,
387    },
388    /// The check could not complete (network error, unexpected status).
389    #[error("credential check failed: {0}")]
390    CheckFailed(String),
391}
392
393/// A form field shown when adding an account of a kind.
394///
395/// # Examples
396///
397/// ```
398/// use ironflow_core::account::{AccountKind, ClaudeSubscriptionKind};
399///
400/// let fields = ClaudeSubscriptionKind::new().form_fields();
401/// assert_eq!(fields[0].name, "token");
402/// assert!(fields[0].secret);
403/// ```
404#[derive(Debug, Clone, Serialize)]
405pub struct AccountFormField {
406    /// Field identifier.
407    pub name: &'static str,
408    /// Human-readable label.
409    pub label: &'static str,
410    /// Whether the value is a secret (write-only, masked).
411    pub secret: bool,
412    /// Help text shown next to the field.
413    pub help: &'static str,
414}
415
416/// Future returned by [`AccountKind::check_credential`].
417pub type AccountCheckFuture<'a> =
418    Pin<Box<dyn Future<Output = Result<CredentialCheck, AccountError>> + Send + 'a>>;
419
420/// A kind of Provider Account.
421///
422/// # Examples
423///
424/// ```
425/// use ironflow_core::account::{AccountKind, ClaudeSubscriptionKind};
426///
427/// let kind = ClaudeSubscriptionKind::new();
428/// assert_eq!(kind.id(), "claude_subscription");
429/// let credential = kind.credential("sk-ant-oat01-xxxx");
430/// assert_eq!(credential.env_var(), "CLAUDE_CODE_OAUTH_TOKEN");
431/// ```
432pub trait AccountKind: Send + Sync {
433    /// Stable identifier, stored in the account row.
434    fn id(&self) -> &'static str;
435
436    /// Human-readable name.
437    fn display_name(&self) -> &'static str;
438
439    /// Fields of the add form.
440    fn form_fields(&self) -> Vec<AccountFormField>;
441
442    /// Check the credential format without any network call.
443    ///
444    /// # Errors
445    ///
446    /// Returns [`AccountError::InvalidCredential`] when the format is wrong.
447    fn validate_credential(&self, raw: &str) -> Result<(), AccountError>;
448
449    /// Build the credential injected into provider processes.
450    fn credential(&self, secret: &str) -> AccountCredential;
451
452    /// Check the credential against the provider.
453    ///
454    /// # Errors
455    ///
456    /// Returns [`AccountError::Unauthorized`] when the provider rejects it and
457    /// [`AccountError::CheckFailed`] when the check cannot complete.
458    fn check_credential<'a>(&'a self, secret: &'a str) -> AccountCheckFuture<'a>;
459}
460
461const CLAUDE_TOKEN_PREFIX: &str = "sk-ant-oat01-";
462const CLAUDE_TOKEN_MIN_LEN: usize = 40;
463const CHECK_TIMEOUT: Duration = Duration::from_secs(15);
464const CHECK_MODEL: &str = "claude-haiku-4-5-20251001";
465const MAX_ERROR_BODY: usize = 200;
466
467/// Claude Pro/Max subscription, authenticated by a `claude setup-token` token.
468///
469/// # Examples
470///
471/// ```
472/// use ironflow_core::account::{AccountKind, ClaudeSubscriptionKind};
473///
474/// let kind = ClaudeSubscriptionKind::new();
475/// assert_eq!(kind.id(), ClaudeSubscriptionKind::ID);
476/// ```
477#[derive(Debug, Clone)]
478pub struct ClaudeSubscriptionKind {
479    api_base: String,
480    client: Client,
481}
482
483impl ClaudeSubscriptionKind {
484    /// Kind identifier.
485    pub const ID: &'static str = "claude_subscription";
486
487    /// Environment variable the Claude CLI reads the token from.
488    pub const TOKEN_ENV: &'static str = "CLAUDE_CODE_OAUTH_TOKEN";
489
490    /// Kind talking to `https://api.anthropic.com`.
491    ///
492    /// # Examples
493    ///
494    /// See [`ClaudeSubscriptionKind`].
495    pub fn new() -> Self {
496        Self::with_api_base("https://api.anthropic.com")
497    }
498
499    /// Kind talking to another API base URL (tests, proxies).
500    ///
501    /// # Examples
502    ///
503    /// ```
504    /// use ironflow_core::account::ClaudeSubscriptionKind;
505    ///
506    /// let kind = ClaudeSubscriptionKind::with_api_base("http://127.0.0.1:8080");
507    /// # let _ = kind;
508    /// ```
509    pub fn with_api_base(api_base: &str) -> Self {
510        Self {
511            api_base: api_base.trim_end_matches('/').to_string(),
512            client: Client::new(),
513        }
514    }
515}
516
517impl Default for ClaudeSubscriptionKind {
518    fn default() -> Self {
519        Self::new()
520    }
521}
522
523impl AccountKind for ClaudeSubscriptionKind {
524    fn id(&self) -> &'static str {
525        Self::ID
526    }
527
528    fn display_name(&self) -> &'static str {
529        "Claude subscription (Pro/Max)"
530    }
531
532    fn form_fields(&self) -> Vec<AccountFormField> {
533        vec![AccountFormField {
534            name: "token",
535            label: "OAuth token",
536            secret: true,
537            help: "Run `claude setup-token` and paste the sk-ant-oat01-... token",
538        }]
539    }
540
541    fn validate_credential(&self, raw: &str) -> Result<(), AccountError> {
542        let token = raw.trim();
543        let well_formed = token.starts_with(CLAUDE_TOKEN_PREFIX)
544            && token.len() >= CLAUDE_TOKEN_MIN_LEN
545            && token
546                .chars()
547                .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-');
548        if well_formed {
549            Ok(())
550        } else {
551            Err(AccountError::InvalidCredential(
552                "expected a `claude setup-token` token (sk-ant-oat01-...)".to_string(),
553            ))
554        }
555    }
556
557    fn credential(&self, secret: &str) -> AccountCredential {
558        AccountCredential::new(Self::TOKEN_ENV, secret.trim().to_string())
559    }
560
561    fn check_credential<'a>(&'a self, secret: &'a str) -> AccountCheckFuture<'a> {
562        Box::pin(async move {
563            let body = json!({
564                "model": CHECK_MODEL,
565                "max_tokens": 1,
566                "system": "You are Claude Code, Anthropic's official CLI for Claude.",
567                "messages": [{"role": "user", "content": "ping"}],
568            });
569            let response = self
570                .client
571                .post(format!("{}/v1/messages", self.api_base))
572                .timeout(CHECK_TIMEOUT)
573                .bearer_auth(secret.trim())
574                .header("anthropic-version", "2023-06-01")
575                .header("anthropic-beta", "oauth-2025-04-20")
576                .json(&body)
577                .send()
578                .await
579                .map_err(|e| AccountError::CheckFailed(e.without_url().to_string()))?;
580
581            let status = response.status();
582            let now = Utc::now();
583            if status.is_success() {
584                return Ok(CredentialCheck::Valid {
585                    windows: windows_from_headers(response.headers(), now, false),
586                });
587            }
588            if status == StatusCode::TOO_MANY_REQUESTS {
589                return Ok(CredentialCheck::Limited {
590                    windows: windows_from_headers(response.headers(), now, true),
591                });
592            }
593            if status == StatusCode::UNAUTHORIZED || status == StatusCode::FORBIDDEN {
594                return Err(AccountError::Unauthorized {
595                    status: status.as_u16(),
596                });
597            }
598            let text = match response.text().await {
599                Ok(text) => text,
600                Err(e) => format!("<unreadable body: {e}>"),
601            };
602            let truncated: String = text.chars().take(MAX_ERROR_BODY).collect();
603            Err(AccountError::CheckFailed(format!(
604                "HTTP {}: {truncated}",
605                status.as_u16()
606            )))
607        })
608    }
609}
610
611fn header_str<'h>(headers: &'h HeaderMap, name: &str) -> Option<&'h str> {
612    headers
613        .get(name)
614        .and_then(|v| v.to_str().ok())
615        .map(str::trim)
616}
617
618fn unix_seconds(value: &str) -> Option<DateTime<Utc>> {
619    let secs = value.parse::<i64>().ok()?;
620    Utc.timestamp_opt(secs, 0).single()
621}
622
623/// Read the unified rate-limit headers into windows.
624///
625/// With `limited` set and no window header, builds one rejected `five_hour`
626/// window from the generic reset headers.
627fn windows_from_headers(
628    headers: &HeaderMap,
629    now: DateTime<Utc>,
630    limited: bool,
631) -> Vec<AccountWindow> {
632    let mut windows = Vec::new();
633    for (suffix, name) in [("5h", "five_hour"), ("7d", "seven_day")] {
634        let prefix = format!("anthropic-ratelimit-unified-{suffix}");
635        let utilization = header_str(headers, &format!("{prefix}-utilization"))
636            .and_then(|v| v.parse::<f64>().ok());
637        let status = header_str(headers, &format!("{prefix}-status"))
638            .and_then(|v| v.parse::<WindowStatus>().ok());
639        if utilization.is_none() && status.is_none() {
640            continue;
641        }
642        let status = status.unwrap_or(WindowStatus::Allowed);
643        let utilization = utilization
644            .map(|u| if u > 1.0 { u / 100.0 } else { u })
645            .unwrap_or(if status == WindowStatus::Rejected {
646                1.0
647            } else {
648                0.0
649            })
650            .clamp(0.0, 1.0);
651        windows.push(AccountWindow {
652            window: name.to_string(),
653            utilization,
654            resets_at: header_str(headers, &format!("{prefix}-reset")).and_then(unix_seconds),
655            status,
656            model_scope: None,
657            observed_at: now,
658        });
659    }
660
661    if limited && windows.is_empty() {
662        let resets_at = header_str(headers, "anthropic-ratelimit-unified-reset")
663            .and_then(unix_seconds)
664            .or_else(|| {
665                header_str(headers, RETRY_AFTER.as_str())
666                    .and_then(|v| v.parse::<i64>().ok())
667                    .map(|secs| now + TimeDelta::seconds(secs))
668            });
669        windows.push(AccountWindow {
670            window: "five_hour".to_string(),
671            utilization: 1.0,
672            resets_at,
673            status: WindowStatus::Rejected,
674            model_scope: None,
675            observed_at: now,
676        });
677    }
678    windows
679}
680
681#[cfg(test)]
682mod tests {
683    use super::*;
684    use crate::account_strategy::{AccountCandidate, LeastUtilized, select_account};
685    use tokio::io::{AsyncReadExt, AsyncWriteExt};
686    use tokio::net::TcpListener;
687
688    const GOOD_TOKEN: &str = "sk-ant-oat01-abcdefghijklmnopqrstuvwxyz0123456789_-AB";
689
690    #[test]
691    fn validate_credential_accepts_well_formed_token() {
692        let kind = ClaudeSubscriptionKind::new();
693        assert!(kind.validate_credential(GOOD_TOKEN).is_ok());
694        assert!(
695            kind.validate_credential(&format!("  {GOOD_TOKEN}\n"))
696                .is_ok()
697        );
698    }
699
700    #[test]
701    fn validate_credential_rejects_wrong_prefix() {
702        let kind = ClaudeSubscriptionKind::new();
703        let token = GOOD_TOKEN.replace("oat01", "api03");
704        assert!(matches!(
705            kind.validate_credential(&token),
706            Err(AccountError::InvalidCredential(_))
707        ));
708    }
709
710    #[test]
711    fn validate_credential_rejects_short_token() {
712        let kind = ClaudeSubscriptionKind::new();
713        assert!(kind.validate_credential("sk-ant-oat01-invalid").is_err());
714        assert!(kind.validate_credential("").is_err());
715    }
716
717    #[test]
718    fn validate_credential_rejects_bad_chars() {
719        let kind = ClaudeSubscriptionKind::new();
720        let token = format!("{GOOD_TOKEN}$é");
721        assert!(kind.validate_credential(&token).is_err());
722        let token = format!("{GOOD_TOKEN} abc");
723        assert!(kind.validate_credential(&token).is_err());
724    }
725
726    #[test]
727    fn credential_and_session_debug_are_redacted() {
728        let credential = AccountCredential::new("CLAUDE_CODE_OAUTH_TOKEN", GOOD_TOKEN.to_string());
729        let debug = format!("{credential:?}");
730        assert!(!debug.contains(GOOD_TOKEN));
731        assert!(debug.contains("[REDACTED]"));
732        let session = AccountSession::new(credential, RateLimitRecorder::default());
733        assert!(!format!("{session:?}").contains(GOOD_TOKEN));
734    }
735
736    #[test]
737    fn credential_trims_the_secret() {
738        let kind = ClaudeSubscriptionKind::new();
739        let credential = kind.credential(&format!("{GOOD_TOKEN}\n"));
740        assert_eq!(credential.expose(), GOOD_TOKEN);
741    }
742
743    #[test]
744    fn form_fields_have_one_secret_token() {
745        let fields = ClaudeSubscriptionKind::new().form_fields();
746        assert_eq!(fields.len(), 1);
747        assert_eq!(fields[0].name, "token");
748        assert!(fields[0].secret);
749        assert!(fields[0].help.contains("claude setup-token"));
750    }
751
752    /// Serve one HTTP request on a real local socket with a canned response.
753    async fn stub_server(response: String) -> String {
754        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
755        let addr = listener.local_addr().unwrap();
756        tokio::spawn(async move {
757            let (mut socket, _) = listener.accept().await.unwrap();
758            let mut buf = Vec::new();
759            let mut chunk = [0u8; 4096];
760            loop {
761                let n = socket.read(&mut chunk).await.unwrap();
762                if n == 0 {
763                    break;
764                }
765                buf.extend_from_slice(&chunk[..n]);
766                let text = String::from_utf8_lossy(&buf);
767                if let Some(end) = text.find("\r\n\r\n") {
768                    let content_length = text[..end]
769                        .lines()
770                        .find_map(|l| {
771                            l.to_ascii_lowercase()
772                                .strip_prefix("content-length:")
773                                .map(|v| v.trim().parse::<usize>().unwrap_or(0))
774                        })
775                        .unwrap_or(0);
776                    if buf.len() >= end + 4 + content_length {
777                        break;
778                    }
779                }
780            }
781            socket.write_all(response.as_bytes()).await.unwrap();
782            socket.shutdown().await.unwrap();
783        });
784        format!("http://{addr}")
785    }
786
787    fn http_response(status: &str, headers: &[(&str, String)], body: &str) -> String {
788        let mut out = format!("HTTP/1.1 {status}\r\n");
789        for (name, value) in headers {
790            out.push_str(&format!("{name}: {value}\r\n"));
791        }
792        out.push_str(&format!(
793            "content-length: {}\r\nconnection: close\r\n\r\n{body}",
794            body.len()
795        ));
796        out
797    }
798
799    #[tokio::test]
800    async fn check_credential_valid_reads_unified_windows() {
801        let reset = Utc::now().timestamp() + 3600;
802        let response = http_response(
803            "200 OK",
804            &[
805                (
806                    "anthropic-ratelimit-unified-5h-utilization",
807                    "0.42".to_string(),
808                ),
809                ("anthropic-ratelimit-unified-5h-reset", reset.to_string()),
810                (
811                    "anthropic-ratelimit-unified-5h-status",
812                    "allowed".to_string(),
813                ),
814                (
815                    "anthropic-ratelimit-unified-7d-utilization",
816                    "0.8".to_string(),
817                ),
818                (
819                    "anthropic-ratelimit-unified-7d-status",
820                    "allowed_warning".to_string(),
821                ),
822            ],
823            "{}",
824        );
825        let base = stub_server(response).await;
826        let kind = ClaudeSubscriptionKind::with_api_base(&base);
827        let check = kind.check_credential(GOOD_TOKEN).await.unwrap();
828        let CredentialCheck::Valid { windows } = check else {
829            panic!("expected Valid, got {check:?}");
830        };
831        assert_eq!(windows.len(), 2);
832        assert_eq!(windows[0].window, "five_hour");
833        assert!((windows[0].utilization - 0.42).abs() < 1e-9);
834        assert_eq!(windows[0].resets_at.unwrap().timestamp(), reset);
835        assert_eq!(windows[1].window, "seven_day");
836        assert_eq!(windows[1].status, WindowStatus::AllowedWarning);
837    }
838
839    #[tokio::test]
840    async fn check_credential_429_is_limited_with_rejected_window() {
841        let reset = Utc::now().timestamp() + 1800;
842        let response = http_response(
843            "429 Too Many Requests",
844            &[("anthropic-ratelimit-unified-reset", reset.to_string())],
845            "{\"error\":\"rate_limited\"}",
846        );
847        let base = stub_server(response).await;
848        let kind = ClaudeSubscriptionKind::with_api_base(&base);
849        let check = kind.check_credential(GOOD_TOKEN).await.unwrap();
850        let CredentialCheck::Limited { windows } = check else {
851            panic!("expected Limited, got {check:?}");
852        };
853        assert_eq!(windows.len(), 1);
854        assert_eq!(windows[0].window, "five_hour");
855        assert_eq!(windows[0].status, WindowStatus::Rejected);
856        assert_eq!(windows[0].utilization, 1.0);
857        assert_eq!(windows[0].resets_at.unwrap().timestamp(), reset);
858    }
859
860    #[tokio::test]
861    async fn check_credential_401_is_unauthorized() {
862        let response = http_response("401 Unauthorized", &[], "{}");
863        let base = stub_server(response).await;
864        let kind = ClaudeSubscriptionKind::with_api_base(&base);
865        let err = kind.check_credential(GOOD_TOKEN).await.unwrap_err();
866        assert!(matches!(err, AccountError::Unauthorized { status: 401 }));
867        assert!(!err.to_string().contains(GOOD_TOKEN));
868    }
869
870    #[tokio::test]
871    async fn check_credential_500_is_check_failed() {
872        let response = http_response("500 Internal Server Error", &[], "boom");
873        let base = stub_server(response).await;
874        let kind = ClaudeSubscriptionKind::with_api_base(&base);
875        let err = kind.check_credential(GOOD_TOKEN).await.unwrap_err();
876        let AccountError::CheckFailed(message) = err else {
877            panic!("expected CheckFailed");
878        };
879        assert!(message.contains("HTTP 500"));
880        assert!(message.contains("boom"));
881    }
882
883    #[tokio::test]
884    async fn check_credential_unreachable_is_check_failed() {
885        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
886        let addr = listener.local_addr().unwrap();
887        drop(listener);
888        let kind = ClaudeSubscriptionKind::with_api_base(&format!("http://{addr}"));
889        let err = kind.check_credential(GOOD_TOKEN).await.unwrap_err();
890        assert!(matches!(err, AccountError::CheckFailed(_)));
891    }
892
893    /// A kind limited by a dollar budget, to keep [`AccountKind`] honest
894    /// beyond Claude subscriptions.
895    struct BudgetKind;
896
897    impl AccountKind for BudgetKind {
898        fn id(&self) -> &'static str {
899            "budget"
900        }
901
902        fn display_name(&self) -> &'static str {
903            "Budget"
904        }
905
906        fn form_fields(&self) -> Vec<AccountFormField> {
907            vec![AccountFormField {
908                name: "api_key",
909                label: "API key",
910                secret: true,
911                help: "Paste the API key",
912            }]
913        }
914
915        fn validate_credential(&self, raw: &str) -> Result<(), AccountError> {
916            if raw.is_empty() {
917                Err(AccountError::InvalidCredential("empty".to_string()))
918            } else {
919                Ok(())
920            }
921        }
922
923        fn credential(&self, secret: &str) -> AccountCredential {
924            AccountCredential::new("BUDGET_API_KEY", secret.to_string())
925        }
926
927        fn check_credential<'a>(&'a self, _secret: &'a str) -> AccountCheckFuture<'a> {
928            Box::pin(async {
929                Ok(CredentialCheck::Valid {
930                    windows: vec![budget_window(0.5, WindowStatus::Allowed)],
931                })
932            })
933        }
934    }
935
936    fn budget_window(utilization: f64, status: WindowStatus) -> AccountWindow {
937        AccountWindow {
938            window: "monthly_budget_usd".to_string(),
939            utilization,
940            resets_at: Some(Utc::now() + TimeDelta::days(10)),
941            status,
942            model_scope: None,
943            observed_at: Utc::now(),
944        }
945    }
946
947    #[tokio::test]
948    async fn budget_kind_goes_through_select_account() {
949        let kind = BudgetKind;
950        assert!(kind.validate_credential("").is_err());
951        let check = kind.check_credential("key").await.unwrap();
952        let spent = AccountCandidate {
953            id: "a".to_string(),
954            name: "spent".to_string(),
955            priority: 1,
956            max_concurrency: None,
957            running_steps: 0,
958            windows: vec![budget_window(1.0, WindowStatus::Rejected)],
959        };
960        let fresh = AccountCandidate {
961            id: "b".to_string(),
962            name: "fresh".to_string(),
963            priority: 2,
964            max_concurrency: None,
965            running_steps: 0,
966            windows: check.windows().to_vec(),
967        };
968        let candidates = [spent, fresh];
969        let selected =
970            select_account(&LeastUtilized, &candidates, "any-model", Utc::now()).unwrap();
971        assert_eq!(selected.name, "fresh");
972        assert_eq!(kind.credential("key").env_var(), "BUDGET_API_KEY");
973    }
974}