Skip to main content

ironflow_engine/
accounts.rs

1//! Provider Account selection and injection for agent steps.
2//!
3//! [`AccountAwareProvider`] wraps the worker's [`AgentProvider`]. For each
4//! agent invocation on a provider that can inject a Provider Account
5//! credential for the invocation ([`AgentProvider::account_kind_for`]), it:
6//!
7//! 1. lists the candidate accounts of that kind from the store,
8//! 2. picks one with an [`AccountStrategy`],
9//! 3. reads its credential and runs the invocation under it,
10//! 4. records the usage windows observed during the invocation.
11//!
12//! With no account of the kind in the store, or a provider that cannot
13//! inject one, the invocation runs unchanged with the worker environment.
14//!
15//! # Examples
16//!
17//! ```no_run
18//! use std::sync::Arc;
19//! use ironflow_core::account_strategy::Priority;
20//! use ironflow_core::providers::claude::ClaudeCodeProvider;
21//! use ironflow_engine::accounts::AccountAwareProvider;
22//! use ironflow_store::memory::InMemoryStore;
23//!
24//! let provider = AccountAwareProvider::new(
25//!     Arc::new(ClaudeCodeProvider::new()),
26//!     Arc::new(InMemoryStore::new()),
27//! )
28//! .with_strategy(Arc::new(Priority));
29//! # let _ = provider;
30//! ```
31
32use std::collections::HashMap;
33use std::fmt;
34use std::sync::Arc;
35
36use chrono::Utc;
37use ironflow_core::account::{
38    AccountKind, AccountSession, AccountWindow, ClaudeSubscriptionKind, RateLimitRecorder,
39    WindowStatus,
40};
41use ironflow_core::account_strategy::{
42    AccountCandidate, AccountStrategy, LeastUtilized, select_account,
43};
44use ironflow_core::error::AgentError;
45use ironflow_core::provider::{
46    AgentConfig, AgentOutput, AgentProvider, InvokeFuture, LogSink, ReleaseFuture,
47};
48use ironflow_store::entities::{
49    AccountWindowStatus, NewAccountWindow, NewProviderAccountObservation, ProviderAccount,
50    ProviderAccountCandidate, ProviderAccountWindow,
51};
52use ironflow_store::store::Store;
53use tracing::{debug, info, warn};
54
55/// Convert a stored window into the core representation.
56///
57/// # Examples
58///
59/// ```
60/// use chrono::Utc;
61/// use ironflow_engine::accounts::window_from_store;
62/// use ironflow_store::entities::{AccountWindowStatus, ProviderAccountWindow};
63/// use uuid::Uuid;
64///
65/// let window = window_from_store(&ProviderAccountWindow {
66///     account_id: Uuid::now_v7(),
67///     window: "five_hour".to_string(),
68///     utilization: 0.4,
69///     resets_at: None,
70///     status: AccountWindowStatus::Allowed,
71///     model_scope: None,
72///     observed_at: Utc::now(),
73/// });
74/// assert_eq!(window.window, "five_hour");
75/// ```
76pub fn window_from_store(window: &ProviderAccountWindow) -> AccountWindow {
77    AccountWindow {
78        window: window.window.clone(),
79        utilization: window.utilization,
80        resets_at: window.resets_at,
81        status: match window.status {
82            AccountWindowStatus::Allowed => WindowStatus::Allowed,
83            AccountWindowStatus::AllowedWarning => WindowStatus::AllowedWarning,
84            AccountWindowStatus::Rejected => WindowStatus::Rejected,
85        },
86        model_scope: window.model_scope.clone(),
87        observed_at: window.observed_at,
88    }
89}
90
91/// Convert an observed core window into the store representation.
92///
93/// # Examples
94///
95/// ```
96/// use chrono::Utc;
97/// use ironflow_core::account::{AccountWindow, WindowStatus};
98/// use ironflow_engine::accounts::window_to_store;
99/// use ironflow_store::entities::AccountWindowStatus;
100///
101/// let window = window_to_store(AccountWindow {
102///     window: "seven_day".to_string(),
103///     utilization: 1.0,
104///     resets_at: None,
105///     status: WindowStatus::Rejected,
106///     model_scope: Some("opus".to_string()),
107///     observed_at: Utc::now(),
108/// });
109/// assert_eq!(window.status, AccountWindowStatus::Rejected);
110/// ```
111pub fn window_to_store(window: AccountWindow) -> NewAccountWindow {
112    NewAccountWindow {
113        window: window.window,
114        utilization: window.utilization,
115        resets_at: window.resets_at,
116        status: match window.status {
117            WindowStatus::Allowed => AccountWindowStatus::Allowed,
118            WindowStatus::AllowedWarning => AccountWindowStatus::AllowedWarning,
119            WindowStatus::Rejected => AccountWindowStatus::Rejected,
120        },
121        model_scope: window.model_scope,
122        observed_at: window.observed_at,
123    }
124}
125
126fn to_core_candidate(candidate: &ProviderAccountCandidate) -> AccountCandidate {
127    AccountCandidate {
128        id: candidate.account.id.to_string(),
129        name: candidate.account.name.clone(),
130        priority: candidate.account.priority,
131        max_concurrency: candidate.account.max_concurrency,
132        running_steps: candidate.running_steps,
133        windows: candidate.windows.iter().map(window_from_store).collect(),
134    }
135}
136
137fn resolution_error(message: String) -> AgentError {
138    AgentError::ProcessFailed {
139        exit_code: -1,
140        stderr: message,
141    }
142}
143
144/// An [`AgentProvider`] that runs each invocation under a Provider Account.
145///
146/// See the [module documentation](self).
147pub struct AccountAwareProvider {
148    inner: Arc<dyn AgentProvider>,
149    store: Arc<dyn Store>,
150    strategy: Arc<dyn AccountStrategy>,
151    kinds: HashMap<&'static str, Arc<dyn AccountKind>>,
152}
153
154impl fmt::Debug for AccountAwareProvider {
155    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
156        f.debug_struct("AccountAwareProvider")
157            .field("strategy", &self.strategy.name())
158            .field("kinds", &self.kinds.keys().collect::<Vec<_>>())
159            .finish_non_exhaustive()
160    }
161}
162
163impl AccountAwareProvider {
164    /// Wrap `inner`, reading accounts from `store`, with the
165    /// [`LeastUtilized`] strategy and the [`ClaudeSubscriptionKind`] kind.
166    ///
167    /// # Examples
168    ///
169    /// See the [module documentation](self).
170    pub fn new(inner: Arc<dyn AgentProvider>, store: Arc<dyn Store>) -> Self {
171        let claude: Arc<dyn AccountKind> = Arc::new(ClaudeSubscriptionKind::new());
172        Self {
173            inner,
174            store,
175            strategy: Arc::new(LeastUtilized),
176            kinds: HashMap::from([(claude.id(), claude)]),
177        }
178    }
179
180    /// Replace the account selection strategy.
181    ///
182    /// # Examples
183    ///
184    /// See the [module documentation](self).
185    pub fn with_strategy(mut self, strategy: Arc<dyn AccountStrategy>) -> Self {
186        self.strategy = strategy;
187        self
188    }
189
190    /// Register an additional account kind (or replace one with the same id).
191    ///
192    /// # Examples
193    ///
194    /// ```no_run
195    /// use std::sync::Arc;
196    /// use ironflow_core::account::ClaudeSubscriptionKind;
197    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
198    /// use ironflow_engine::accounts::AccountAwareProvider;
199    /// use ironflow_store::memory::InMemoryStore;
200    ///
201    /// let provider = AccountAwareProvider::new(
202    ///     Arc::new(ClaudeCodeProvider::new()),
203    ///     Arc::new(InMemoryStore::new()),
204    /// )
205    /// .with_kind(Arc::new(ClaudeSubscriptionKind::with_api_base("http://proxy:8080")));
206    /// # let _ = provider;
207    /// ```
208    pub fn with_kind(mut self, kind: Arc<dyn AccountKind>) -> Self {
209        self.kinds.insert(kind.id(), kind);
210        self
211    }
212
213    async fn invoke_inner(
214        &self,
215        config: &AgentConfig,
216        sink: Option<Arc<dyn LogSink>>,
217    ) -> Result<AgentOutput, AgentError> {
218        match sink {
219            Some(sink) => self.inner.invoke_with_logs(config, sink).await,
220            None => self.inner.invoke(config).await,
221        }
222    }
223
224    async fn run(
225        &self,
226        config: &AgentConfig,
227        sink: Option<Arc<dyn LogSink>>,
228    ) -> Result<AgentOutput, AgentError> {
229        let Some(kind_id) = self.inner.account_kind_for(config) else {
230            return self.invoke_inner(config, sink).await;
231        };
232        let Some(kind) = self.kinds.get(kind_id).cloned() else {
233            debug!(
234                kind = kind_id,
235                "no account kind registered, using worker environment"
236            );
237            return self.invoke_inner(config, sink).await;
238        };
239
240        let candidates = self
241            .store
242            .list_provider_account_candidates(kind_id.to_string())
243            .await
244            .map_err(|e| resolution_error(format!("provider account resolution failed: {e}")))?;
245        if candidates.is_empty() {
246            debug!(
247                kind = kind_id,
248                "no provider account for kind, using worker environment"
249            );
250            return self.invoke_inner(config, sink).await;
251        }
252
253        let core_candidates: Vec<AccountCandidate> =
254            candidates.iter().map(to_core_candidate).collect();
255        let now = Utc::now();
256        let Some(selected) =
257            select_account(self.strategy.as_ref(), &core_candidates, &config.model, now)
258        else {
259            let next_reset = core_candidates
260                .iter()
261                .flat_map(|c| c.windows.iter())
262                .filter(|w| w.applies_to(&config.model) && w.is_exhausted(now))
263                .filter_map(|w| w.resets_at)
264                .min()
265                .map_or_else(|| "unknown".to_string(), |at| at.to_rfc3339());
266            return Err(resolution_error(format!(
267                "no provider account available for {kind_id}: all limited or at max_concurrency (next reset {next_reset})"
268            )));
269        };
270        let account: &ProviderAccount = &candidates
271            .iter()
272            .find(|c| c.account.id.to_string() == selected.id)
273            .ok_or_else(|| resolution_error("selected provider account vanished".to_string()))?
274            .account;
275
276        let missing = || {
277            resolution_error(format!(
278                "credential of provider account '{}' is missing",
279                account.name
280            ))
281        };
282        let secret = match self.store.get_secret(&account.secret_key).await {
283            Ok(Some(secret)) => secret,
284            Ok(None) => return Err(missing()),
285            Err(e) => {
286                warn!(account = %account.name, error = %e, "failed to read provider account credential");
287                return Err(missing());
288            }
289        };
290
291        info!(
292            account = %account.name,
293            strategy = self.strategy.name(),
294            model = %config.model,
295            "selected provider account"
296        );
297
298        let recorder = RateLimitRecorder::default();
299        // Rate-limit events only appear in stream-json, hence verbose.
300        let account_config = config
301            .clone()
302            .verbose(true)
303            .account_session(AccountSession::new(
304                kind.credential(&secret.value),
305                recorder.clone(),
306            ));
307
308        let result = self.invoke_inner(&account_config, sink).await;
309
310        let windows = recorder.take();
311        let auth_failed = matches!(
312            result,
313            Err(AgentError::Api {
314                status: Some(401 | 403),
315                ..
316            })
317        );
318        if !windows.is_empty() || auth_failed {
319            let observation = NewProviderAccountObservation {
320                windows: windows.into_iter().map(window_to_store).collect(),
321                auth_failed,
322            };
323            if let Err(e) = self
324                .store
325                .record_provider_account_observation(account.id, observation)
326                .await
327            {
328                warn!(account = %account.name, error = %e, "failed to record provider account usage");
329            }
330        }
331
332        result.map(|mut output| {
333            output.account_id = Some(account.id.to_string());
334            output
335        })
336    }
337}
338
339impl AgentProvider for AccountAwareProvider {
340    fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
341        Box::pin(self.run(config, None))
342    }
343
344    fn invoke_with_logs<'a>(
345        &'a self,
346        config: &'a AgentConfig,
347        log_sink: Arc<dyn LogSink>,
348    ) -> InvokeFuture<'a> {
349        Box::pin(self.run(config, Some(log_sink)))
350    }
351
352    fn release_run<'a>(&'a self, run_id: &'a str) -> ReleaseFuture<'a> {
353        self.inner.release_run(run_id)
354    }
355
356    fn account_kind(&self) -> Option<&'static str> {
357        self.inner.account_kind()
358    }
359
360    fn account_kind_for(&self, config: &AgentConfig) -> Option<&'static str> {
361        self.inner.account_kind_for(config)
362    }
363}
364
365#[cfg(test)]
366mod tests {
367    use std::sync::Mutex;
368
369    use chrono::TimeDelta;
370    use ironflow_core::providers::router::{ProviderMatcher, ProviderRouter};
371    use ironflow_store::crypto::KeyRing;
372    use ironflow_store::entities::{NewProviderAccount, provider_account_secret_key};
373    use ironflow_store::memory::InMemoryStore;
374    use ironflow_store::provider_account_store::ProviderAccountStore;
375    use ironflow_store::secret_store::SecretStore;
376    use serde_json::json;
377    use uuid::Uuid;
378
379    use super::*;
380
381    const TOKEN: &str = "sk-ant-oat01-test-token-abcdefghijklmnopqrstuvwxyz";
382
383    /// What the test provider does once invoked.
384    #[derive(Clone, Copy)]
385    enum Outcome {
386        Succeed,
387        FailApi(u16),
388    }
389
390    /// A provider that behaves like the Claude transports: it reads the
391    /// injected credential and reports a rate-limit window to the recorder.
392    struct RecordingProvider {
393        kind: Option<&'static str>,
394        outcome: Outcome,
395        seen: Mutex<Vec<(Option<String>, bool)>>,
396    }
397
398    impl RecordingProvider {
399        fn new(kind: Option<&'static str>, outcome: Outcome) -> Self {
400            Self {
401                kind,
402                outcome,
403                seen: Mutex::new(Vec::new()),
404            }
405        }
406
407        fn seen(&self) -> Vec<(Option<String>, bool)> {
408            self.seen.lock().unwrap().clone()
409        }
410    }
411
412    impl AgentProvider for RecordingProvider {
413        fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
414            Box::pin(async move {
415                let credential = config
416                    .account
417                    .as_ref()
418                    .map(|s| s.credential().expose().to_string());
419                self.seen.lock().unwrap().push((credential, config.verbose));
420                if let Some(session) = &config.account {
421                    session.recorder().record(AccountWindow {
422                        window: "five_hour".to_string(),
423                        utilization: 0.42,
424                        resets_at: Some(Utc::now() + TimeDelta::hours(2)),
425                        status: WindowStatus::Allowed,
426                        model_scope: None,
427                        observed_at: Utc::now(),
428                    });
429                }
430                match self.outcome {
431                    Outcome::Succeed => Ok(AgentOutput::new(json!("done"))),
432                    Outcome::FailApi(status) => Err(AgentError::Api {
433                        status: Some(status),
434                        code: None,
435                        message: "API Error".to_string(),
436                    }),
437                }
438            })
439        }
440
441        fn account_kind(&self) -> Option<&'static str> {
442            self.kind
443        }
444    }
445
446    fn store_with_key() -> InMemoryStore {
447        let mut store = InMemoryStore::new();
448        let spec = format!("1:{}", "aa".repeat(32));
449        store.set_key_ring(KeyRing::from_spec(&spec, Some(1)).unwrap());
450        store
451    }
452
453    async fn add_account(store: &InMemoryStore, name: &str, priority: i32) -> ProviderAccount {
454        let id = Uuid::now_v7();
455        let secret_key = provider_account_secret_key(id);
456        store
457            .set_secret(&secret_key, &format!("{TOKEN}-{name}"))
458            .await
459            .unwrap();
460        store
461            .create_provider_account(NewProviderAccount {
462                id,
463                name: name.to_string(),
464                display_name: name.to_string(),
465                kind: ClaudeSubscriptionKind::ID.to_string(),
466                secret_key,
467                enabled: true,
468                priority,
469                tags: Vec::new(),
470                max_concurrency: None,
471                alert_threshold: 0.8,
472                expires_at: Utc::now() + TimeDelta::days(30),
473                plan: None,
474                created_by: None,
475            })
476            .await
477            .unwrap()
478    }
479
480    fn wrap(inner: Arc<RecordingProvider>, store: &Arc<InMemoryStore>) -> AccountAwareProvider {
481        let store: Arc<dyn Store> = store.clone();
482        AccountAwareProvider::new(inner, store)
483    }
484
485    #[tokio::test]
486    async fn account_aware_provider_injects_selected_account() {
487        let store = Arc::new(store_with_key());
488        add_account(&store, "busy", 10).await;
489        let busy = store
490            .find_provider_account_by_name("busy")
491            .await
492            .unwrap()
493            .unwrap();
494        store
495            .record_provider_account_observation(
496                busy.id,
497                NewProviderAccountObservation {
498                    windows: vec![NewAccountWindow {
499                        window: "five_hour".to_string(),
500                        utilization: 0.9,
501                        resets_at: Some(Utc::now() + TimeDelta::hours(1)),
502                        status: AccountWindowStatus::Allowed,
503                        model_scope: None,
504                        observed_at: Utc::now(),
505                    }],
506                    auth_failed: false,
507                },
508            )
509            .await
510            .unwrap();
511        add_account(&store, "idle", 20).await;
512
513        let inner = Arc::new(RecordingProvider::new(
514            Some(ClaudeSubscriptionKind::ID),
515            Outcome::Succeed,
516        ));
517        let provider = wrap(inner.clone(), &store);
518        provider.invoke(&AgentConfig::new("hello")).await.unwrap();
519
520        let seen = inner.seen();
521        assert_eq!(seen.len(), 1);
522        assert_eq!(seen[0].0.as_deref(), Some(format!("{TOKEN}-idle").as_str()));
523        assert!(seen[0].1, "verbose must be forced for rate_limit_event");
524    }
525
526    #[tokio::test]
527    async fn account_is_injected_behind_a_router() {
528        let store = Arc::new(store_with_key());
529        let account = add_account(&store, "perso", 10).await;
530        let claude = Arc::new(RecordingProvider::new(
531            Some(ClaudeSubscriptionKind::ID),
532            Outcome::Succeed,
533        ));
534        let router = ProviderRouter::new(claude.clone());
535        let dyn_store: Arc<dyn Store> = store.clone();
536        let provider = AccountAwareProvider::new(Arc::new(router), dyn_store);
537
538        let output = provider
539            .invoke(&AgentConfig::new("p").model("sonnet"))
540            .await
541            .unwrap();
542
543        assert_eq!(output.account_id, Some(account.id.to_string()));
544        let seen = claude.seen();
545        assert_eq!(seen.len(), 1);
546        assert_eq!(
547            seen[0].0.as_deref(),
548            Some(format!("{TOKEN}-perso").as_str())
549        );
550    }
551
552    #[tokio::test]
553    async fn mixed_router_injects_an_account_only_on_claude_routes() {
554        let store = Arc::new(store_with_key());
555        let account = add_account(&store, "perso", 10).await;
556        let claude = Arc::new(RecordingProvider::new(
557            Some(ClaudeSubscriptionKind::ID),
558            Outcome::Succeed,
559        ));
560        let http = Arc::new(RecordingProvider::new(None, Outcome::Succeed));
561        let router = ProviderRouter::new(claude.clone())
562            .route(ProviderMatcher::ModelPrefix("gpt-".into()), http.clone());
563        let dyn_store: Arc<dyn Store> = store.clone();
564        let provider = AccountAwareProvider::new(Arc::new(router), dyn_store);
565
566        let claude_output = provider
567            .invoke(&AgentConfig::new("p").model("sonnet"))
568            .await
569            .unwrap();
570        assert_eq!(claude_output.account_id, Some(account.id.to_string()));
571        assert!(claude.seen()[0].0.is_some());
572
573        let http_output = provider
574            .invoke(&AgentConfig::new("p").model("gpt-5"))
575            .await
576            .unwrap();
577        assert_eq!(http_output.account_id, None);
578        assert_eq!(http.seen(), vec![(None, false)]);
579    }
580
581    #[tokio::test]
582    async fn account_aware_provider_passthrough_without_accounts() {
583        let store = Arc::new(store_with_key());
584        let inner = Arc::new(RecordingProvider::new(
585            Some(ClaudeSubscriptionKind::ID),
586            Outcome::Succeed,
587        ));
588        let provider = wrap(inner.clone(), &store);
589        let output = provider.invoke(&AgentConfig::new("hello")).await.unwrap();
590        assert_eq!(output.account_id, None);
591        assert_eq!(inner.seen(), vec![(None, false)]);
592    }
593
594    #[tokio::test]
595    async fn account_aware_provider_passthrough_for_kindless_provider() {
596        let store = Arc::new(store_with_key());
597        add_account(&store, "perso", 10).await;
598        let inner = Arc::new(RecordingProvider::new(None, Outcome::Succeed));
599        let provider = wrap(inner.clone(), &store);
600        let output = provider.invoke(&AgentConfig::new("hello")).await.unwrap();
601        assert_eq!(output.account_id, None);
602        assert_eq!(inner.seen(), vec![(None, false)]);
603        assert_eq!(provider.account_kind(), None);
604    }
605
606    #[tokio::test]
607    async fn account_aware_provider_records_windows_on_error() {
608        let store = Arc::new(store_with_key());
609        let account = add_account(&store, "perso", 10).await;
610        let inner = Arc::new(RecordingProvider::new(
611            Some(ClaudeSubscriptionKind::ID),
612            Outcome::FailApi(500),
613        ));
614        let provider = wrap(inner, &store);
615        let err = provider
616            .invoke(&AgentConfig::new("hello"))
617            .await
618            .unwrap_err();
619        assert!(matches!(
620            err,
621            AgentError::Api {
622                status: Some(500),
623                ..
624            }
625        ));
626
627        let windows = store
628            .list_provider_account_windows(vec![account.id])
629            .await
630            .unwrap();
631        assert_eq!(windows.len(), 1);
632        assert_eq!(windows[0].window, "five_hour");
633        let stored = store
634            .get_provider_account(account.id)
635            .await
636            .unwrap()
637            .unwrap();
638        assert!(stored.auth_failed_at.is_none());
639    }
640
641    #[tokio::test]
642    async fn account_aware_provider_marks_auth_failed_on_401() {
643        let store = Arc::new(store_with_key());
644        let account = add_account(&store, "perso", 10).await;
645        let inner = Arc::new(RecordingProvider::new(
646            Some(ClaudeSubscriptionKind::ID),
647            Outcome::FailApi(401),
648        ));
649        let provider = wrap(inner, &store);
650        provider
651            .invoke(&AgentConfig::new("hello"))
652            .await
653            .unwrap_err();
654
655        let stored = store
656            .get_provider_account(account.id)
657            .await
658            .unwrap()
659            .unwrap();
660        assert!(stored.auth_failed_at.is_some());
661        let candidates = store
662            .list_provider_account_candidates(ClaudeSubscriptionKind::ID.to_string())
663            .await
664            .unwrap();
665        assert!(candidates.is_empty(), "a rejected token is not a candidate");
666    }
667
668    #[tokio::test]
669    async fn account_aware_provider_fails_when_all_exhausted() {
670        let store = Arc::new(store_with_key());
671        let account = add_account(&store, "perso", 10).await;
672        let reset = Utc::now() + TimeDelta::hours(1);
673        store
674            .record_provider_account_observation(
675                account.id,
676                NewProviderAccountObservation {
677                    windows: vec![NewAccountWindow {
678                        window: "five_hour".to_string(),
679                        utilization: 1.0,
680                        resets_at: Some(reset),
681                        status: AccountWindowStatus::Rejected,
682                        model_scope: None,
683                        observed_at: Utc::now(),
684                    }],
685                    auth_failed: false,
686                },
687            )
688            .await
689            .unwrap();
690        let inner = Arc::new(RecordingProvider::new(
691            Some(ClaudeSubscriptionKind::ID),
692            Outcome::Succeed,
693        ));
694        let provider = wrap(inner.clone(), &store);
695        let err = provider
696            .invoke(&AgentConfig::new("hello"))
697            .await
698            .unwrap_err();
699        let AgentError::ProcessFailed { stderr, .. } = err else {
700            panic!("expected ProcessFailed");
701        };
702        assert!(stderr.contains("no provider account available"));
703        assert!(stderr.contains(&reset.to_rfc3339()));
704        assert!(inner.seen().is_empty(), "the agent must not run");
705    }
706
707    #[tokio::test]
708    async fn account_aware_provider_fails_when_credential_missing() {
709        let store = Arc::new(store_with_key());
710        let account = add_account(&store, "perso", 10).await;
711        store.delete_secret(&account.secret_key).await.unwrap();
712        let inner = Arc::new(RecordingProvider::new(
713            Some(ClaudeSubscriptionKind::ID),
714            Outcome::Succeed,
715        ));
716        let provider = wrap(inner, &store);
717        let err = provider
718            .invoke(&AgentConfig::new("hello"))
719            .await
720            .unwrap_err();
721        let message = err.to_string();
722        assert!(message.contains("credential of provider account 'perso' is missing"));
723        assert!(!message.contains(TOKEN));
724    }
725
726    #[tokio::test]
727    async fn account_aware_provider_sets_output_account_id() {
728        let store = Arc::new(store_with_key());
729        let account = add_account(&store, "perso", 10).await;
730        let inner = Arc::new(RecordingProvider::new(
731            Some(ClaudeSubscriptionKind::ID),
732            Outcome::Succeed,
733        ));
734        let provider = wrap(inner, &store);
735        let output = provider.invoke(&AgentConfig::new("hello")).await.unwrap();
736        assert_eq!(output.account_id, Some(account.id.to_string()));
737
738        let windows = store
739            .list_provider_account_windows(vec![account.id])
740            .await
741            .unwrap();
742        assert_eq!(windows.len(), 1);
743        assert!((windows[0].utilization - 0.42).abs() < 1e-9);
744    }
745}