use std::sync::Arc;
use crate::accounts::SelectedSubscriptionAccount;
use crate::app_state::AppState;
use crate::config::UpstreamProvider;
use crate::subscription::{SubscriptionProvider, SubscriptionReader, SubscriptionToken};
use super::{ModelRouteError, available_provider_for_model, credential_state};
#[derive(Debug, Clone)]
pub struct ValidatedSubscription {
pub provider: SubscriptionProvider,
reader: Option<SubscriptionReader>,
selection: CredentialSelection,
requires_live_catalog: bool,
required_model: Option<String>,
}
#[derive(Debug, Clone)]
enum CredentialSelection {
Ready {
cache: Arc<crate::refresh::TokenCache>,
account: String,
baseline: SubscriptionToken,
selected: Box<SelectedSubscriptionAccount>,
},
AccountPool,
}
impl ValidatedSubscription {
pub(crate) async fn for_dispatch(&self) -> Result<SelectedSubscriptionAccount, String> {
let CredentialSelection::Ready {
cache,
account,
baseline,
selected,
} = &self.selection
else {
return Err(format!(
"the {} account pool requires request routing context",
self.provider
));
};
let current = cache
.load_authoritative(self.provider, account)
.await?
.ok_or_else(|| {
format!(
"failed to reload {} credentials from the registered store",
self.provider
)
})?;
if !durably_equivalent(self.provider, ¤t, baseline)
&& !durably_equivalent(self.provider, ¤t, &selected.token)
{
return Err(format!(
"the {} credential changed after its model catalog was validated; retry after discovery completes",
self.provider
));
}
Ok(selected.as_ref().clone())
}
pub(crate) async fn for_dispatch_with_context(
&self,
state: &AppState,
context: &crate::accounts::RoutingContext,
) -> Result<SelectedSubscriptionAccount, String> {
if matches!(self.selection, CredentialSelection::Ready { .. }) {
return self.for_dispatch().await;
}
let router = state
.account_router
.as_ref()
.filter(|router| router.provider() == self.provider)
.ok_or_else(|| format!("no {} account pool is configured", self.provider))?;
let selected = router
.select_subscription_where_authoritative(
context,
&state.subscription_cache,
|account| {
if !self.requires_live_catalog {
return true;
}
let catalog = state.model_catalogs.status_for(self.provider, account);
catalog.discovered
&& catalog.credential_healthy
&& self
.required_model
.as_ref()
.is_none_or(|model| catalog.routable_models().contains(model))
&& state
.subscription_cache
.evidence_for(self.provider, account)
!= Some(crate::refresh::CredentialEvidence::Rejected)
},
)
.await
.map_err(|error| error.to_string())?;
let snapshot =
subscription_snapshot_for_account(state, self.provider, &selected.name, None).await?;
let catalog = state
.model_catalogs
.status_for(self.provider, &selected.name);
if self.requires_live_catalog && (!catalog.discovered || !catalog.credential_healthy) {
return Err(format!(
"the {} model catalog is not currently routable",
self.provider
));
}
let selected_token = snapshot
.selected_token()
.expect("an account snapshot is immediately ready");
if catalog.discovered && !catalog_belongs_to(selected_token, catalog.account.as_deref()) {
return Err(format!(
"the discovered {} catalog belongs to a different account",
self.provider
));
}
snapshot.for_dispatch().await
}
const fn selected_token(&self) -> Option<&SubscriptionToken> {
match &self.selection {
CredentialSelection::Ready { selected, .. } => Some(&selected.token),
CredentialSelection::AccountPool => None,
}
}
const fn uses_account_pool(&self) -> bool {
matches!(self.selection, CredentialSelection::AccountPool)
}
}
fn durably_equivalent(
provider: SubscriptionProvider,
current: &SubscriptionToken,
expected: &SubscriptionToken,
) -> bool {
current == expected
|| (provider == SubscriptionProvider::Codex
&& current.access_token == expected.access_token
&& current.refresh_token == expected.refresh_token
&& current.account_id == expected.account_id
&& current.resource_url == expected.resource_url)
}
pub struct RoutedState {
pub state: AppState,
pub subscription: Option<ValidatedSubscription>,
}
fn catalog_belongs_to(token: &SubscriptionToken, account: Option<&str>) -> bool {
match (token.account_id.as_deref(), account) {
(Some(current), Some(discovered)) => current == discovered,
(None, None) => true,
_ => false,
}
}
async fn subscription_snapshot_for_account(
state: &AppState,
provider: SubscriptionProvider,
account: &str,
reader: Option<SubscriptionReader>,
) -> Result<ValidatedSubscription, String> {
if let Some(reader) = reader.as_ref() {
state.subscription_cache.register_reader(account, reader);
}
let baseline = state
.subscription_cache
.load_authoritative(provider, account)
.await?
.ok_or_else(|| {
format!("failed to load {provider} credentials from the registered store")
})?;
let token = state
.subscription_cache
.get_fresh_loaded(
&state.client,
provider,
account,
baseline.clone(),
chrono::Utc::now().timestamp_millis(),
)
.await?;
if state.subscription_cache.evidence_for(provider, account)
== Some(crate::refresh::CredentialEvidence::Rejected)
{
return Err(format!(
"the {provider} credential was rejected by its upstream"
));
}
Ok(ValidatedSubscription {
provider,
reader,
selection: CredentialSelection::Ready {
cache: Arc::clone(&state.subscription_cache),
account: account.to_string(),
baseline,
selected: Box::new(SelectedSubscriptionAccount {
name: account.to_string(),
token,
}),
},
requires_live_catalog: false,
required_model: None,
})
}
fn account_pool_matches(state: &AppState, provider: SubscriptionProvider) -> bool {
state
.account_router
.as_ref()
.is_some_and(|router| router.provider() == provider)
}
fn local_routing_catalog(
state: &AppState,
) -> (
crate::model_catalog::ModelCatalogCache,
Vec<SubscriptionProvider>,
) {
let catalog = crate::model_catalog::ModelCatalogCache::new();
let mut healthy = Vec::new();
for provider in SubscriptionProvider::ALL {
let accounts = state
.account_router
.as_ref()
.filter(|router| router.provider() == provider)
.map_or_else(
|| {
if state
.subscription_readers
.iter()
.any(|reader| reader.provider() == provider)
{
vec![crate::credential_recovery_store::PRIMARY_ACCOUNT.to_string()]
} else {
Vec::new()
}
},
|router| {
router
.subscription_readers()
.into_iter()
.map(|(account, _)| account)
.collect::<Vec<_>>()
},
);
let mut provider_healthy = false;
for account in accounts {
let status = state.model_catalogs.status_for(provider, &account);
if status.discovered
&& status.credential_healthy
&& state.subscription_cache.evidence_for(provider, &account)
!= Some(crate::refresh::CredentialEvidence::Rejected)
{
provider_healthy = true;
catalog.record_success_for_account(
provider,
&account,
status.account,
status.models,
);
}
}
if provider_healthy {
healthy.push(provider);
}
}
(catalog, healthy)
}
async fn subscription_candidate(
state: &AppState,
provider: SubscriptionProvider,
requires_live_catalog: bool,
required_model: Option<&str>,
) -> Result<ValidatedSubscription, String> {
if account_pool_matches(state, provider) {
return Ok(ValidatedSubscription {
provider,
reader: None,
selection: CredentialSelection::AccountPool,
requires_live_catalog,
required_model: required_model.map(str::to_string),
});
}
let reader = state
.subscription_readers
.iter()
.find(|reader| reader.provider() == provider)
.cloned()
.ok_or_else(|| format!("no {provider} credential reader is configured"))?;
let mut subscription = subscription_snapshot_for_account(
state,
provider,
crate::credential_recovery_store::PRIMARY_ACCOUNT,
Some(reader),
)
.await?;
subscription.requires_live_catalog = requires_live_catalog;
subscription.required_model = required_model.map(str::to_string);
Ok(subscription)
}
async fn validated_catalog_subscription(
state: &AppState,
provider: SubscriptionProvider,
model: &str,
) -> Option<ValidatedSubscription> {
if !account_pool_matches(state, provider) {
let catalog = state.model_catalogs.status(provider);
if !catalog.discovered || !catalog.credential_healthy {
return None;
}
}
let subscription = subscription_candidate(state, provider, true, Some(model))
.await
.ok()?;
let catalog = state.model_catalogs.status(provider);
if subscription
.selected_token()
.is_some_and(|token| !catalog_belongs_to(token, catalog.account.as_deref()))
{
return None;
}
Some(subscription)
}
fn routed_subscription_state(
state: &AppState,
subscription: ValidatedSubscription,
model: Option<&str>,
) -> RoutedState {
let mut routed = state.clone();
routed.upstream_provider = match subscription.provider {
SubscriptionProvider::Claude => UpstreamProvider::Anthropic,
SubscriptionProvider::Codex => UpstreamProvider::Codex,
SubscriptionProvider::Gemini => UpstreamProvider::Gemini,
SubscriptionProvider::Qwen => UpstreamProvider::Qwen,
};
if !subscription.uses_account_pool() {
routed.account_router = None;
}
if let Some(reader) = subscription.reader.clone() {
routed.subscription_reader = Some(reader);
}
if subscription.provider != SubscriptionProvider::Claude
&& let Some(model) = model
{
routed.bridge_model = Some(model.to_string());
}
RoutedState {
state: routed,
subscription: Some(subscription),
}
}
pub async fn route_subscription_model(
state: &AppState,
model: &str,
) -> Result<RoutedState, ModelRouteError> {
let candidates = SubscriptionProvider::ALL
.into_iter()
.filter(|provider| {
state
.model_catalogs
.models(*provider)
.iter()
.any(|candidate| candidate == model)
})
.collect::<Vec<_>>();
let has_catalog_candidate = !candidates.is_empty();
if !has_catalog_candidate {
let (catalog, healthy) = local_routing_catalog(state);
return match available_provider_for_model(model, &healthy, &catalog) {
Err(error) => Err(error),
Ok(_) => unreachable!("a model absent from the complete catalog appeared locally"),
};
}
let relevant = super::provider_for_model(model, &state.model_catalogs)
.map_or(candidates, |provider| vec![provider]);
let validated = futures_util::future::join_all(
relevant
.into_iter()
.map(|provider| validated_catalog_subscription(state, provider, model)),
)
.await
.into_iter()
.flatten()
.collect::<Vec<_>>();
let healthy = validated
.iter()
.map(|subscription| subscription.provider)
.collect::<Vec<_>>();
let provider = available_provider_for_model(model, &healthy, &state.model_catalogs)?;
let subscription = validated
.into_iter()
.find(|subscription| subscription.provider == provider)
.ok_or_else(|| {
let cause = credential_state(provider, &state.model_catalogs)
.unwrap_or_else(|| format!("no usable {provider} credential is available"));
ModelRouteError::NotFound(format!(
"model '{model}' has no healthy {provider} credential: {cause}"
))
})?;
Ok(routed_subscription_state(state, subscription, Some(model)))
}
pub async fn route_pinned_subscription(
state: &AppState,
provider: SubscriptionProvider,
) -> Result<RoutedState, ModelRouteError> {
let catalog = state.model_catalogs.status(provider);
if !account_pool_matches(state, provider)
&& !state
.subscription_readers
.iter()
.any(|reader| reader.provider() == provider)
{
if catalog.discovered && catalog.account.is_some() {
return Err(ModelRouteError::NotFound(format!(
"the discovered {provider} catalog owner cannot be validated without a credential reader"
)));
}
return Ok(RoutedState {
state: state.clone(),
subscription: None,
});
}
let subscription = subscription_candidate(state, provider, false, None)
.await
.map_err(ModelRouteError::NotFound)?;
if catalog.discovered
&& subscription
.selected_token()
.is_some_and(|token| !catalog_belongs_to(token, catalog.account.as_deref()))
{
return Err(ModelRouteError::NotFound(format!(
"the discovered {provider} catalog belongs to a different account"
)));
}
Ok(routed_subscription_state(state, subscription, None))
}