use std::sync::Arc;
use axum::Json;
use axum::extract::{Path, Query, State};
use axum::http::StatusCode;
use axum::response::IntoResponse;
use meerkat_contracts::{
WireAuthProfile, WireAuthProfileCleared, WireAuthProfileCreated, WireAuthProfileDetail,
WireAuthProfilesList, WireAuthStatusDetail, WireBackendProfile, WireBindingIdentity,
WireDeviceStart, WireLoginReady, WireLoginStart, WireOAuthProvider, WireProviderBinding,
WireRealmConnectionSet, WireRealmList, WireRealmSummary,
};
use meerkat_core::connection::{BindingId, ConnectionTargetError, ProfileId, RealmId};
use meerkat_core::handles::LeaseKey;
use meerkat_core::{
AuthBindingRef, CredentialSourceSpec, Provider, RealmConnectionSet, ResolvedConnectionTarget,
};
use meerkat_providers::NormalizedAuthMethod;
use meerkat_providers::auth_oauth::{OAuthError, PkcePair, exchange_authorization_code_with_state};
use meerkat_providers::auth_store::{
CredentialMutationError, CredentialMutationOutcome, PersistedAuthMode, PersistedTokens,
ProviderAuthPersistence, RefreshCoordinator, TokenKey, TokenStore,
credential_source_uses_persisted_store, persisted_auth_mode_is_oauth_login,
};
use meerkat_providers::oauth_flow::{
OAuthDevicePollLease, OAuthFlowError, OAuthProviderIdentity, oauth_provider_resolution,
validate_oauth_login_binding,
};
use crate::AppState;
async fn load_config(state: &AppState) -> Result<meerkat_core::Config, (StatusCode, String)> {
let head = state
.config_runtime
.get()
.await
.map(|snap| snap.config)
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("Failed to load config: {e}"),
)
})?;
meerkat_core::EffectiveConfigReader::new(std::sync::Arc::clone(&state.realm_config_source))
.effective_config_over_head(&state.realm, head)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("Failed to load config: {e}"),
)
})
}
async fn resolve_realm(
state: &AppState,
realm_id: &RealmId,
) -> Result<RealmConnectionSet, (StatusCode, String)> {
let config = load_config(state).await?;
let section = config
.realm
.get(realm_id.as_str())
.ok_or_else(|| (StatusCode::NOT_FOUND, format!("Unknown realm: {realm_id}")))?;
RealmConnectionSet::from_config(realm_id.as_str(), section).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("Realm config invalid: {e}"),
)
})
}
async fn resolve_binding_identity(
state: &AppState,
realm_id: &RealmId,
binding_id: &BindingId,
profile_id: Option<&ProfileId>,
) -> Result<
(
AuthBindingRef,
meerkat_core::ProviderBinding,
meerkat_core::AuthProfile,
),
(StatusCode, String),
> {
let realm = resolve_realm(state, realm_id).await?;
let auth_binding = AuthBindingRef {
realm: realm_id.clone(),
binding: binding_id.clone(),
profile: profile_id.cloned(),
origin: meerkat_core::BindingOrigin::Configured,
};
match realm.lookup_auth_binding(&auth_binding) {
Ok((binding, _, auth_profile)) => Ok((auth_binding, binding.clone(), auth_profile.clone())),
Err(e) => {
let config = load_config(state).await?;
match meerkat_core::connection::resolve_write_owner(&config, realm_id, binding_id) {
Err(inherited @ meerkat_core::connection::WriteOwnerError::Inherited { .. }) => {
Err((StatusCode::CONFLICT, inherited.to_string()))
}
_ => Err((
StatusCode::NOT_FOUND,
format!("Unknown auth identity {realm_id}:{binding_id}: {e}"),
)),
}
}
}
}
async fn resolve_binding_identity_for_read(
state: &AppState,
realm_id: &RealmId,
binding_id: &BindingId,
profile_id: Option<&ProfileId>,
) -> Result<
(
AuthBindingRef,
meerkat_core::ProviderBinding,
meerkat_core::AuthProfile,
),
(StatusCode, String),
> {
let config = load_config(state).await?;
let requested = AuthBindingRef {
realm: realm_id.clone(),
binding: binding_id.clone(),
profile: profile_id.cloned(),
origin: meerkat_core::BindingOrigin::Configured,
};
let target = meerkat_core::resolve_explicit_auth_binding_target(&config, &requested)
.map_err(|error| (target_error_status(&error), error.to_string()))?;
Ok((target.auth_binding, target.binding, target.auth_profile))
}
fn target_error_status(error: &ConnectionTargetError) -> StatusCode {
match error {
ConnectionTargetError::UnknownRealm(_)
| ConnectionTargetError::MissingDefaultBinding { .. }
| ConnectionTargetError::BindingInvalid { .. } => StatusCode::NOT_FOUND,
ConnectionTargetError::RealmConfigInvalid { .. } => StatusCode::INTERNAL_SERVER_ERROR,
ConnectionTargetError::MissingRealm
| ConnectionTargetError::InvalidRealmId { .. }
| ConnectionTargetError::InvalidBindingId { .. }
| ConnectionTargetError::ProviderMismatch { .. }
| ConnectionTargetError::AmbiguousCredentialAccountBindings { .. }
| ConnectionTargetError::RealmChain(_) => StatusCode::BAD_REQUEST,
}
}
fn host_auth_service(state: &AppState) -> meerkat::HostAuthService {
meerkat::HostAuthService::new(
state.provider_auth_persistence.clone(),
state.runtime_adapter.provider_auth_runtime_authority(),
)
}
fn host_auth_error_response(error: meerkat::HostAuthError) -> axum::response::Response {
let status = match &error {
meerkat::HostAuthError::Target(error) => target_error_status(error),
meerkat::HostAuthError::WriteOwner(_)
| meerkat::HostAuthError::OAuthTarget(_)
| meerkat::HostAuthError::BrowserFlowUnsupported(_)
| meerkat::HostAuthError::DeviceFlowUnsupported(_) => StatusCode::BAD_REQUEST,
meerkat::HostAuthError::OAuthFlow(
OAuthFlowError::Missing
| OAuthFlowError::ProviderMismatch { .. }
| OAuthFlowError::RedirectUriMismatch
| OAuthFlowError::TargetMismatch { .. }
| OAuthFlowError::DevicePollInProgress
| OAuthFlowError::DeviceCodeAlreadyAdmitted
| OAuthFlowError::DeviceExpiryOutOfRange,
) => StatusCode::BAD_REQUEST,
meerkat::HostAuthError::OAuthExchange(_) => StatusCode::BAD_GATEWAY,
meerkat::HostAuthError::OAuthFlow(
OAuthFlowError::RegistryProjectionMissing { .. }
| OAuthFlowError::StateGenerationFailed
| OAuthFlowError::LifecycleRejected { .. }
| OAuthFlowError::PersistenceFailed { .. },
)
| meerkat::HostAuthError::CredentialMutation(_)
| meerkat::HostAuthError::TokenStore(_)
| meerkat::HostAuthError::Factory(_)
| meerkat::HostAuthError::PersistenceUnavailable
| meerkat::HostAuthError::Lifecycle(_)
| meerkat::HostAuthError::StatusRehydrate(_) => StatusCode::INTERNAL_SERVER_ERROR,
};
(
status,
Json(serde_json::json!({ "error": error.to_string() })),
)
.into_response()
}
fn require_managed_store_source(
binding_id: &BindingId,
auth_profile: &meerkat_core::AuthProfile,
) -> Result<(), (StatusCode, String)> {
if matches!(&auth_profile.source, CredentialSourceSpec::ManagedStore) {
return Ok(());
}
Err((
StatusCode::BAD_REQUEST,
format!(
"binding {binding_id} resolves auth profile '{}' with source '{}'; \
POST /auth/profiles can only persist credentials for source.kind = 'managed_store'",
auth_profile.id,
source_kind_label(&auth_profile.source),
),
))
}
fn source_kind_label(source: &CredentialSourceSpec) -> &'static str {
match source {
CredentialSourceSpec::InlineSecret { .. } => "inline_secret",
CredentialSourceSpec::ManagedStore => "managed_store",
CredentialSourceSpec::Env { .. } => "env",
CredentialSourceSpec::ExternalResolver { .. } => "external_resolver",
CredentialSourceSpec::PlatformDefault => "platform_default",
CredentialSourceSpec::Command { .. } => "command",
CredentialSourceSpec::FileDescriptor { .. } => "file_descriptor",
}
}
#[cfg(test)]
fn oauth_device_state_error(err: OAuthFlowError) -> (StatusCode, String) {
match err {
OAuthFlowError::Missing => (
StatusCode::BAD_REQUEST,
"oauth device code is missing or expired".to_string(),
),
OAuthFlowError::RegistryProjectionMissing { .. } => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("oauth device registry projection failed: {err}"),
),
other => (
StatusCode::BAD_REQUEST,
format!("oauth device state verification failed: {other}"),
),
}
}
#[cfg(test)]
fn oauth_terminal_device_consume_error(err: OAuthFlowError) -> (StatusCode, String) {
match err {
OAuthFlowError::LifecycleRejected { .. } => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("oauth device terminal consume failed: {err}"),
),
other => oauth_device_state_error(other),
}
}
#[cfg(test)]
fn release_uncredentialed_terminal_oauth_lifecycle(
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
) {
release_uncredentialed_terminal_oauth_lifecycle_for_identity(
auth_lease,
&meerkat_core::AuthCredentialIdentity::from_auth_binding(auth_binding),
);
}
#[cfg(test)]
fn release_uncredentialed_terminal_oauth_lifecycle_for_identity(
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
credential_identity: &meerkat_core::AuthCredentialIdentity,
) {
let lease_key = meerkat_core::handles::LeaseKey::from_credential_identity(credential_identity);
if !auth_lease.snapshot(&lease_key).credential_present {
let _ = auth_lease.release_credential_lifecycle(&lease_key);
}
}
#[cfg(test)]
fn consume_terminal_device_flow_for_identity(
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
credential_identity: &meerkat_core::AuthCredentialIdentity,
poll_lease: OAuthDevicePollLease,
) -> Result<(), (StatusCode, String)> {
poll_lease.consume().map(|_| ()).map_err(|err| {
release_uncredentialed_terminal_oauth_lifecycle_for_identity(
auth_lease,
credential_identity,
);
oauth_terminal_device_consume_error(err)
})
}
#[cfg(test)]
fn consume_terminal_device_flow(
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
poll_lease: OAuthDevicePollLease,
) -> Result<(), (StatusCode, String)> {
poll_lease.consume().map(|_| ()).map_err(|err| {
release_uncredentialed_terminal_oauth_lifecycle(auth_lease, auth_binding);
oauth_terminal_device_consume_error(err)
})
}
#[cfg(test)]
fn finish_device_flow_poll(poll_lease: OAuthDevicePollLease) -> Result<(), (StatusCode, String)> {
poll_lease.finish().map_err(oauth_device_state_error)
}
#[cfg(test)]
fn verify_terminal_device_flow(
poll_lease: &OAuthDevicePollLease,
) -> Result<(), (StatusCode, String)> {
poll_lease
.verify()
.map(|_| ())
.map_err(oauth_device_state_error)
}
#[cfg(test)]
struct PreparedTokenCommitSnapshot {
key: TokenKey,
lease_key: meerkat_core::handles::LeaseKey,
previous: Option<PersistedTokens>,
}
#[cfg(test)]
struct TokenCommitSnapshot {
key: TokenKey,
lease_key: meerkat_core::handles::LeaseKey,
previous: Option<PersistedTokens>,
previous_lifecycle: meerkat_core::handles::AuthLeaseSnapshot,
previous_lifecycle_restore: meerkat_core::handles::AuthLeaseRestoreSnapshot,
lifecycle_transition: meerkat_core::handles::AuthLeaseTransition,
}
#[cfg(test)]
async fn prepare_token_commit_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
) -> Result<PreparedTokenCommitSnapshot, (StatusCode, String)> {
let key = TokenKey::from_auth_binding(auth_binding);
let previous = meerkat_core::rehydrate_durable_predecessor_for_mutation(
token_store,
auth_lease,
auth_binding,
chrono::Utc::now(),
)
.await
.map_err(|error| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("durable credential predecessor rehydrate failed: {error}"),
)
})?;
Ok(PreparedTokenCommitSnapshot {
key,
lease_key: meerkat_core::handles::LeaseKey::from_auth_binding(auth_binding),
previous,
})
}
#[cfg(test)]
async fn save_tokens_and_publish_lifecycle_commit_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
tokens: &PersistedTokens,
) -> Result<TokenCommitSnapshot, (StatusCode, String)> {
let key = TokenKey::from_auth_binding(auth_binding);
let lease_key = meerkat_core::handles::LeaseKey::from_auth_binding(auth_binding);
let previous = meerkat_core::rehydrate_durable_predecessor_for_mutation(
token_store,
auth_lease,
auth_binding,
chrono::Utc::now(),
)
.await
.map_err(|error| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("durable credential predecessor rehydrate failed: {error}"),
)
})?;
let previous_lifecycle_restore = auth_lease.capture_auth_lifecycle_restore_snapshot(&lease_key);
let previous_lifecycle = previous_lifecycle_restore.snapshot().clone();
let transition =
meerkat_core::publish_token_lifecycle_acquired(auth_lease, auth_binding, tokens).map_err(
|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("AuthMachine lifecycle acquire failed: {e}"),
)
},
)?;
let commit = TokenCommitSnapshot {
key,
lease_key,
previous,
previous_lifecycle,
previous_lifecycle_restore,
lifecycle_transition: transition,
};
save_marked_token_commit_unlocked(token_store, auth_lease, &commit, tokens, "").await?;
Ok(commit)
}
#[cfg(test)]
async fn save_marked_token_commit_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
commit: &TokenCommitSnapshot,
tokens: &PersistedTokens,
failure_context: &str,
) -> Result<(), (StatusCode, String)> {
let committed_tokens = match meerkat_core::mark_tokens_lifecycle_published_for_transition(
&commit.key,
tokens,
&commit.lifecycle_transition,
) {
Ok(committed_tokens) => committed_tokens,
Err(e) => {
let message = match rollback_token_commit(token_store, auth_lease, commit).await {
Ok(()) => format!(
"AuthMachine lifecycle marker handoff failed{failure_context}: {e}; acquired lease rolled back"
),
Err(rollback_error) => format!(
"AuthMachine lifecycle marker handoff failed{failure_context}: {e}; acquired lease rollback failed: {rollback_error}"
),
};
return Err((StatusCode::INTERNAL_SERVER_ERROR, message));
}
};
if let Err(e) = token_store.save(&commit.key, &committed_tokens).await {
let message = match rollback_token_commit(token_store, auth_lease, commit).await {
Ok(()) => {
format!("TokenStore save failed{failure_context}: {e}; acquired lease rolled back")
}
Err(rollback_error) => {
format!(
"TokenStore save failed{failure_context}: {e}; acquired lease rollback failed: {rollback_error}"
)
}
};
return Err((StatusCode::INTERNAL_SERVER_ERROR, message));
}
Ok(())
}
#[cfg(test)]
async fn save_tokens_and_publish_lifecycle(
persistence: ProviderAuthPersistence,
auth_lease: meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: AuthBindingRef,
tokens: PersistedTokens,
) -> Result<(), (StatusCode, String)> {
let key = TokenKey::from_auth_binding(&auth_binding);
let token_store = persistence.token_store();
let refresh_coordinator = persistence.refresh_coordinator();
let mutation_store = Arc::clone(&token_store);
let outcome = refresh_coordinator
.with_exclusive_mutation(
key,
Box::new(move || {
Box::pin(async move {
let lease_key =
meerkat_core::handles::LeaseKey::from_auth_binding(&auth_binding);
let _guard = meerkat_core::acquire_auth_login_lifecycle_guard(&lease_key).await;
let commit = save_tokens_and_publish_lifecycle_commit_unlocked(
mutation_store.as_ref(),
&auth_lease,
&auth_binding,
&tokens,
)
.await
.map_err(|(_, message)| CredentialMutationError::Operation(message))?;
let committed = meerkat_core::mark_tokens_lifecycle_published_for_transition(
&commit.key,
&tokens,
&commit.lifecycle_transition,
)
.map_err(|error| CredentialMutationError::AuthLifecycle(error.to_string()))?;
Ok(CredentialMutationOutcome::Persisted(committed))
})
}),
)
.await
.map_err(|error| (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()))?;
match outcome {
CredentialMutationOutcome::Persisted(_) => Ok(()),
CredentialMutationOutcome::Cleared => Err((
StatusCode::INTERNAL_SERVER_ERROR,
"credential save transaction returned cleared outcome".to_string(),
)),
}
}
#[cfg(test)]
async fn rollback_token_commit(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
commit: &TokenCommitSnapshot,
) -> Result<(), String> {
match &commit.previous {
Some(previous) => match commit.previous_lifecycle.phase {
Some(phase) if phase != meerkat_core::handles::AuthLeasePhase::Released => {
auth_lease
.release_credential_lifecycle(&commit.lease_key)
.map_err(|e| format!("AuthMachine lifecycle rollback release failed: {e}"))?;
token_store
.save(&commit.key, previous)
.await
.map_err(|e| format!("TokenStore rollback save failed: {e}"))?;
let restored_transition = meerkat_core::restore_token_lifecycle_snapshot(
auth_lease,
&commit.previous_lifecycle_restore,
)
.map_err(|e| format!("AuthMachine lifecycle rollback failed: {e}"))?;
if let Some(restored_transition) = restored_transition {
let restored_previous =
meerkat_core::mark_tokens_lifecycle_published_for_transition(
&commit.key,
previous,
&restored_transition,
)
.map_err(|e| format!("AuthMachine rollback marker handoff failed: {e}"))?;
token_store
.save(&commit.key, &restored_previous)
.await
.map_err(|e| format!("TokenStore rollback marker save failed: {e}"))?;
}
}
_ => {
auth_lease
.release_credential_lifecycle(&commit.lease_key)
.map_err(|e| format!("AuthMachine lifecycle rollback release failed: {e}"))?;
token_store
.save(&commit.key, previous)
.await
.map_err(|e| format!("TokenStore rollback save failed: {e}"))?;
}
},
None => {
auth_lease
.release_credential_lifecycle(&commit.lease_key)
.map_err(|e| format!("AuthMachine lifecycle rollback release failed: {e}"))?;
token_store
.clear(&commit.key)
.await
.map_err(|e| format!("TokenStore rollback clear failed: {e}"))?;
}
}
Ok(())
}
#[cfg(test)]
async fn save_prepared_tokens_after_terminal_consume_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
tokens: &PersistedTokens,
prepared: PreparedTokenCommitSnapshot,
) -> Result<(), (StatusCode, String)> {
let previous_lifecycle_restore =
auth_lease.capture_auth_lifecycle_restore_snapshot(&prepared.lease_key);
let previous_lifecycle = previous_lifecycle_restore.snapshot().clone();
let transition =
meerkat_core::publish_token_lifecycle_acquired(auth_lease, auth_binding, tokens).map_err(
|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("AuthMachine lifecycle acquire failed after OAuth consume: {e}"),
)
},
)?;
let commit = TokenCommitSnapshot {
key: prepared.key,
lease_key: prepared.lease_key,
previous: prepared.previous,
previous_lifecycle,
previous_lifecycle_restore,
lifecycle_transition: transition,
};
save_marked_token_commit_unlocked(
token_store,
auth_lease,
&commit,
tokens,
" after OAuth consume",
)
.await
}
#[cfg(test)]
async fn save_tokens_and_consume_device_flow_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
tokens: &PersistedTokens,
poll_lease: OAuthDevicePollLease,
) -> Result<(), (StatusCode, String)> {
if !poll_lease.terminal_flow_state_is_authmachine_owned() {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
"consume_oauth_device_flow requires an AuthMachine-owned OAuth device poll lease"
.to_string(),
));
}
verify_terminal_device_flow(&poll_lease)?;
let prepared = prepare_token_commit_unlocked(token_store, auth_lease, auth_binding).await?;
consume_terminal_device_flow(auth_lease, auth_binding, poll_lease)?;
save_prepared_tokens_after_terminal_consume_unlocked(
token_store,
auth_lease,
auth_binding,
tokens,
prepared,
)
.await?;
Ok(())
}
#[cfg(test)]
async fn save_tokens_and_consume_device_flow(
persistence: ProviderAuthPersistence,
auth_lease: meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: AuthBindingRef,
tokens: PersistedTokens,
poll_lease: OAuthDevicePollLease,
) -> Result<(), (StatusCode, String)> {
let key = TokenKey::from_auth_binding(&auth_binding);
let token_store = persistence.token_store();
let refresh_coordinator = persistence.refresh_coordinator();
let load_key = key.clone();
let outcome = refresh_coordinator
.with_exclusive_mutation(
key,
Box::new(move || {
Box::pin(async move {
let lease_key =
meerkat_core::handles::LeaseKey::from_auth_binding(&auth_binding);
let _guard = meerkat_core::acquire_auth_login_lifecycle_guard(&lease_key).await;
save_tokens_and_consume_device_flow_unlocked(
token_store.as_ref(),
&auth_lease,
&auth_binding,
&tokens,
poll_lease,
)
.await
.map_err(|(_, message)| CredentialMutationError::Operation(message))?;
let committed = token_store
.load(&load_key)
.await
.map_err(|error| CredentialMutationError::TokenStore(error.to_string()))?
.ok_or_else(|| {
CredentialMutationError::TokenStore(
"successful device login left no persisted credential".to_string(),
)
})?;
Ok(CredentialMutationOutcome::Persisted(committed))
})
}),
)
.await
.map_err(|error| (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()))?;
match outcome {
CredentialMutationOutcome::Persisted(_) => Ok(()),
CredentialMutationOutcome::Cleared => Err((
StatusCode::INTERNAL_SERVER_ERROR,
"device-login transaction returned cleared outcome".to_string(),
)),
}
}
#[cfg(test)]
struct BrowserFlowConsume<'a> {
authority: &'a dyn meerkat_providers::oauth_flow::OAuthFlowAuthority,
state: &'a str,
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
redirect_uri: &'a str,
}
#[cfg(test)]
struct OwnedBrowserFlowConsume {
authority: Arc<dyn meerkat_providers::oauth_flow::OAuthFlowAuthority>,
state: String,
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
redirect_uri: String,
}
#[cfg(test)]
async fn save_tokens_and_consume_browser_flow_unlocked(
token_store: &dyn TokenStore,
auth_lease: &meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: &AuthBindingRef,
tokens: &PersistedTokens,
flow: BrowserFlowConsume<'_>,
) -> Result<(), (StatusCode, String)> {
if !flow.authority.terminal_flow_state_is_authmachine_owned() {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
"consume_oauth_browser_flow requires an AuthMachine-owned OAuth flow authority"
.to_string(),
));
}
let prepared = prepare_token_commit_unlocked(token_store, auth_lease, auth_binding).await?;
flow.authority
.consume(
flow.state,
&meerkat_core::AuthCredentialIdentity::from_auth_binding(auth_binding),
flow.provider,
flow.redirect_uri,
)
.map_err(|err| {
release_uncredentialed_terminal_oauth_lifecycle(auth_lease, auth_binding);
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("oauth state terminal consume failed: {err}"),
)
})?;
save_prepared_tokens_after_terminal_consume_unlocked(
token_store,
auth_lease,
auth_binding,
tokens,
prepared,
)
.await?;
Ok(())
}
#[cfg(test)]
async fn save_tokens_and_consume_browser_flow(
persistence: ProviderAuthPersistence,
auth_lease: meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: AuthBindingRef,
tokens: PersistedTokens,
flow: OwnedBrowserFlowConsume,
) -> Result<(), (StatusCode, String)> {
let key = TokenKey::from_auth_binding(&auth_binding);
let token_store = persistence.token_store();
let refresh_coordinator = persistence.refresh_coordinator();
let load_key = key.clone();
let outcome = refresh_coordinator
.with_exclusive_mutation(
key,
Box::new(move || {
Box::pin(async move {
let lease_key =
meerkat_core::handles::LeaseKey::from_auth_binding(&auth_binding);
let _guard = meerkat_core::acquire_auth_login_lifecycle_guard(&lease_key).await;
save_tokens_and_consume_browser_flow_unlocked(
token_store.as_ref(),
&auth_lease,
&auth_binding,
&tokens,
BrowserFlowConsume {
authority: flow.authority.as_ref(),
state: &flow.state,
provider: flow.provider,
redirect_uri: &flow.redirect_uri,
},
)
.await
.map_err(|(_, message)| CredentialMutationError::Operation(message))?;
let committed = token_store
.load(&load_key)
.await
.map_err(|error| CredentialMutationError::TokenStore(error.to_string()))?
.ok_or_else(|| {
CredentialMutationError::TokenStore(
"successful browser login left no persisted credential".to_string(),
)
})?;
Ok(CredentialMutationOutcome::Persisted(committed))
})
}),
)
.await
.map_err(|error| (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()))?;
match outcome {
CredentialMutationOutcome::Persisted(_) => Ok(()),
CredentialMutationOutcome::Cleared => Err((
StatusCode::INTERNAL_SERVER_ERROR,
"browser-login transaction returned cleared outcome".to_string(),
)),
}
}
#[cfg(test)]
async fn clear_tokens_and_publish_lifecycle(
persistence: ProviderAuthPersistence,
auth_lease: meerkat_core::handles::GeneratedAuthLeaseHandle,
auth_binding: AuthBindingRef,
) -> Result<(), (StatusCode, String)> {
meerkat_core::clear_tokens_and_publish_lifecycle_released_coordinated(
persistence,
auth_lease,
auth_binding,
)
.await
.map_err(|error| (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()))
}
pub async fn list_realms(State(state): State<AppState>) -> impl IntoResponse {
match load_config(&state).await {
Ok(config) => {
let realms: Vec<WireRealmSummary> = config
.realm
.iter()
.map(|(realm_id, section)| WireRealmSummary {
realm_id: realm_id.clone(),
default_binding: section.default_binding.clone(),
backend_count: section.backend.len(),
auth_profile_count: section.auth.len(),
binding_count: section.binding.len(),
})
.collect();
(StatusCode::OK, Json(WireRealmList { realms })).into_response()
}
Err((status, msg)) => (status, Json(serde_json::json!({ "error": msg }))).into_response(),
}
}
pub async fn get_realm(
State(state): State<AppState>,
Path(realm_id): Path<RealmId>,
) -> impl IntoResponse {
match resolve_realm(&state, &realm_id).await {
Ok(realm) => {
let wire = WireRealmConnectionSet::from(&realm);
(StatusCode::OK, Json(wire)).into_response()
}
Err((status, msg)) => (status, Json(serde_json::json!({ "error": msg }))).into_response(),
}
}
#[derive(serde::Deserialize)]
pub struct RealmQuery {
pub realm_id: RealmId,
#[serde(default)]
pub profile_id: Option<ProfileId>,
}
pub async fn list_auth_profiles(
State(state): State<AppState>,
Query(query): Query<RealmQuery>,
) -> impl IntoResponse {
match resolve_realm(&state, &query.realm_id).await {
Ok(realm) => {
let profiles: Vec<WireAuthProfile> = realm
.auth_profiles
.values()
.map(WireAuthProfile::from)
.collect();
let backends: Vec<WireBackendProfile> = realm
.backends
.values()
.map(WireBackendProfile::from)
.collect();
let bindings: Vec<WireProviderBinding> = realm
.bindings
.values()
.map(WireProviderBinding::from)
.collect();
(
StatusCode::OK,
Json(WireAuthProfilesList {
realm_id: realm.realm_id.to_string(),
auth_profiles: profiles,
backend_profiles: backends,
bindings,
}),
)
.into_response()
}
Err((status, msg)) => (status, Json(serde_json::json!({ "error": msg }))).into_response(),
}
}
pub use meerkat_contracts::wire::RestAuthProfileCreateRequest as CreateAuthProfileBody;
pub async fn create_auth_profile(
State(state): State<AppState>,
Json(body): Json<CreateAuthProfileBody>,
) -> impl IntoResponse {
let (auth_binding, binding, auth_profile) = match resolve_binding_identity(
&state,
&body.realm_id,
&body.binding_id,
body.profile_id.as_ref(),
)
.await
{
Ok(v) => v,
Err((status, msg)) => {
return (status, Json(serde_json::json!({ "error": msg }))).into_response();
}
};
if meerkat_core::Provider::parse_strict(&body.provider) != Some(auth_profile.provider) {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({
"error": format!(
"binding {} resolves provider '{}' not '{}'",
body.binding_id,
auth_profile.provider.as_str(),
body.provider,
),
})),
)
.into_response();
}
if body.auth_method != auth_profile.auth_method {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({
"error": format!(
"binding {} resolves auth_method '{}' not '{}'",
body.binding_id,
auth_profile.auth_method,
body.auth_method,
),
})),
)
.into_response();
}
if let Err((status, msg)) = require_managed_store_source(&body.binding_id, &auth_profile) {
return (status, Json(serde_json::json!({ "error": msg }))).into_response();
}
let auth_mode = match NormalizedAuthMethod::from_auth_profile(&auth_profile)
.and_then(NormalizedAuthMethod::persisted_auth_mode)
{
Some(mode) if meerkat_core::persisted_auth_mode_is_directly_creatable(mode) => mode,
_ => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({
"error": format!(
"auth_method '{}' cannot be created via the REST endpoint. \
OAuth methods use /auth/login/start + /auth/login/complete; \
managed_store and external_resolver are configured via TOML.",
auth_profile.auth_method,
),
})),
)
.into_response();
}
};
let tokens = PersistedTokens {
auth_mode,
primary_secret: Some(body.secret),
refresh_token: None,
id_token: None,
expires_at: None,
last_refresh: Some(chrono::Utc::now()),
scopes: Vec::new(),
account_id: None,
metadata: serde_json::Value::Null,
};
if let Err(error) = meerkat_providers::browser_login::save_tokens_and_publish_lifecycle(
state.provider_auth_persistence.clone(),
state.auth_lease.clone(),
binding.credential_identity(&auth_binding),
tokens.clone(),
)
.await
{
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": error.to_string() })),
)
.into_response();
}
tracing::info!(
target: "meerkat::auth::audit",
binding_key = ?auth_binding,
action = "create_profile",
provider = %auth_profile.provider.as_str(),
auth_method = %auth_profile.auth_method,
"binding-scoped auth credentials stored via REST"
);
(
StatusCode::CREATED,
Json(WireAuthProfileCreated {
identity: WireBindingIdentity::from(&auth_binding),
profile_id: auth_profile.id.clone(),
provider: auth_profile.provider.as_str().to_string(),
auth_method: auth_profile.auth_method.clone(),
stored: true,
}),
)
.into_response()
}
pub async fn get_auth_profile(
State(state): State<AppState>,
Path(binding_id): Path<BindingId>,
Query(query): Query<RealmQuery>,
) -> impl IntoResponse {
match resolve_binding_identity(
&state,
&query.realm_id,
&binding_id,
query.profile_id.as_ref(),
)
.await
{
Ok((auth_binding, binding, auth_profile)) => (
StatusCode::OK,
Json(WireAuthProfileDetail {
auth_binding: auth_binding.into(),
binding_id: binding.id.clone(),
profile_id: auth_profile.id.clone(),
auth_profile: WireAuthProfile::from(&auth_profile),
}),
)
.into_response(),
Err((status, msg)) => (status, Json(serde_json::json!({ "error": msg }))).into_response(),
}
}
pub async fn delete_auth_profile(
State(state): State<AppState>,
Path(binding_id): Path<BindingId>,
Query(query): Query<RealmQuery>,
) -> impl IntoResponse {
let (auth_binding, binding, auth_profile) = match resolve_binding_identity(
&state,
&query.realm_id,
&binding_id,
query.profile_id.as_ref(),
)
.await
{
Ok(v) => v,
Err((status, msg)) => {
return (status, Json(serde_json::json!({ "error": msg }))).into_response();
}
};
if let Err(error) =
meerkat_core::clear_tokens_and_publish_lifecycle_released_coordinated_for_identity(
state.provider_auth_persistence.clone(),
state.auth_lease.clone(),
binding.credential_identity(&auth_binding),
)
.await
{
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": error.to_string() })),
)
.into_response();
}
tracing::info!(
target: "meerkat::auth::audit",
binding_key = ?auth_binding,
action = "delete_profile",
"binding-scoped auth credentials deleted via REST"
);
(
StatusCode::NO_CONTENT,
Json(WireAuthProfileCleared {
identity: WireBindingIdentity::from(&auth_binding),
profile_id: auth_profile.id.clone(),
cleared: true,
}),
)
.into_response()
}
pub use meerkat_contracts::wire::RestAuthBindingTestRequest as TestBindingBody;
pub async fn test_auth_binding(
State(state): State<AppState>,
Path(binding_id): Path<BindingId>,
Json(body): Json<TestBindingBody>,
) -> impl IntoResponse {
match resolve_realm(&state, &body.realm_id).await {
Ok(realm) => {
let env = meerkat_providers::ResolverEnvironment::with_process_env()
.with_provider_auth_persistence(state.provider_auth_persistence.clone())
.with_auth_lease_handle(state.auth_lease.clone());
let auth_binding = AuthBindingRef {
realm: body.realm_id.clone(),
binding: binding_id.clone(),
profile: body.profile_id.clone(),
origin: meerkat_core::BindingOrigin::Configured,
};
match state
.provider_registry
.resolve(&realm, &auth_binding, &env)
.await
{
Ok(conn) => {
let (auth_binding, binding, auth_profile) = match resolve_binding_identity(
&state,
&body.realm_id,
&binding_id,
body.profile_id.as_ref(),
)
.await
{
Ok(v) => v,
Err((status, msg)) => {
return (status, Json(serde_json::json!({ "error": msg })))
.into_response();
}
};
(
StatusCode::OK,
Json(serde_json::json!({
"state": "valid",
"auth_binding": &auth_binding,
"binding_id": &binding.id,
"profile_id": &auth_profile.id,
"provider": conn.provider.as_str(),
"backend_profile_id": &conn.backend_profile.id,
"has_credential": conn.resolved_secret().is_some()
|| conn.resolved_authorizer().is_some(),
})),
)
.into_response()
}
Err(e) => (
StatusCode::UNPROCESSABLE_ENTITY,
Json(serde_json::json!({
"state": "error",
"error": format!("Binding resolution failed: {e}"),
})),
)
.into_response(),
}
}
Err((status, msg)) => (status, Json(serde_json::json!({ "error": msg }))).into_response(),
}
}
fn parse_auth_identity_triple(
realm_id: &str,
binding_id: &str,
profile_id: Option<&str>,
) -> Result<(RealmId, BindingId, Option<ProfileId>), String> {
let realm_id = RealmId::parse(realm_id).map_err(|e| format!("invalid realm_id: {e}"))?;
let binding_id =
BindingId::parse(binding_id).map_err(|e| format!("invalid binding_id: {e}"))?;
let profile_id = profile_id
.map(|p| ProfileId::parse(p).map_err(|e| format!("invalid profile_id: {e}")))
.transpose()?;
Ok((realm_id, binding_id, profile_id))
}
pub use meerkat_contracts::wire::LoginStartParams as LoginStartBody;
pub async fn start_login(
State(state): State<AppState>,
Json(body): Json<LoginStartBody>,
) -> impl IntoResponse {
let (realm_id, binding_id, profile_id) = match parse_auth_identity_triple(
&body.realm_id,
&body.binding_id,
body.profile_id.as_deref(),
) {
Ok(v) => v,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "error": e })),
)
.into_response();
}
};
let target = meerkat::HostAuthTarget {
provider: body.provider.identity(),
realm_id,
binding_id,
profile_id,
};
let config = match load_config(&state).await {
Ok(config) => config,
Err((status, message)) => {
return (status, Json(serde_json::json!({ "error": message }))).into_response();
}
};
match host_auth_service(&state)
.login_start(&config, &target, body.redirect_uri)
.await
{
Ok(started) => (
StatusCode::OK,
Json(WireLoginStart {
authorize_url: started.authorize_url,
state: started.state,
redirect_uri: started.redirect_uri,
provider: body.provider,
}),
)
.into_response(),
Err(error) => host_auth_error_response(error),
}
}
pub use meerkat_contracts::wire::LoginCompleteParams as LoginCompleteBody;
pub async fn complete_login(
State(state): State<AppState>,
Json(body): Json<LoginCompleteBody>,
) -> impl IntoResponse {
let (realm_id, binding_id, profile_id) = match parse_auth_identity_triple(
&body.realm_id,
&body.binding_id,
body.profile_id.as_deref(),
) {
Ok(v) => v,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "error": e })),
)
.into_response();
}
};
let target = meerkat::HostAuthTarget {
provider: body.provider.identity(),
realm_id,
binding_id,
profile_id,
};
let config = match load_config(&state).await {
Ok(config) => config,
Err((status, message)) => {
return (status, Json(serde_json::json!({ "error": message }))).into_response();
}
};
match host_auth_service(&state)
.login_complete(&config, &target, body.redirect_uri, body.state, &body.code)
.await
{
Ok(completed) => {
tracing::info!(
target: "meerkat::auth::audit",
binding_key = ?completed.auth_binding,
action = "login_oauth_complete",
provider = %body.provider,
has_refresh_token = %completed.has_refresh_token,
"OAuth login completed via REST"
);
(
StatusCode::OK,
Json(WireLoginReady {
state: None,
identity: WireBindingIdentity::from(&completed.auth_binding),
profile_id: completed.profile_id,
provider: body.provider,
expires_at: completed.expires_at.map(|expires| expires.to_rfc3339()),
has_refresh_token: completed.has_refresh_token,
scopes: completed.scopes,
}),
)
.into_response()
}
Err(error) => host_auth_error_response(error),
}
}
pub use meerkat_contracts::wire::DeviceStartParams as DeviceStartBody;
pub async fn start_device_login(
State(state): State<AppState>,
Json(body): Json<DeviceStartBody>,
) -> impl IntoResponse {
let (realm_id, binding_id, profile_id) = match parse_auth_identity_triple(
&body.realm_id,
&body.binding_id,
body.profile_id.as_deref(),
) {
Ok(v) => v,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "error": e })),
)
.into_response();
}
};
let target = meerkat::HostAuthTarget {
provider: body.provider.identity(),
realm_id,
binding_id,
profile_id,
};
let config = match load_config(&state).await {
Ok(config) => config,
Err((status, message)) => {
return (status, Json(serde_json::json!({ "error": message }))).into_response();
}
};
match host_auth_service(&state)
.device_start(&config, &target)
.await
{
Ok(device) => (
StatusCode::OK,
Json(WireDeviceStart {
device_code: device.device_code,
user_code: device.user_code,
verification_uri: device.verification_uri,
verification_uri_complete: device.verification_uri_complete,
expires_in: device.expires_in,
interval: device.interval,
provider: body.provider,
}),
)
.into_response(),
Err(error) => host_auth_error_response(error),
}
}
pub use meerkat_contracts::wire::DeviceCompleteParams as DeviceCompleteBody;
pub async fn complete_device_login(
State(state): State<AppState>,
Json(body): Json<DeviceCompleteBody>,
) -> impl IntoResponse {
let (realm_id, binding_id, profile_id) = match parse_auth_identity_triple(
&body.realm_id,
&body.binding_id,
body.profile_id.as_deref(),
) {
Ok(v) => v,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "error": e })),
)
.into_response();
}
};
let target = meerkat::HostAuthTarget {
provider: body.provider.identity(),
realm_id,
binding_id,
profile_id,
};
let config = match load_config(&state).await {
Ok(config) => config,
Err((status, message)) => {
return (status, Json(serde_json::json!({ "error": message }))).into_response();
}
};
match host_auth_service(&state)
.device_poll(&config, &target, &body.device_code)
.await
{
Ok(meerkat::HostAuthDevicePoll::Pending) => (
StatusCode::ACCEPTED,
Json(serde_json::json!({ "state": "pending" })),
)
.into_response(),
Ok(meerkat::HostAuthDevicePoll::SlowDown) => (
StatusCode::TOO_MANY_REQUESTS,
Json(serde_json::json!({ "state": "slow_down" })),
)
.into_response(),
Ok(meerkat::HostAuthDevicePoll::AccessDenied) => (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "state": "access_denied" })),
)
.into_response(),
Ok(meerkat::HostAuthDevicePoll::Expired) => (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "state": "expired" })),
)
.into_response(),
Ok(meerkat::HostAuthDevicePoll::Ready(completed)) => {
tracing::info!(
target: "meerkat::auth::audit",
binding_key = ?completed.auth_binding,
action = "login_device_complete",
provider = %body.provider,
has_refresh_token = %completed.has_refresh_token,
"OAuth device-flow login completed via REST"
);
(
StatusCode::OK,
Json(WireLoginReady {
state: Some("ready".to_string()),
identity: WireBindingIdentity::from(&completed.auth_binding),
profile_id: completed.profile_id,
provider: body.provider,
expires_at: completed.expires_at.map(|expires| expires.to_rfc3339()),
has_refresh_token: completed.has_refresh_token,
scopes: completed.scopes,
}),
)
.into_response()
}
Err(error) => host_auth_error_response(error),
}
}
pub async fn get_auth_status(
State(state): State<AppState>,
Path(binding_id): Path<BindingId>,
Query(query): Query<RealmQuery>,
) -> impl IntoResponse {
let (auth_binding, binding, auth_profile) = match resolve_binding_identity_for_read(
&state,
&query.realm_id,
&binding_id,
query.profile_id.as_ref(),
)
.await
{
Ok(v) => v,
Err((status, msg)) => {
return (status, Json(serde_json::json!({ "error": msg }))).into_response();
}
};
let credential_identity = binding.credential_identity(&auth_binding);
let lease_key = LeaseKey::from_credential_identity(&credential_identity);
let now = chrono::Utc::now();
if let Err(err) = state.auth_lease.observe_credential_freshness(
&lease_key,
now.timestamp().max(0) as u64,
meerkat_core::handles::AUTH_LEASE_TTL_REFRESH_WINDOW_SECS,
) {
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": format!("AuthMachine freshness observation failed: {err}")
})),
)
.into_response();
}
let mut snapshot = state.auth_lease.snapshot(&lease_key);
let expected_mode = NormalizedAuthMethod::from_auth_profile(&auth_profile)
.and_then(NormalizedAuthMethod::persisted_auth_mode);
let source_uses_store = credential_source_uses_persisted_store(&auth_profile.source);
let oauth_mode = expected_mode
.map(persisted_auth_mode_is_oauth_login)
.unwrap_or(false);
let phase = meerkat_core::AuthStatusPhase::from_lease_snapshot(now, &snapshot);
let mut stored = None;
if source_uses_store {
if phase.is_no_live_lease() {
if let Some(expected_mode) = expected_mode {
match meerkat_core::rehydrate_marked_tokens_for_status_for_identity(
state.token_store().as_ref(),
&state.auth_lease,
&credential_identity,
expected_mode,
now,
)
.await
{
Ok(Some(rehydrated)) => {
stored = Some(rehydrated);
snapshot = state.auth_lease.snapshot(&lease_key);
}
Ok(None) => {}
Err(err) => {
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": format!("TokenStore rehydration failed: {err}")
})),
)
.into_response();
}
}
}
} else {
stored = match state
.token_store()
.load(&TokenKey::from_credential_identity(&credential_identity))
.await
{
Ok(stored) => stored,
Err(err) => {
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": format!("TokenStore load failed: {err}")
})),
)
.into_response();
}
};
}
}
if stored
.as_ref()
.is_some_and(|tokens| Some(tokens.auth_mode) != expected_mode)
{
stored = None;
}
let oauth_source_rejected = expected_mode
.map(|mode| persisted_auth_mode_is_oauth_login(mode) && !source_uses_store)
.unwrap_or(false);
let token_matches_binding = if source_uses_store {
match stored.as_ref() {
Some(tokens) => Some(tokens.auth_mode) == expected_mode,
None => !oauth_mode,
}
} else {
true
};
let marker_projection_snapshot;
let (projection_tokens, projection_snapshot) = if token_matches_binding {
marker_projection_snapshot = stored.as_ref().filter(|_| oauth_mode).and_then(|tokens| {
meerkat_core::oauth_status_projection_snapshot_from_newer_marker(&snapshot, tokens)
});
(
if source_uses_store && !oauth_source_rejected {
stored.as_ref()
} else {
None
},
marker_projection_snapshot.as_ref().unwrap_or(&snapshot),
)
} else {
(None, &snapshot)
};
let projection =
meerkat_core::project_published_auth_status(now, projection_tokens, projection_snapshot);
let tokens = projection.tokens;
(
StatusCode::OK,
Json(WireAuthStatusDetail {
identity: WireBindingIdentity::from(&auth_binding),
profile_id: auth_profile.id.clone(),
provider: auth_profile.provider.as_str().to_string(),
auth_method: auth_profile.auth_method.clone(),
state: projection.phase,
expires_at: projection.expires_at.map(|e| e.to_rfc3339()),
last_refresh_at: tokens.and_then(|t| t.last_refresh.map(|e| e.to_rfc3339())),
account_id: tokens.and_then(|t| t.account_id.clone()),
has_refresh_token: tokens.map(|t| t.refresh_token.is_some()).unwrap_or(false),
}),
)
.into_response()
}
pub async fn logout(
State(state): State<AppState>,
Path(binding_id): Path<BindingId>,
Query(query): Query<RealmQuery>,
) -> impl IntoResponse {
let (auth_binding, binding, auth_profile) = match resolve_binding_identity(
&state,
&query.realm_id,
&binding_id,
query.profile_id.as_ref(),
)
.await
{
Ok(v) => v,
Err((status, msg)) => {
return (status, Json(serde_json::json!({ "error": msg }))).into_response();
}
};
if let Err(error) =
meerkat_core::clear_tokens_and_publish_lifecycle_released_coordinated_for_identity(
state.provider_auth_persistence.clone(),
state.auth_lease.clone(),
binding.credential_identity(&auth_binding),
)
.await
{
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": error.to_string() })),
)
.into_response();
}
tracing::info!(
target: "meerkat::auth::audit",
binding_key = ?auth_binding,
action = "logout",
"binding-scoped auth credentials logged out via REST"
);
(
StatusCode::OK,
Json(WireAuthProfileCleared {
identity: WireBindingIdentity::from(&auth_binding),
profile_id: auth_profile.id.clone(),
cleared: true,
}),
)
.into_response()
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
use axum::body::to_bytes;
use meerkat_core::handles::{AuthLeaseHandle, AuthLeasePhase, AuthLeaseTransition, LeaseKey};
use meerkat_providers::auth_store::{EphemeralTokenStore, FileTokenStore, InMemoryCoordinator};
use meerkat_runtime::RuntimeAuthLeaseHandle;
use meerkat_runtime::handles::RuntimeOAuthFlowHandle;
fn test_refresh_coordinator() -> Arc<dyn RefreshCoordinator> {
Arc::new(InMemoryCoordinator::new())
}
fn managed_auth_binding() -> AuthBindingRef {
AuthBindingRef {
realm: RealmId::parse("rest-test").unwrap(),
binding: BindingId::parse("managed").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
}
}
fn openai_auth_binding() -> AuthBindingRef {
AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
}
}
fn google_auth_binding() -> AuthBindingRef {
AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_google").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
}
}
fn generated_auth_transition_for_test(
lease_key: &LeaseKey,
expires_at: u64,
) -> AuthLeaseTransition {
let handle = RuntimeAuthLeaseHandle::new();
handle.acquire_lease(lease_key, expires_at).unwrap()
}
fn generated_auth_lease_handle_for_test(
handle: Arc<RuntimeAuthLeaseHandle>,
) -> meerkat_core::handles::GeneratedAuthLeaseHandle {
meerkat_runtime::protocol_auth_lease_lifecycle_publication::generated_auth_lease_handle(
handle,
)
.expect("test AuthLeaseHandle must be certified by generated AuthMachine authority")
}
fn mark_tokens_lifecycle_published_for_test(
tokens: &PersistedTokens,
generation: u64,
credential_published_at_millis: Option<u64>,
) -> PersistedTokens {
let _ = (generation, credential_published_at_millis);
let key = TokenKey::from_auth_binding(&openai_auth_binding());
let lease_key = LeaseKey::from_auth_binding(&openai_auth_binding());
let transition = generated_auth_transition_for_test(
&lease_key,
meerkat_core::persisted_token_expires_at_epoch_secs(tokens),
);
meerkat_core::mark_tokens_lifecycle_published_for_transition(&key, tokens, &transition)
.expect("runtime AuthMachine transition marks fixture tokens")
}
fn api_key_tokens() -> PersistedTokens {
api_key_tokens_with_secret("sk-test")
}
fn api_key_tokens_with_secret(secret: &str) -> PersistedTokens {
PersistedTokens {
auth_mode: PersistedAuthMode::ApiKey,
primary_secret: Some(secret.to_string()),
refresh_token: None,
id_token: None,
expires_at: None,
last_refresh: Some(chrono::Utc::now()),
scopes: Vec::new(),
account_id: None,
metadata: serde_json::Value::Null,
}
}
fn chatgpt_oauth_tokens_with_secret(secret: &str) -> PersistedTokens {
PersistedTokens {
auth_mode: PersistedAuthMode::ChatgptOauth,
primary_secret: Some(secret.to_string()),
refresh_token: Some(format!("{secret}-refresh")),
id_token: None,
expires_at: Some(chrono::Utc::now() + chrono::Duration::hours(1)),
last_refresh: Some(chrono::Utc::now()),
scopes: Vec::new(),
account_id: Some("acct-1".into()),
metadata: serde_json::Value::Null,
}
}
struct RejectDeviceConsumeLifecycle;
impl meerkat_providers::oauth_flow::OAuthDevicePollLifecycle for RejectDeviceConsumeLifecycle {
fn device_flow_state_is_authmachine_owned(&self) -> bool {
true
}
fn finish_device_poll(
&self,
_target: &meerkat_core::AuthCredentialIdentity,
_device_code: &str,
) -> Result<(), OAuthFlowError> {
Ok(())
}
fn consume_device_flow(
&self,
_target: &meerkat_core::AuthCredentialIdentity,
_device_code: &str,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
) -> Result<(), OAuthFlowError> {
Err(OAuthFlowError::LifecycleRejected {
operation: "consume_oauth_device_flow",
detail: "test rejection".to_string(),
})
}
fn expire_device_flow(
&self,
_target: &meerkat_core::AuthCredentialIdentity,
_device_code: &str,
) -> Result<(), OAuthFlowError> {
Ok(())
}
}
struct RejectBrowserConsumeAuthority;
impl meerkat_providers::oauth_flow::OAuthFlowAuthority for RejectBrowserConsumeAuthority {
fn terminal_flow_state_is_authmachine_owned(&self) -> bool {
true
}
fn start(
&self,
_target: meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
_redirect_uri: String,
_pkce_verifier: String,
) -> Result<String, OAuthFlowError> {
unreachable!("browser consume rollback test only consumes")
}
fn verify(
&self,
_state: &str,
_target: &meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
_redirect_uri: &str,
) -> Result<meerkat_providers::oauth_flow::OAuthFlowRecord, OAuthFlowError> {
unreachable!("browser consume rollback test only consumes")
}
fn consume(
&self,
_state: &str,
_target: &meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
_redirect_uri: &str,
) -> Result<meerkat_providers::oauth_flow::OAuthFlowRecord, OAuthFlowError> {
Err(OAuthFlowError::LifecycleRejected {
operation: "consume_oauth_browser_flow",
detail: "test rejection".to_string(),
})
}
fn admit_device_code(
&self,
_target: meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
_device_code: String,
_expires_in: std::time::Duration,
) -> Result<(), OAuthFlowError> {
unreachable!("browser consume rollback test only consumes")
}
fn verify_device_code(
&self,
_device_code: &str,
_target: &meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
) -> Result<meerkat_providers::oauth_flow::OAuthDeviceFlowRecord, OAuthFlowError> {
unreachable!("browser consume rollback test only consumes")
}
fn begin_device_code_poll(
&self,
_device_code: &str,
_target: &meerkat_core::AuthCredentialIdentity,
_provider: meerkat_providers::oauth_flow::OAuthProviderIdentity,
) -> Result<OAuthDevicePollLease, OAuthFlowError> {
unreachable!("browser consume rollback test only consumes")
}
}
struct LoadFailingTokenStore;
#[async_trait::async_trait]
impl TokenStore for LoadFailingTokenStore {
async fn load(
&self,
_key: &TokenKey,
) -> Result<Option<PersistedTokens>, meerkat_providers::auth_store::TokenStoreError>
{
Err(meerkat_providers::auth_store::TokenStoreError::Serde(
"malformed token fixture".into(),
))
}
async fn save(
&self,
_key: &TokenKey,
_tokens: &PersistedTokens,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
Ok(())
}
async fn clear(
&self,
_key: &TokenKey,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
Ok(())
}
async fn list(
&self,
) -> Result<Vec<TokenKey>, meerkat_providers::auth_store::TokenStoreError> {
Ok(Vec::new())
}
fn backend_name(&self) -> &'static str {
"load_failing"
}
}
struct SaveCountingTokenStore {
inner: EphemeralTokenStore,
save_count: std::sync::atomic::AtomicUsize,
}
impl SaveCountingTokenStore {
fn new() -> Self {
Self {
inner: EphemeralTokenStore::new(),
save_count: std::sync::atomic::AtomicUsize::new(0),
}
}
fn save_count(&self) -> usize {
self.save_count.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[async_trait::async_trait]
impl TokenStore for SaveCountingTokenStore {
async fn load(
&self,
key: &TokenKey,
) -> Result<Option<PersistedTokens>, meerkat_providers::auth_store::TokenStoreError>
{
self.inner.load(key).await
}
async fn save(
&self,
key: &TokenKey,
tokens: &PersistedTokens,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
self.save_count
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
self.inner.save(key, tokens).await
}
async fn clear(
&self,
key: &TokenKey,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
self.inner.clear(key).await
}
async fn list(
&self,
) -> Result<Vec<TokenKey>, meerkat_providers::auth_store::TokenStoreError> {
self.inner.list().await
}
fn backend_name(&self) -> &'static str {
"save_counting"
}
}
struct SaveFailingTokenStore;
#[async_trait::async_trait]
impl TokenStore for SaveFailingTokenStore {
async fn load(
&self,
_key: &TokenKey,
) -> Result<Option<PersistedTokens>, meerkat_providers::auth_store::TokenStoreError>
{
Ok(None)
}
async fn save(
&self,
_key: &TokenKey,
_tokens: &PersistedTokens,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
Err(meerkat_providers::auth_store::TokenStoreError::Io(
"durable save refused".into(),
))
}
async fn clear(
&self,
_key: &TokenKey,
) -> Result<(), meerkat_providers::auth_store::TokenStoreError> {
Ok(())
}
async fn list(
&self,
) -> Result<Vec<TokenKey>, meerkat_providers::auth_store::TokenStoreError> {
Ok(Vec::new())
}
fn backend_name(&self) -> &'static str {
"save_failing"
}
}
fn config_with_openai_managed_store_binding() -> meerkat_core::Config {
let mut config = meerkat_core::Config::default();
let mut section = meerkat_core::RealmConfigSection::default();
section.backend.insert(
"openai_backend".into(),
meerkat_core::BackendProfileConfig {
provider: "openai".into(),
backend_kind: "openai_api".into(),
base_url: None,
options: serde_json::json!({}),
server: None,
},
);
section.auth.insert(
"openai_managed".into(),
meerkat_core::AuthProfileConfig {
provider: "openai".into(),
auth_method: "api_key".into(),
source: CredentialSourceSpec::ManagedStore,
constraints: Default::default(),
metadata_defaults: Default::default(),
},
);
section.binding.insert(
"default_openai".into(),
meerkat_core::ProviderBindingConfig {
backend_profile: "openai_backend".into(),
auth_profile: "openai_managed".into(),
credential_account: None,
default_model: None,
policy: Default::default(),
provider_default: false,
},
);
section.default_binding = Some("default_openai".into());
config.realm.insert("dev".into(), section);
config
}
fn config_with_openai_oauth_binding(source: CredentialSourceSpec) -> meerkat_core::Config {
let mut config = meerkat_core::Config::default();
let mut section = meerkat_core::RealmConfigSection::default();
section.backend.insert(
"chatgpt_backend".into(),
meerkat_core::BackendProfileConfig {
provider: "openai".into(),
backend_kind: "chatgpt_backend".into(),
base_url: None,
options: serde_json::json!({}),
server: None,
},
);
section.auth.insert(
"openai_oauth".into(),
meerkat_core::AuthProfileConfig {
provider: "openai".into(),
auth_method: "managed_chatgpt_oauth".into(),
source,
constraints: Default::default(),
metadata_defaults: Default::default(),
},
);
section.binding.insert(
"default_openai".into(),
meerkat_core::ProviderBindingConfig {
backend_profile: "chatgpt_backend".into(),
auth_profile: "openai_oauth".into(),
credential_account: None,
default_model: None,
policy: Default::default(),
provider_default: false,
},
);
section.default_binding = Some("default_openai".into());
config.realm.insert("dev".into(), section);
config
}
fn config_with_openai_oauth_wrong_backend_binding() -> meerkat_core::Config {
let mut config = config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore);
config
.realm
.get_mut("dev")
.unwrap()
.backend
.get_mut("chatgpt_backend")
.unwrap()
.backend_kind = "openai_api".into();
config
}
fn config_with_openai_external_authorizer_binding() -> meerkat_core::Config {
let mut config = meerkat_core::Config::default();
let mut section = meerkat_core::RealmConfigSection::default();
section.backend.insert(
"openai_backend".into(),
meerkat_core::BackendProfileConfig {
provider: "openai".into(),
backend_kind: "openai_api".into(),
base_url: None,
options: serde_json::json!({}),
server: None,
},
);
section.auth.insert(
"openai_external".into(),
meerkat_core::AuthProfileConfig {
provider: "openai".into(),
auth_method: "external_authorizer".into(),
source: CredentialSourceSpec::ExternalResolver {
handle: "external-openai".into(),
},
constraints: Default::default(),
metadata_defaults: Default::default(),
},
);
section.binding.insert(
"default_openai".into(),
meerkat_core::ProviderBindingConfig {
backend_profile: "openai_backend".into(),
auth_profile: "openai_external".into(),
credential_account: None,
default_model: None,
policy: Default::default(),
provider_default: false,
},
);
section.default_binding = Some("default_openai".into());
config.realm.insert("dev".into(), section);
config
}
fn config_with_google_api_key_binding() -> meerkat_core::Config {
let mut config = meerkat_core::Config::default();
let mut section = meerkat_core::RealmConfigSection::default();
section.backend.insert(
"google_backend".into(),
meerkat_core::BackendProfileConfig {
provider: "gemini".into(),
backend_kind: "google_genai".into(),
base_url: None,
options: serde_json::json!({}),
server: None,
},
);
section.auth.insert(
"google_api_key".into(),
meerkat_core::AuthProfileConfig {
provider: "gemini".into(),
auth_method: "api_key".into(),
source: CredentialSourceSpec::ManagedStore,
constraints: Default::default(),
metadata_defaults: Default::default(),
},
);
section.binding.insert(
"default_google".into(),
meerkat_core::ProviderBindingConfig {
backend_profile: "google_backend".into(),
auth_profile: "google_api_key".into(),
credential_account: None,
default_model: None,
policy: Default::default(),
provider_default: false,
},
);
section.default_binding = Some("default_google".into());
config.realm.insert("dev".into(), section);
config
}
fn config_with_google_oauth_wrong_backend_binding() -> meerkat_core::Config {
let mut config = config_with_google_api_key_binding();
config
.realm
.get_mut("dev")
.unwrap()
.auth
.get_mut("google_api_key")
.unwrap()
.auth_method = "google_oauth".into();
config
}
async fn auth_status_detail(response: impl IntoResponse) -> WireAuthStatusDetail {
let response = response.into_response();
assert_eq!(response.status(), StatusCode::OK);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
serde_json::from_slice(&body).unwrap()
}
fn install_ephemeral_auth_state(state: &mut AppState) {
let auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
state
.runtime_adapter
.set_runtime_auth_lease_handle(Arc::clone(&auth_lease));
state.auth_lease = state.runtime_adapter.generated_auth_lease_handle();
state.provider_auth_persistence = ProviderAuthPersistence::new(
Arc::new(EphemeralTokenStore::new()),
Arc::new(InMemoryCoordinator::new()),
);
}
#[tokio::test]
async fn rest_test_auth_binding_is_read_only_and_does_not_advance_lease() {
use meerkat_providers::{
ProviderAuthError, ProviderClientError, ProviderRuntime, ProviderRuntimeRegistry,
ResolvedConnection, ResolverEnvironment, StaticLease, ValidatedBinding,
};
struct ResolvingOpenAiRuntime {
expires_at: chrono::DateTime<chrono::Utc>,
saw_auth_lease_handle: Arc<std::sync::atomic::AtomicBool>,
saw_refresh_coordinator: Arc<std::sync::atomic::AtomicBool>,
expected_refresh_coordinator:
Arc<dyn meerkat_providers::auth_store::RefreshCoordinator>,
}
#[async_trait::async_trait]
impl ProviderRuntime for ResolvingOpenAiRuntime {
fn provider_id(&self) -> Provider {
Provider::OpenAI
}
async fn resolve_binding(
&self,
binding: &ValidatedBinding,
env: &ResolverEnvironment,
) -> Result<ResolvedConnection, ProviderAuthError> {
self.saw_auth_lease_handle.store(
env.auth_lease_handle.is_some(),
std::sync::atomic::Ordering::SeqCst,
);
self.saw_refresh_coordinator.store(
env.provider_auth_persistence()
.map(ProviderAuthPersistence::refresh_coordinator)
.is_some_and(|actual| {
Arc::ptr_eq(&actual, &self.expected_refresh_coordinator)
}),
std::sync::atomic::Ordering::SeqCst,
);
Ok(ResolvedConnection {
provider: Provider::OpenAI,
backend: binding.backend(),
backend_profile: Arc::clone(binding.backend_profile()),
credential_identity: binding.credential_identity().clone(),
auth_lease: Arc::new(StaticLease::inline_secret(
"sk-rest-test".to_string(),
meerkat_core::AuthMetadata::default(),
Some(self.expires_at),
"test",
)),
})
}
fn build_client(
&self,
_connection: ResolvedConnection,
) -> Result<Arc<dyn meerkat_client::LlmClient>, ProviderClientError> {
Err(ProviderClientError::MissingFeature("test"))
}
}
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let expires_at = chrono::Utc::now() + chrono::Duration::hours(2);
let saw_auth_lease_handle = Arc::new(std::sync::atomic::AtomicBool::new(false));
let saw_refresh_coordinator = Arc::new(std::sync::atomic::AtomicBool::new(false));
let expected_refresh_coordinator = state.provider_auth_persistence.refresh_coordinator();
state.provider_registry = Arc::new(ProviderRuntimeRegistry::empty().with_runtime(
Arc::new(ResolvingOpenAiRuntime {
expires_at,
saw_auth_lease_handle: Arc::clone(&saw_auth_lease_handle),
saw_refresh_coordinator: Arc::clone(&saw_refresh_coordinator),
expected_refresh_coordinator,
}),
));
let response = test_auth_binding(
State(state.clone()),
Path(BindingId::parse("default_openai").unwrap()),
Json(TestBindingBody {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::OK);
assert!(
saw_auth_lease_handle.load(std::sync::atomic::Ordering::SeqCst),
"auth test resolution must receive the AuthMachine authority handle"
);
assert!(
saw_refresh_coordinator.load(std::sync::atomic::Ordering::SeqCst),
"auth test resolution must receive backend-derived refresh authority"
);
let auth_binding = AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
};
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(
snapshot.phase, None,
"test surface must not transition the AuthMachine lease into any phase"
);
assert!(
!snapshot.credential_present,
"test surface must not mark the lease credential as present"
);
assert_eq!(
snapshot.expires_at, None,
"test surface must not stamp lease freshness/expiry"
);
}
#[tokio::test]
async fn rest_login_start_is_isolated_between_runtime_authorities() {
let temp = tempfile::tempdir().unwrap();
let unrelated_temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
let unrelated_state = AppState::load_from(unrelated_temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore),
None,
)
.await
.unwrap();
unrelated_state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore),
None,
)
.await
.unwrap();
let redirect_uri = "http://127.0.0.1:0/callback";
let response = start_login(
State(state.clone()),
Json(LoginStartBody {
provider: WireOAuthProvider::OpenAi,
redirect_uri: redirect_uri.to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::OK);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let result: serde_json::Value = serde_json::from_slice(&body).unwrap();
let state_token = result["state"].as_str().expect("state");
assert!(matches!(
unrelated_state.oauth_flow_authority().consume(
state_token,
&credential_identity(&openai_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
),
Err(OAuthFlowError::LifecycleRejected {
operation: "verify_oauth_browser_flow",
..
})
));
let flow = state
.oauth_flow_authority()
.consume(
state_token,
&credential_identity(&openai_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("starting runtime owns the flow");
assert!(!flow.pkce_verifier.is_empty());
}
#[tokio::test]
async fn rest_login_start_records_flow_on_runtime_authority() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore),
None,
)
.await
.unwrap();
let redirect_uri = "http://127.0.0.1:0/callback";
let response = start_login(
State(state.clone()),
Json(LoginStartBody {
provider: WireOAuthProvider::OpenAi,
redirect_uri: redirect_uri.to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::OK);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let result: serde_json::Value = serde_json::from_slice(&body).unwrap();
let state_token = result["state"].as_str().expect("state");
let flow = state
.runtime_adapter
.oauth_flow_authority()
.consume(
state_token,
&credential_identity(&openai_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("runtime AuthMachine authority owns the REST login flow");
assert!(!flow.pkce_verifier.is_empty());
}
#[tokio::test]
async fn rest_login_start_rejects_same_provider_non_oauth_target_before_state_admission() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let response = start_login(
State(state.clone()),
Json(LoginStartBody {
provider: WireOAuthProvider::OpenAi,
redirect_uri: "http://127.0.0.1:0/callback".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("auth_method 'api_key'")
);
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&openai_auth_binding()));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_login_start_rejects_oauth_method_with_external_source_before_state_admission() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ExternalResolver {
handle: "external-chatgpt".into(),
}),
None,
)
.await
.unwrap();
let response = start_login(
State(state.clone()),
Json(LoginStartBody {
provider: WireOAuthProvider::OpenAi,
redirect_uri: "http://127.0.0.1:0/callback".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("source 'external_resolver'")
);
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&openai_auth_binding()));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_login_start_rejects_oauth_method_with_wrong_backend_before_state_admission() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_openai_oauth_wrong_backend_binding(), None)
.await
.unwrap();
let response = start_login(
State(state.clone()),
Json(LoginStartBody {
provider: WireOAuthProvider::OpenAi,
redirect_uri: "http://127.0.0.1:0/callback".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("backend_kind 'openai_api'")
);
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&openai_auth_binding()));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_device_start_rejects_same_provider_non_oauth_target_before_state_admission() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_google_api_key_binding(), None)
.await
.unwrap();
let response = start_device_login(
State(state.clone()),
Json(DeviceStartBody {
provider: WireOAuthProvider::Google,
realm_id: "dev".to_string(),
binding_id: "default_google".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("auth_method 'api_key'")
);
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&google_auth_binding()));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_device_start_rejects_oauth_method_with_wrong_backend_before_state_admission() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_google_oauth_wrong_backend_binding(), None)
.await
.unwrap();
let response = start_device_login(
State(state.clone()),
Json(DeviceStartBody {
provider: WireOAuthProvider::Google,
realm_id: "dev".to_string(),
binding_id: "default_google".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("backend_kind 'google_genai'")
);
let snapshot = state
.auth_lease
.snapshot(&LeaseKey::from_auth_binding(&google_auth_binding()));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_login_complete_rejects_same_provider_non_oauth_target_before_state_lookup() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let response = complete_login(
State(state),
Json(LoginCompleteBody {
provider: WireOAuthProvider::OpenAi,
code: "provider-code".to_string(),
state: "missing-state".to_string(),
redirect_uri: "http://127.0.0.1:0/callback".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("auth_method 'api_key'")
);
}
#[tokio::test]
async fn rest_login_complete_rejects_oauth_method_with_wrong_backend_before_state_lookup() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_openai_oauth_wrong_backend_binding(), None)
.await
.unwrap();
let response = complete_login(
State(state),
Json(LoginCompleteBody {
provider: WireOAuthProvider::OpenAi,
code: "provider-code".to_string(),
state: "missing-state".to_string(),
redirect_uri: "http://127.0.0.1:0/callback".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_openai".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("backend_kind 'openai_api'")
);
}
#[tokio::test]
async fn rest_device_complete_rejects_same_provider_non_oauth_target_before_state_lookup() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_google_api_key_binding(), None)
.await
.unwrap();
let response = complete_device_login(
State(state),
Json(DeviceCompleteBody {
provider: WireOAuthProvider::Google,
device_code: "missing-device-code".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_google".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("auth_method 'api_key'")
);
}
#[tokio::test]
async fn rest_device_complete_rejects_oauth_method_with_wrong_backend_before_state_lookup() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.config_runtime
.set(config_with_google_oauth_wrong_backend_binding(), None)
.await
.unwrap();
let response = complete_device_login(
State(state),
Json(DeviceCompleteBody {
provider: WireOAuthProvider::Google,
device_code: "missing-device-code".to_string(),
realm_id: "dev".to_string(),
binding_id: "default_google".to_string(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let error: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("backend_kind 'google_genai'")
);
}
#[tokio::test]
async fn rest_device_completion_poll_drop_releases_runtime_authority() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.oauth_flow_authority()
.admit_device_code(
credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
"device-code".to_string(),
std::time::Duration::from_secs(600),
)
.expect("device code admitted");
let poll = state
.oauth_flow_authority()
.begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("device completion poll begins");
assert!(matches!(
state.oauth_flow_authority().begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
),
Err(OAuthFlowError::LifecycleRejected {
operation: "begin_oauth_device_poll",
..
})
));
drop(poll);
state
.oauth_flow_authority()
.begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("dropped REST completion poll releases runtime authority");
}
#[tokio::test]
async fn rest_device_completion_poll_abort_releases_runtime_authority() {
let temp = tempfile::tempdir().unwrap();
let state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
state
.oauth_flow_authority()
.admit_device_code(
credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
"device-code".to_string(),
std::time::Duration::from_secs(600),
)
.expect("device code admitted");
let authority = state.oauth_flow_authority();
let (begun_tx, begun_rx) = tokio::sync::oneshot::channel();
let task = tokio::spawn(async move {
let _poll = authority
.begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("device completion poll begins");
begun_tx.send(()).expect("signal poll start");
std::future::pending::<()>().await;
});
begun_rx.await.expect("poll lease was acquired");
assert!(matches!(
state.oauth_flow_authority().begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
),
Err(OAuthFlowError::LifecycleRejected {
operation: "begin_oauth_device_poll",
..
})
));
task.abort();
let _ = task.await;
state
.oauth_flow_authority()
.begin_device_code_poll(
"device-code",
&credential_identity(&managed_auth_binding()),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("aborted REST completion poll releases runtime authority");
}
#[tokio::test]
async fn rest_token_write_helpers_publish_auth_machine_lifecycle() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let tokens = api_key_tokens();
save_tokens_and_publish_lifecycle(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
tokens,
)
.await
.unwrap();
assert!(store.load(&key).await.unwrap().is_some());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, Some(AuthLeasePhase::Valid));
clear_tokens_and_publish_lifecycle(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
)
.await
.unwrap();
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
assert_eq!(snapshot.expires_at, None);
}
#[tokio::test]
async fn rest_token_commit_is_single_marked_save() {
let store = Arc::new(SaveCountingTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
save_tokens_and_publish_lifecycle(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
)
.await
.unwrap();
assert_eq!(
store.save_count(),
1,
"acquire-first commit must persist the credential in a single marked save"
);
let stored = store.load(&key).await.unwrap().unwrap();
assert!(
meerkat_core::tokens_lifecycle_published(&stored),
"the only durable write must already carry the proof-of-acquisition marker"
);
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, Some(AuthLeasePhase::Valid));
}
#[tokio::test]
async fn rest_save_failure_after_acquire_releases_lease() {
let store = Arc::new(SaveFailingTokenStore);
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let err = save_tokens_and_publish_lifecycle(
ProviderAuthPersistence::new(store, test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("TokenStore save failed"), "{}", err.1);
assert!(err.1.contains("acquired lease rolled back"), "{}", err.1);
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(
snapshot.phase, None,
"save failure after acquire must release the freshly acquired lease"
);
assert!(!snapshot.credential_present);
}
#[tokio::test]
async fn rest_orphan_unmarked_token_is_rejected_on_rehydration() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let orphan = chatgpt_oauth_tokens_with_secret("orphan-access");
assert!(!meerkat_core::tokens_lifecycle_published(&orphan));
store.save(&key, &orphan).await.unwrap();
let restored = meerkat_core::rehydrate_marked_tokens_for_status(
store.as_ref(),
&auth_lease,
&auth_binding,
PersistedAuthMode::ChatgptOauth,
chrono::Utc::now(),
)
.await
.unwrap();
assert!(
restored.is_none(),
"unmarked durable tokens must be rejected on rehydration"
);
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
assert!(!snapshot.credential_present);
}
#[tokio::test]
async fn rest_ready_device_consume_failure_does_not_commit_token() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let registry = meerkat_providers::oauth_flow::OAuthFlowRegistry::new(
std::time::Duration::from_secs(600),
);
registry
.admit_device_code(
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
"device-code".to_string(),
std::time::Duration::from_secs(600),
)
.expect("device code admitted");
let poll_lease = registry
.begin_device_code_poll(
"device-code",
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("device poll lease begins")
.with_lifecycle(Arc::new(RejectDeviceConsumeLifecycle));
let err = save_tokens_and_consume_device_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
poll_lease,
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_device_flow"));
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_browser_consume_failure_does_not_commit_token() {
let temp = tempfile::tempdir().unwrap();
let store = Arc::new(FileTokenStore::new(temp.path().join("tokens")));
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let authority: Arc<dyn meerkat_providers::oauth_flow::OAuthFlowAuthority> =
Arc::new(RejectBrowserConsumeAuthority);
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
OwnedBrowserFlowConsume {
authority,
state: "browser-state".to_string(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: "http://127.0.0.1/callback".to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_browser_flow"));
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_browser_consume_failure_does_not_save_before_durable_claim() {
let store = Arc::new(SaveCountingTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
OwnedBrowserFlowConsume {
authority: Arc::new(RejectBrowserConsumeAuthority),
state: "browser-state".to_string(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: "http://127.0.0.1/callback".to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_browser_flow"));
assert_eq!(
store.save_count(),
0,
"token material must not be saved before winning the durable OAuth consume claim"
);
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_raw_registry_browser_success_cannot_commit_tokens() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let redirect_uri = "http://127.0.0.1/callback";
let registry = Arc::new(meerkat_providers::oauth_flow::OAuthFlowRegistry::new(
std::time::Duration::from_secs(600),
));
let state = registry
.start(
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"registry-only-verifier".to_string(),
)
.expect("raw registry admits browser flow");
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
OwnedBrowserFlowConsume {
authority: registry.clone(),
state: state.clone(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: redirect_uri.to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("AuthMachine-owned OAuth flow authority"));
assert!(store.load(&key).await.unwrap().is_none());
assert_eq!(
auth_lease
.snapshot(&LeaseKey::from_auth_binding(&auth_binding))
.phase,
None
);
registry
.verify(
&state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("raw registry success path is inert and remains unconsumed");
}
#[tokio::test]
async fn rest_raw_registry_device_success_cannot_commit_tokens() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let registry = meerkat_providers::oauth_flow::OAuthFlowRegistry::new(
std::time::Duration::from_secs(600),
);
registry
.admit_device_code(
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
"registry-only-device-code".to_string(),
std::time::Duration::from_secs(600),
)
.expect("raw registry admits device flow");
let poll_lease = registry
.begin_device_code_poll(
"registry-only-device-code",
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("raw registry begins device poll");
let err = save_tokens_and_consume_device_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
poll_lease,
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("AuthMachine-owned OAuth device poll lease"));
assert!(store.load(&key).await.unwrap().is_none());
assert_eq!(
auth_lease
.snapshot(&LeaseKey::from_auth_binding(&auth_binding))
.phase,
None
);
registry
.verify_device_code(
"registry-only-device-code",
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("raw registry device success path is inert and remains unconsumed");
}
#[tokio::test]
async fn rest_browser_consume_failure_does_not_reauthorize_stale_previous_tokens() {
let store = Arc::new(EphemeralTokenStore::new());
let auth_lease =
generated_auth_lease_handle_for_test(Arc::new(RuntimeAuthLeaseHandle::new()));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let stale = api_key_tokens_with_secret("sk-stale");
store.save(&key, &stale).await.unwrap();
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens_with_secret("sk-new"),
OwnedBrowserFlowConsume {
authority: Arc::new(RejectBrowserConsumeAuthority),
state: "browser-state".to_string(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: "http://127.0.0.1/callback".to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_browser_flow"));
let stored = store.load(&key).await.unwrap().unwrap();
assert_eq!(stored.primary_secret.as_deref(), Some("sk-stale"));
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
let projection = meerkat_core::project_published_auth_status(
chrono::Utc::now(),
Some(&stored),
&snapshot,
);
assert_eq!(
projection.phase,
meerkat_core::AuthStatusPhase::MissingCredential
);
assert!(projection.tokens.is_none());
}
#[tokio::test]
async fn rest_stale_previous_rollback_preserves_newer_oauth_flow() {
let store = Arc::new(EphemeralTokenStore::new());
let raw_auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
let auth_lease = generated_auth_lease_handle_for_test(Arc::clone(&raw_auth_lease));
let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
std::time::Duration::from_secs(600),
Arc::clone(&raw_auth_lease),
);
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let redirect_uri = "http://127.0.0.1/callback";
store
.save(&key, &api_key_tokens_with_secret("sk-stale"))
.await
.unwrap();
let mut failed_tokens = api_key_tokens_with_secret("sk-new");
failed_tokens.expires_at =
Some(chrono::DateTime::from_timestamp(1_800_000_000, 0).unwrap());
let commit = save_tokens_and_publish_lifecycle_commit_unlocked(
store.as_ref(),
&auth_lease,
&auth_binding,
&failed_tokens,
)
.await
.unwrap();
let state = meerkat_providers::oauth_flow::OAuthFlowAuthority::start(
&authority,
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"new-verifier".to_string(),
)
.expect("newer OAuth flow admitted after rollback snapshot");
rollback_token_commit(store.as_ref(), &auth_lease, &commit)
.await
.expect("rollback clears credential lifecycle without clobbering OAuth flow");
let stored = store.load(&key).await.unwrap().unwrap();
assert_eq!(stored.primary_secret.as_deref(), Some("sk-stale"));
let record = meerkat_providers::oauth_flow::OAuthFlowAuthority::verify(
&authority,
&state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("newer OAuth flow remains authoritative");
assert_eq!(record.pkce_verifier, "new-verifier");
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, Some(AuthLeasePhase::ReauthRequired));
assert_eq!(snapshot.expires_at, None);
let projection = meerkat_core::project_published_auth_status(
chrono::Utc::now(),
Some(&stored),
&snapshot,
);
assert_eq!(
projection.phase,
meerkat_core::AuthStatusPhase::MissingCredential
);
assert!(projection.tokens.is_none());
}
#[tokio::test]
async fn rest_oauth_only_previous_rollback_does_not_reauthorize_stale_tokens() {
let store = EphemeralTokenStore::new();
let raw_auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
let auth_lease = generated_auth_lease_handle_for_test(Arc::clone(&raw_auth_lease));
let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
std::time::Duration::from_secs(600),
Arc::clone(&raw_auth_lease),
);
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let redirect_uri = "http://127.0.0.1/callback";
store
.save(&key, &api_key_tokens_with_secret("sk-stale"))
.await
.unwrap();
let state = meerkat_providers::oauth_flow::OAuthFlowAuthority::start(
&authority,
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"existing-flow-verifier".to_string(),
)
.expect("existing OAuth flow admitted before rollback snapshot");
let previous_snapshot = auth_lease.snapshot(&lease_key);
assert_eq!(
previous_snapshot.phase,
Some(AuthLeasePhase::ReauthRequired)
);
assert!(!previous_snapshot.credential_present);
let mut failed_tokens = api_key_tokens_with_secret("sk-new");
failed_tokens.expires_at =
Some(chrono::DateTime::from_timestamp(1_800_000_000, 0).unwrap());
let commit = save_tokens_and_publish_lifecycle_commit_unlocked(
&store,
&auth_lease,
&auth_binding,
&failed_tokens,
)
.await
.unwrap();
rollback_token_commit(&store, &auth_lease, &commit)
.await
.expect("rollback restores token bytes without restoring credential authority");
let stored = store.load(&key).await.unwrap().unwrap();
assert_eq!(stored.primary_secret.as_deref(), Some("sk-stale"));
let record = meerkat_providers::oauth_flow::OAuthFlowAuthority::verify(
&authority,
&state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("existing OAuth flow remains authoritative");
assert_eq!(record.pkce_verifier, "existing-flow-verifier");
let snapshot = auth_lease.snapshot(&lease_key);
assert_eq!(snapshot.phase, Some(AuthLeasePhase::ReauthRequired));
assert!(!snapshot.credential_present);
let projection = meerkat_core::project_published_auth_status(
chrono::Utc::now(),
Some(&stored),
&snapshot,
);
assert_eq!(
projection.phase,
meerkat_core::AuthStatusPhase::MissingCredential
);
assert!(projection.tokens.is_none());
}
#[tokio::test]
async fn rest_real_device_consume_failure_releases_credential_lifecycle() {
let store = Arc::new(EphemeralTokenStore::new());
let raw_auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
let auth_lease = generated_auth_lease_handle_for_test(Arc::clone(&raw_auth_lease));
let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
std::time::Duration::from_secs(600),
Arc::clone(&raw_auth_lease),
);
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
meerkat_providers::oauth_flow::OAuthFlowAuthority::admit_device_code(
&authority,
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
"device-code".to_string(),
std::time::Duration::from_secs(600),
)
.expect("device flow admitted by runtime authority");
let poll_lease = meerkat_providers::oauth_flow::OAuthFlowAuthority::begin_device_code_poll(
&authority,
"device-code",
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::GoogleCodeAssist,
)
.expect("device poll lease begins through runtime authority");
meerkat_providers::oauth_flow::OAuthDevicePollLifecycle::expire_device_flow(
raw_auth_lease.as_ref(),
&credential_identity(&auth_binding),
"device-code",
)
.expect("test removes the AuthMachine flow membership");
let err = save_tokens_and_consume_device_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
poll_lease,
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_device_flow"));
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_real_browser_missing_consume_releases_credential_lifecycle() {
let store = Arc::new(EphemeralTokenStore::new());
let raw_auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
let auth_lease = generated_auth_lease_handle_for_test(Arc::clone(&raw_auth_lease));
let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
std::time::Duration::from_secs(600),
Arc::clone(&raw_auth_lease),
));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let redirect_uri = "http://127.0.0.1/callback";
let state = meerkat_providers::oauth_flow::OAuthFlowAuthority::start(
authority.as_ref(),
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"verifier".to_string(),
)
.expect("browser flow admitted by runtime authority");
meerkat_providers::oauth_flow::OAuthFlowAuthority::verify(
authority.as_ref(),
&state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("browser flow verifies through runtime authority");
meerkat_providers::oauth_flow::OAuthFlowAuthority::consume(
authority.as_ref(),
&state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("test pre-consumes the browser flow");
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
OwnedBrowserFlowConsume {
authority,
state: state.clone(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: redirect_uri.to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("oauth state terminal consume failed"));
assert!(err.1.contains("verify_oauth_browser_flow"));
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, None);
}
#[tokio::test]
async fn rest_terminal_consume_failure_preserves_other_browser_flow() {
let store = Arc::new(EphemeralTokenStore::new());
let raw_auth_lease = Arc::new(RuntimeAuthLeaseHandle::new());
let auth_lease = generated_auth_lease_handle_for_test(Arc::clone(&raw_auth_lease));
let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
std::time::Duration::from_secs(600),
Arc::clone(&raw_auth_lease),
));
let auth_binding = managed_auth_binding();
let key = TokenKey::from_auth_binding(&auth_binding);
let redirect_uri = "http://127.0.0.1/callback";
let old_state = meerkat_providers::oauth_flow::OAuthFlowAuthority::start(
authority.as_ref(),
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"old-verifier".to_string(),
)
.expect("old browser flow admitted");
let new_state = meerkat_providers::oauth_flow::OAuthFlowAuthority::start(
authority.as_ref(),
credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri.to_string(),
"new-verifier".to_string(),
)
.expect("new browser flow admitted");
let rejecting: Arc<dyn meerkat_providers::oauth_flow::OAuthFlowAuthority> =
Arc::new(RejectBrowserConsumeAuthority);
let err = save_tokens_and_consume_browser_flow(
ProviderAuthPersistence::new(store.clone(), test_refresh_coordinator()),
auth_lease.clone(),
auth_binding.clone(),
api_key_tokens(),
OwnedBrowserFlowConsume {
authority: rejecting,
state: old_state.clone(),
provider: meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri: redirect_uri.to_string(),
},
)
.await
.unwrap_err();
assert_eq!(err.0, StatusCode::INTERNAL_SERVER_ERROR);
assert!(err.1.contains("consume_oauth_browser_flow"));
assert!(store.load(&key).await.unwrap().is_none());
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, Some(AuthLeasePhase::ReauthRequired));
let record = meerkat_providers::oauth_flow::OAuthFlowAuthority::verify(
authority.as_ref(),
&new_state,
&credential_identity(&auth_binding),
meerkat_providers::oauth_flow::OAuthProviderIdentity::OpenAiChatGpt,
redirect_uri,
)
.expect("rollback preserves other admitted browser flow");
assert_eq!(record.pkce_verifier, "new-verifier");
let snapshot = auth_lease.snapshot(&LeaseKey::from_auth_binding(&auth_binding));
assert_eq!(snapshot.phase, Some(AuthLeasePhase::ReauthRequired));
}
#[tokio::test]
async fn rest_auth_status_hides_stale_token_without_auth_machine_lifecycle() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let auth_binding = AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
};
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&PersistedTokens {
auth_mode: PersistedAuthMode::ApiKey,
primary_secret: Some("sk-stale".into()),
refresh_token: Some("refresh-stale".into()),
id_token: None,
expires_at: Some(chrono::DateTime::from_timestamp(1_800_000_000, 0).unwrap()),
last_refresh: Some(chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap()),
scopes: Vec::new(),
account_id: Some("acct-stale".into()),
metadata: serde_json::Value::Null,
},
)
.await
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(
detail.state,
meerkat_core::AuthStatusPhase::MissingCredential
);
assert_eq!(detail.expires_at, None);
assert_eq!(detail.last_refresh_at, None);
assert_eq!(detail.account_id, None);
assert!(!detail.has_refresh_token);
}
#[tokio::test]
async fn rest_auth_status_rehydrates_marked_oauth_token_after_restart() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore),
None,
)
.await
.unwrap();
let auth_binding = openai_auth_binding();
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let tokens = chatgpt_oauth_tokens_with_secret("fresh-chatgpt-access");
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&mark_tokens_lifecycle_published_for_test(&tokens, 1, None),
)
.await
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state.clone()),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
assert!(detail.expires_at.is_some());
assert_eq!(detail.account_id.as_deref(), Some("acct-1"));
assert!(detail.has_refresh_token);
let snapshot = state.auth_lease.snapshot(&lease_key);
assert_eq!(snapshot.phase, Some(AuthLeasePhase::Valid));
assert!(snapshot.credential_present);
}
#[tokio::test]
async fn rest_auth_status_reports_lease_phase_when_token_is_missing() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let auth_binding = AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
};
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
assert!(detail.expires_at.is_some());
assert!(!detail.has_refresh_token);
}
#[tokio::test]
async fn rest_auth_status_hides_wrong_mode_token_without_hiding_lifecycle() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ManagedStore),
None,
)
.await
.unwrap();
let auth_binding = openai_auth_binding();
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&PersistedTokens::api_key("sk-stale-api-key"),
)
.await
.unwrap();
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
assert!(detail.expires_at.is_some());
assert_eq!(detail.last_refresh_at, None);
assert_eq!(detail.account_id, None);
assert!(!detail.has_refresh_token);
}
#[tokio::test]
async fn rest_auth_status_hides_wrong_source_oauth_token_without_hiding_lifecycle() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(
config_with_openai_oauth_binding(CredentialSourceSpec::ExternalResolver {
handle: "external-chatgpt".into(),
}),
None,
)
.await
.unwrap();
let auth_binding = openai_auth_binding();
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&PersistedTokens {
auth_mode: PersistedAuthMode::ChatgptOauth,
primary_secret: Some("fresh-chatgpt-access".into()),
refresh_token: Some("rt".into()),
id_token: None,
expires_at: Some(chrono::Utc::now() + chrono::Duration::hours(1)),
last_refresh: Some(chrono::Utc::now()),
scopes: Vec::new(),
account_id: Some("acct-stale".into()),
metadata: serde_json::Value::Null,
},
)
.await
.unwrap();
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
assert!(detail.expires_at.is_some());
assert_eq!(detail.last_refresh_at, None);
assert_eq!(detail.account_id, None);
assert!(!detail.has_refresh_token);
}
#[tokio::test]
async fn rest_auth_status_ignores_stale_token_for_non_persisted_source_without_hiding_lifecycle()
{
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(config_with_openai_external_authorizer_binding(), None)
.await
.unwrap();
let auth_binding = openai_auth_binding();
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&PersistedTokens {
auth_mode: PersistedAuthMode::ApiKey,
primary_secret: Some("sk-stale".into()),
refresh_token: Some("refresh-stale".into()),
id_token: None,
expires_at: Some(chrono::Utc::now() + chrono::Duration::hours(1)),
last_refresh: Some(chrono::Utc::now()),
scopes: Vec::new(),
account_id: Some("acct-stale".into()),
metadata: serde_json::Value::Null,
},
)
.await
.unwrap();
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let detail = auth_status_detail(
get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
assert!(detail.expires_at.is_some());
assert_eq!(detail.last_refresh_at, None);
assert_eq!(detail.account_id, None);
assert!(!detail.has_refresh_token);
}
#[tokio::test]
async fn rest_auth_status_fails_closed_when_token_load_fails() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state.provider_auth_persistence = ProviderAuthPersistence::new(
Arc::new(LoadFailingTokenStore),
Arc::new(InMemoryCoordinator::new()),
);
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let auth_binding = AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
};
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let response = get_auth_status(
State(state),
Path(BindingId::parse("default_openai").unwrap()),
Query(RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
}),
)
.await
.into_response();
assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR);
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let body = String::from_utf8(body.to_vec()).unwrap();
assert!(
body.contains("TokenStore load failed"),
"error body must carry the token-store fault, got: {body}"
);
}
#[tokio::test]
async fn rest_auth_status_observes_session_runtime_binding_lifecycle() {
let temp = tempfile::tempdir().unwrap();
let mut state = AppState::load_from(temp.path().to_path_buf())
.await
.unwrap();
install_ephemeral_auth_state(&mut state);
state
.config_runtime
.set(config_with_openai_managed_store_binding(), None)
.await
.unwrap();
let auth_binding = AuthBindingRef {
realm: RealmId::parse("dev").unwrap(),
binding: BindingId::parse("default_openai").unwrap(),
profile: None,
origin: meerkat_core::connection::BindingOrigin::Configured,
};
let lease_key = LeaseKey::from_auth_binding(&auth_binding);
let now = chrono::Utc::now().timestamp() as u64;
let query = || RealmQuery {
realm_id: RealmId::parse("dev").unwrap(),
profile_id: None,
};
let binding_id = || BindingId::parse("default_openai").unwrap();
state
.token_store()
.save(
&TokenKey::from_auth_binding(&auth_binding),
&PersistedTokens::api_key("sk-live"),
)
.await
.unwrap();
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings
.auth_lease()
.acquire_lease(&lease_key, now + 3600)
.unwrap();
let detail = auth_status_detail(
get_auth_status(State(state.clone()), Path(binding_id()), Query(query())).await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Valid);
let bindings = state
.runtime_adapter
.prepare_bindings(meerkat_core::SessionId::new())
.await
.unwrap();
bindings.auth_lease().begin_refresh(&lease_key).unwrap();
let detail = auth_status_detail(
get_auth_status(State(state.clone()), Path(binding_id()), Query(query())).await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::Expiring);
bindings
.auth_lease()
.complete_refresh(&lease_key, now + 7200, now)
.unwrap();
bindings
.auth_lease()
.mark_reauth_required(&lease_key)
.unwrap();
let detail = auth_status_detail(
get_auth_status(State(state), Path(binding_id()), Query(query())).await,
)
.await;
assert_eq!(detail.state, meerkat_core::AuthStatusPhase::ReauthRequired);
}
#[test]
fn login_complete_body_requires_explicit_identity() {
let err = serde_json::from_value::<LoginCompleteBody>(serde_json::json!({
"provider": "anthropic",
"code": "code",
"state": "state",
"redirect_uri": "http://127.0.0.1:0/callback"
}))
.unwrap_err();
assert!(err.to_string().contains("realm_id"));
}
#[test]
fn device_complete_body_requires_explicit_identity() {
let err = serde_json::from_value::<DeviceCompleteBody>(serde_json::json!({
"provider": "anthropic",
"device_code": "device-code"
}))
.unwrap_err();
assert!(err.to_string().contains("realm_id"));
}
}
#[cfg(test)]
fn credential_identity(binding: &AuthBindingRef) -> meerkat_core::AuthCredentialIdentity {
meerkat_core::AuthCredentialIdentity::from_auth_binding(binding)
}