mod common;
use common::*;
use backbone_integrations::application::service::integrations_oauth::{
AuthorizeRequest, CompleteRequest, IntegrationsOauthConfig, IntegrationsOauthService,
};
use backbone_integrations::infrastructure::http::endpoint_guard::{
EndpointOverrides, IdentityClaims, OAuthClientConfig, OAuthClientConfigs, ProviderRegistry,
ProviderEndpointOverride, ReqwestOAuthTransport, TokenResponse, TransportFailure,
TransportFailureKind, ValidatedEndpoints,
};
use backbone_integrations::infrastructure::jobs::refresh_oauth_credentials::{
refresh_oauth_credentials, RefreshSchedule,
};
use backbone_integrations::presentation::http::{create_oauth_routes, OAuthPrincipal};
use axum::body::Body;
use axum::http::{Method, Request, StatusCode};
use base64::Engine;
use chrono::{Duration, Utc};
use sqlx::PgPool;
use std::collections::BTreeMap;
use std::sync::Arc;
use tower::ServiceExt;
use uuid::Uuid;
fn registry() -> ProviderRegistry {
ProviderRegistry::with_builtin()
}
fn clients() -> OAuthClientConfigs {
let mut map = BTreeMap::new();
for provider in registry().providers() {
map.insert(
provider.to_string(),
OAuthClientConfig { client_id: format!("client-{provider}"), client_secret: None },
);
}
OAuthClientConfigs(map)
}
fn overrides_for(provider: &str, o: ProviderEndpointOverride) -> EndpointOverrides {
let mut map = BTreeMap::new();
map.insert(provider.to_string(), o);
EndpointOverrides(map)
}
const STATE_SECRET: &str = "probe-state-secret-0123456789abcdef0123456789abcdef";
fn oauth_config() -> IntegrationsOauthConfig {
IntegrationsOauthConfig {
public_base: Some("https://api.example.test".into()),
clients: clients(),
..Default::default()
}
}
fn oauth_service(
pool: &PgPool,
transport: FakeTransport,
store: FakeStore,
) -> IntegrationsOauthService {
std::env::set_var("INTEGRATIONS_OAUTH_STATE_SECRET", STATE_SECRET);
IntegrationsOauthService::build(
pool.clone(),
oauth_config(),
Arc::new(transport),
Arc::new(store),
)
.expect("oauth service build")
}
fn decode_state(state: &str) -> serde_json::Value {
let (payload, _) = state.split_once('.').expect("state payload.mac shape");
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(payload)
.expect("state payload is base64url");
serde_json::from_slice(&bytes).expect("state payload is JSON")
}
fn mint_state(account_id: Uuid, provider: &str, nonce: &str, exp: i64) -> String {
use hmac::{Hmac, Mac};
use sha2::Sha256;
let payload = serde_json::json!({ "account_id": account_id, "provider": provider, "nonce": nonce, "exp": exp });
let mut mac = Hmac::<Sha256>::new_from_slice(STATE_SECRET.as_bytes()).unwrap();
mac.update(payload.to_string().as_bytes());
format!(
"{}.{}",
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(payload.to_string().as_bytes()),
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
)
}
async fn account_rows(pool: &PgPool, ids: &[Uuid]) -> Vec<(Uuid, String, String)> {
sqlx::query_as(
"SELECT id, status::text, account_ref FROM integrations.integration_accounts WHERE id = ANY($1) ORDER BY account_ref",
)
.bind(ids)
.fetch_all(pool)
.await
.expect("snapshot account rows")
}
#[tokio::test]
async fn ioa8_malicious_overrides_are_refused_at_resolution() {
let reg = registry();
let cases: &[(&str, &str)] = &[
("http scheme", "http://accounts.google.com/o/oauth2/v2/auth"),
("off-allowlist host", "https://evil.example.com/oauth2/auth"),
("allowlist-suffix trick", "https://accounts.google.com.evil.example.com/auth"),
("userinfo component", "https://user@accounts.google.com/o/oauth2/v2/auth"),
("IP literal", "https://127.0.0.1/oauth2/auth"),
("non-443 port", "https://accounts.google.com:8443/o/oauth2/v2/auth"),
("empty path", "https://accounts.google.com/"),
("query string", "https://accounts.google.com/auth?next=/steal"),
];
for (label, url) in cases {
let ov = overrides_for(
"gmail",
ProviderEndpointOverride { authorize: Some((*url).into()), token: None, userinfo: None },
);
let resolved = ValidatedEndpoints::resolve(®, "gmail", &ov);
assert!(resolved.is_err(), "{label}: malicious override accepted ({url})");
}
}
#[tokio::test]
async fn ioa8_registry_defaults_resolve_cleanly() {
let reg = registry();
let providers = reg.providers();
assert_eq!(providers.len(), 4, "exactly the four providers of the one generation");
for provider in providers {
ValidatedEndpoints::resolve(®, provider, &EndpointOverrides::default())
.unwrap_or_else(|e| panic!("{}: registry default rejected: {e:?}", provider));
}
}
#[tokio::test]
async fn ioa8_legitimate_override_accepted() {
let reg = registry();
let ov = overrides_for(
"gmail",
ProviderEndpointOverride {
authorize: Some("https://accounts.google.com/custom/auth".into()),
token: None,
userinfo: None,
},
);
let endpoints = ValidatedEndpoints::resolve(®, "gmail", &ov).expect("legitimate override refused");
assert_eq!(endpoints.authorize.as_str(), "https://accounts.google.com/custom/auth");
}
#[tokio::test]
async fn ioa8_real_transport_has_no_redirect_policy() {
let transport = ReqwestOAuthTransport::new().expect("transport construction");
assert!(
transport.redirect_policy_is_none(),
"the OAuth transport must never follow redirects"
);
}
#[tokio::test]
async fn ioa8_cross_provider_host_refused() {
let reg = registry();
let ov = overrides_for(
"gmail",
ProviderEndpointOverride {
authorize: Some("https://login.microsoftonline.com/common/oauth2/v2.0/authorize".into()),
token: None,
userinfo: None,
},
);
assert!(
ValidatedEndpoints::resolve(®, "gmail", &ov).is_err(),
"a Microsoft host must not serve as gmail's authorize endpoint"
);
}
static SWEEP_LOCK: std::sync::OnceLock<tokio::sync::Mutex<()>> = std::sync::OnceLock::new();
async fn sweep_gate() -> tokio::sync::MutexGuard<'static, ()> {
SWEEP_LOCK.get_or_init(|| tokio::sync::Mutex::new(())).lock().await
}
async fn seed_active_account(
pool: &PgPool,
company: Uuid,
account_ref: &str,
due_in: i64,
store: &FakeStore,
) -> Uuid {
let id = Uuid::new_v4();
sqlx::query(
r#"INSERT INTO integrations.integration_accounts
(id, provider, account_ref, status, scopes, expires_at)
VALUES ($1, 'gmail', $2, 'active', '', now() + make_interval(secs => $3))"#,
)
.bind(id)
.bind(account_ref)
.bind(due_in)
.execute(pool)
.await
.expect("seed account");
store
.issue(
company,
"gmail",
account_ref,
PURPOSE_OAUTH_TOKEN,
TokenBundle::new(
"SEED-ACCESS".into(),
Some("SEED-REFRESH".into()),
Utc::now() + Duration::hours(24),
None,
),
Utc::now() + Duration::hours(24),
)
.await
.expect("seed credential");
id
}
async fn account_state(pool: &PgPool, id: Uuid) -> (String, Option<chrono::DateTime<Utc>>, Option<chrono::DateTime<Utc>>) {
sqlx::query_as(
"SELECT status::text, expires_at, last_refreshed_at FROM integrations.integration_accounts WHERE id=$1",
)
.bind(id)
.fetch_one(pool)
.await
.expect("account row")
}
#[tokio::test]
async fn ioa9_due_account_rotates_with_new_honest_expiry() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::happy("user@example.com", "client-gmail", "nonce-x");
let account =
seed_active_account(&pool, company, "due-user@example.com", 300, &store).await;
let before_expires_at = Utc::now() + Duration::seconds(300);
let report = refresh_oauth_credentials(
&pool, company, ®istry(), &EndpointOverrides::default(), &clients(), &store, &transport,
&RefreshSchedule::default(),
)
.await
.expect("refresh run");
assert_eq!(report.refreshed, 1, "one due account refreshed: {report:?}");
assert_eq!(store.rotate_count(), 1, "the rotation rode the credential port");
match store.calls().iter().find(|c| matches!(c, StoreCall::Rotated { .. })) {
Some(StoreCall::Rotated { expires_at, .. }) => {
let got = *expires_at;
let expect = Utc::now() + Duration::seconds(86_400);
assert!(
got + Duration::seconds(5) > expect && expect + Duration::seconds(5) > got,
"rotated expiry must be now + expires_in (got {got:?}, want ~{expect:?})"
);
}
other => panic!("no rotate call recorded: {other:?}"),
}
let (status, expires_at, last_refreshed_at) = account_state(&pool, account).await;
assert_eq!(status, "active");
let mirror = expires_at.expect("active account keeps its expiry mirror");
assert!(
mirror > before_expires_at + Duration::hours(23),
"the mirror advanced to the new honest expiry (got {mirror:?})"
);
assert!(last_refreshed_at.is_some(), "last_refreshed_at recorded");
assert_eq!(store.revoke_count(), 0, "rotation is not a revoke");
}
#[tokio::test]
async fn ioa9_not_due_account_untouched() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::happy("user@example.com", "client-gmail", "nonce-x");
let account = seed_active_account(&pool, company, "far-user@example.com", 86_400, &store).await;
let report = refresh_oauth_credentials(
&pool, company, ®istry(), &EndpointOverrides::default(), &clients(), &store, &transport,
&RefreshSchedule::default(),
)
.await
.expect("refresh run");
assert_eq!(report.refreshed, 0, "nothing due: {report:?}");
assert_eq!(store.rotate_count(), 0);
let (_, expires_at, last_refreshed_at) = account_state(&pool, account).await;
assert!(expires_at.is_some() && last_refreshed_at.is_none(), "row untouched");
}
#[tokio::test]
async fn ioa9_concurrent_runs_claim_disjoint_sets() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::happy("user@example.com", "client-gmail", "nonce-x");
let n = 6;
let mut refs = Vec::new();
for i in 0..n {
let r = format!("race-{i}@example.com");
seed_active_account(&pool, company, &r, 300, &store).await;
refs.push(r);
}
let schedule = RefreshSchedule::default();
let reg = registry();
let ov = EndpointOverrides::default();
let cl = clients();
let (a, b) = tokio::join!(
refresh_oauth_credentials(&pool, company, ®, &ov, &cl, &store, &transport, &schedule),
refresh_oauth_credentials(&pool, company, ®, &ov, &cl, &store, &transport, &schedule),
);
let (a, b) = (a.expect("run a"), b.expect("run b"));
assert_eq!(a.refreshed + b.refreshed, n, "every due account refreshed exactly once: {a:?} + {b:?}");
let rotations: Vec<(String, String)> = store
.calls()
.iter()
.filter_map(|c| match c {
StoreCall::Rotated { account_ref, provider, .. } => Some((provider.clone(), account_ref.clone())),
_ => None,
})
.collect();
assert_eq!(rotations.len(), n, "N rotations, no more");
let distinct: std::collections::HashSet<_> = rotations.iter().collect();
assert_eq!(distinct.len(), n, "no account rotated twice — claims were disjoint");
assert!(refs.iter().all(|r| distinct.contains(&("gmail".to_string(), r.clone()))));
}
#[tokio::test]
async fn ioa9_invalid_grant_expires_account() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::happy("dead@example.com", "client-gmail", "nonce-x");
transport.fail_with(TransportFailure::new(TransportFailureKind::InvalidGrant, "refresh token expired"));
let account = seed_active_account(&pool, company, "dead@example.com", 300, &store).await;
let report = refresh_oauth_credentials(
&pool, company, ®istry(), &EndpointOverrides::default(), &clients(), &store, &transport,
&RefreshSchedule::default(),
)
.await
.expect("refresh run");
assert_eq!(report.expired, 1, "the dead refresh grant expired its account: {report:?}");
assert_eq!(store.rotate_count(), 0);
assert_eq!(store.revoke_count(), 0, "expiry is not a store revoke — lazy store expiry owns it");
let (status, _, _) = account_state(&pool, account).await;
assert_eq!(status, "expired", "terminal expired; the user must reconnect");
}
#[tokio::test]
async fn ioa9_store_outage_skips_not_expires() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::happy("flaky@example.com", "client-gmail", "nonce-x");
let account = seed_active_account(&pool, company, "flaky@example.com", 300, &store).await;
store.fail_with(OAuthCredentialFailure::transport("store unreachable"));
let report = refresh_oauth_credentials(
&pool, company, ®istry(), &EndpointOverrides::default(), &clients(), &store, &transport,
&RefreshSchedule::default(),
)
.await
.expect("refresh run");
assert_eq!(report.skipped, 1, "a transport-class store failure is retryable: {report:?}");
assert_eq!(report.expired, 0);
let (status, _, _) = account_state(&pool, account).await;
assert_eq!(status, "active", "an outage must not expire a healthy account");
sqlx::query("UPDATE integrations.integration_accounts SET status='revoked' WHERE id=$1")
.bind(account)
.execute(&pool)
.await
.unwrap();
}
#[tokio::test]
async fn ioa9_unstoreable_answer_skips() {
let _gate = sweep_gate().await;
let pool = pool().await;
let company = Uuid::new_v4();
let store = FakeStore::new();
let transport = FakeTransport::default();
transport.set_token_response(TokenResponse {
access_token: "NO-EXPIRY-TOKEN".into(),
refresh_token: Some("NO-EXPIRY-REFRESH".into()),
expires_in: None,
scope: None,
id_token: None,
token_type: Some("Bearer".into()),
});
let account = seed_active_account(&pool, company, "eternal@example.com", 300, &store).await;
let report = refresh_oauth_credentials(
&pool, company, ®istry(), &EndpointOverrides::default(), &clients(), &store, &transport,
&RefreshSchedule::default(),
)
.await
.expect("refresh run");
assert_eq!(store.rotate_count(), 0, "a permanent token is not a storeable value");
assert_eq!(report.refreshed, 0, "{report:?}");
let (status, _, _) = account_state(&pool, account).await;
assert_eq!(status, "active", "unstoreable is a provider-side refusal, not an account expiry");
sqlx::query("UPDATE integrations.integration_accounts SET status='revoked' WHERE id=$1")
.bind(account)
.execute(&pool)
.await
.unwrap();
}
#[tokio::test]
async fn ioa1_one_flow_serves_all_four_providers() {
let pool = pool().await;
let transport = FakeTransport::default();
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store);
let reg = registry();
let mut envelope_keys: Option<Vec<String>> = None;
let mut ids: Vec<Uuid> = Vec::new();
for provider in reg.providers() {
let account_ref = if provider.contains("calendar") { "calendar-user-42".to_string() } else { format!("user-{provider}@example.com") };
let resp = svc
.authorize(AuthorizeRequest { provider: provider.into(), account_ref, scopes: None })
.await
.unwrap_or_else(|e| panic!("{provider}: authorize through the one flow failed: {e}"));
ids.push(resp.account_id);
let claims = decode_state(&state_of(&resp.authorize_url));
let mut keys: Vec<String> = claims.as_object().unwrap().keys().cloned().collect();
keys.sort();
match &envelope_keys {
None => envelope_keys = Some(keys),
Some(expected) => assert_eq!(&keys, expected, "{provider}: state envelope diverged from the one generation"),
}
assert_eq!(claims["account_id"].as_str().unwrap(), resp.account_id.to_string(), "{provider}: state binds the account");
assert_eq!(claims["provider"].as_str().unwrap(), provider, "{provider}: state names its provider");
let url = url::Url::parse(&resp.authorize_url).unwrap();
let params: Vec<String> = {
let mut v: Vec<String> = url.query_pairs().map(|(k, _)| k.to_string()).collect();
v.sort();
v
};
for must in ["access_type", "client_id", "prompt", "redirect_uri", "response_type", "scope", "state"] {
assert!(params.contains(&must.to_string()), "{provider}: authorize URL misses {must}: {params:?}");
}
assert_eq!(url.query_pairs().find(|(k, _)| k == "response_type").unwrap().1, "code");
assert_eq!(
url.query_pairs().find(|(k, _)| k == "redirect_uri").unwrap().1,
"https://api.example.test/api/v1/integrations/oauth/callback",
"the redirect target is the fixed deployment constant"
);
let host = url.host_str().unwrap();
let google_family = host.ends_with("google.com");
let microsoft_family = host.ends_with("microsoftonline.com");
let in_family = if provider.starts_with("google") || provider == "gmail" {
google_family
} else {
microsoft_family
};
assert!(in_family, "{provider}: authorize host {host} outside its registry family");
let pkce_in_url = params.iter().any(|k| k == "code_challenge");
assert_eq!(
pkce_in_url, reg.lookup(provider).unwrap().pkce_s256,
"{provider}: PKCE presence must come from registry data, not flow code"
);
}
let rows = account_rows(&pool, &ids).await;
assert_eq!(rows.len(), 4);
assert!(rows.iter().all(|(_, s, _)| s == "pending"));
}
#[tokio::test]
async fn ioa1_mail_providers_require_email_account_ref() {
let pool = pool().await;
let svc = oauth_service(&pool, FakeTransport::default(), FakeStore::new());
let r = svc
.authorize(AuthorizeRequest { provider: "gmail".into(), account_ref: "not-an-email".into(), scopes: None })
.await;
assert!(r.is_err(), "a gmail account_ref that is not an address must be refused at initiation");
}
fn state_of(authorize_url: &str) -> String {
url::Url::parse(authorize_url)
.unwrap()
.query_pairs()
.find(|(k, _)| k == "state")
.expect("authorize URL carries the state")
.1
.to_string()
}
#[tokio::test]
async fn ioa2_state_fence_rejects_every_forgery_class() {
let pool = pool().await;
let company = Uuid::new_v4();
let transport = FakeTransport::happy("fence@example.com", "client-gmail", "nonce-x");
let store = FakeStore::new();
let svc = oauth_service(&pool, transport.clone(), store.clone());
let resp = svc
.authorize(AuthorizeRequest { provider: "gmail".into(), account_ref: "fence@example.com".into(), scopes: None })
.await
.expect("authorize");
let good_state = state_of(&resp.authorize_url);
let snapshot = account_rows(&pool, &[resp.account_id]).await;
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: "garbage".into() }).await;
assert!(matches!(r, Err(backbone_integrations::application::service::integrations_oauth::OauthError::State(_))), "unsigned garbage state accepted: {r:?}");
let tampered = flip_a_byte(&good_state, 0);
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: tampered }).await;
assert!(matches!(r, Err(backbone_integrations::application::service::integrations_oauth::OauthError::State(_))), "tampered state accepted: {r:?}");
let expired = mint_state(resp.account_id, "gmail", "nonce-x", Utc::now().timestamp() - 100);
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: expired }).await;
assert!(matches!(r, Err(backbone_integrations::application::service::integrations_oauth::OauthError::State(_))), "expired state accepted: {r:?}");
let alien = mint_state(resp.account_id, "outlook", "nonce-x", Utc::now().timestamp() + 300);
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: alien }).await;
assert!(matches!(r, Err(backbone_integrations::application::service::integrations_oauth::OauthError::State(_))), "cross-provider state accepted: {r:?}");
assert_eq!(account_rows(&pool, &[resp.account_id]).await, snapshot, "a forged state changed rows");
assert!(store.calls().is_empty(), "a forged state touched the credential store");
assert_eq!(transport.call_count(), 0, "a forged state reached the network");
}
fn flip_a_byte(state: &str, seg: usize) -> String {
let parts: Vec<&str> = state.splitn(2, '.').collect();
let mut half = parts[seg].to_string();
let mut chars: Vec<char> = half.chars().collect();
let i = chars.len() / 2;
chars[i] = match chars[i] {
'A' => 'B',
c => char::from(c as u8 - 1),
};
half = chars.into_iter().collect();
if seg == 0 { format!("{half}.{}", parts[1]) } else { format!("{}.{half}", parts[0]) }
}
#[tokio::test]
async fn ioa3_flipped_signature_byte_rejected_no_side_effects() {
let pool = pool().await;
let company = Uuid::new_v4();
let transport = FakeTransport::happy("csrf@example.com", "client-gmail", "nonce-x");
let store = FakeStore::new();
let svc = oauth_service(&pool, transport.clone(), store.clone());
let resp = svc
.authorize(AuthorizeRequest { provider: "gmail".into(), account_ref: "csrf@example.com".into(), scopes: None })
.await
.expect("authorize");
let state = state_of(&resp.authorize_url);
let flipped = flip_a_byte(&state, 1);
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: flipped }).await;
assert!(r.is_err(), "a flipped signature byte passed verification");
let rows = account_rows(&pool, &[resp.account_id]).await;
assert_eq!(rows.len(), 1, "the pending account is the only row");
assert_eq!(rows[0].1, "pending", "the account stays pending");
assert!(store.calls().is_empty(), "no credential was issued for a forged state");
assert_eq!(transport.call_count(), 0);
}
#[tokio::test]
async fn ioa4_control_matching_audience_and_nonce_complete() {
let pool = pool().await;
let company = Uuid::new_v4();
let resp_authorize = run_authorize(&pool, "match@example.com").await;
let claims = decode_state(&resp_authorize.state);
let nonce = claims["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("match@example.com", "client-gmail", &nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let out = svc
.complete(company, CompleteRequest { code: "AUTH-CODE".into(), state: resp_authorize.state })
.await
.expect("the control round trip must complete");
assert_eq!(out.status, "active");
assert_eq!(store.issue_count(), 1);
}
#[tokio::test]
async fn ioa4_realistic_echo_request_carries_nonce_idtoken_completes() {
let pool = pool().await;
let company = Uuid::new_v4();
let svc = oauth_service(&pool, FakeTransport::default(), FakeStore::new());
let resp = svc
.authorize(AuthorizeRequest { provider: "gmail".into(), account_ref: "echo@example.com".into(), scopes: None })
.await
.expect("authorize");
let url = url::Url::parse(&resp.authorize_url).unwrap();
let request_nonce = url
.query_pairs()
.find(|(k, _)| k == "nonce")
.expect("the authorize request must carry the nonce pair (the provider echoes what it receives)")
.1
.to_string();
assert!(!request_nonce.is_empty(), "the nonce pair must be non-empty");
let state = state_of(&resp.authorize_url);
assert_eq!(
decode_state(&state)["nonce"].as_str().unwrap(),
request_nonce,
"the request nonce and the signed state's nonce are the same mint"
);
let transport = FakeTransport::realistic_idtoken_echo("echo@example.com", "client-gmail", &request_nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let out = svc
.complete(company, CompleteRequest { code: "AUTH-CODE".into(), state })
.await
.expect("the realistic-echo round trip must complete via the id_token decode path");
assert_eq!(out.status, "active");
assert_eq!(store.issue_count(), 1);
}
#[tokio::test]
async fn ioa4_audience_and_nonce_mismatches_rejected_zero_writes() {
let pool = pool().await;
let company = Uuid::new_v4();
let a = run_authorize(&pool, "aud@example.com").await;
let nonce = decode_state(&a.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("aud@example.com", "attacker-client-id", &nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: a.state }).await;
assert!(r.is_err(), "token with the attacker's audience accepted");
assert_eq!(account_rows(&pool, &[a.account_id]).await.len(), 1);
assert_eq!(account_rows(&pool, &[a.account_id]).await[0].1, "pending", "audience mismatch must leave the account pending");
assert!(store.calls().is_empty() && store.rotate_count() == 0, "audience mismatch wrote to the store");
let b = run_authorize(&pool, "nonce@example.com").await;
let transport = FakeTransport::happy("nonce@example.com", "client-gmail", "a-different-nonce");
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: b.state }).await;
assert!(r.is_err(), "replayed nonce accepted");
assert_eq!(account_rows(&pool, &[a.account_id, b.account_id]).await[1].1, "pending", "nonce mismatch must leave the account pending");
assert!(store.calls().is_empty() && store.rotate_count() == 0, "nonce mismatch wrote to the store");
}
#[tokio::test]
async fn ioa5_identity_email_mismatch_rejected() {
let pool = pool().await;
let company = Uuid::new_v4();
let a = run_authorize(&pool, "victim@example.com").await;
let nonce = decode_state(&a.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("attacker@example.com", "client-gmail", &nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: a.state }).await;
assert!(r.is_err(), "identity email ≠ account_ref accepted (the token-substitution link attack)");
assert_eq!(account_rows(&pool, &[a.account_id]).await[0].1, "pending", "the victim's account stays pending");
assert!(store.calls().is_empty(), "no credential issued for a substituted identity");
let b = run_authorize(&pool, "unverified@example.com").await;
let nonce_b = decode_state(&b.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("unverified@example.com", "client-gmail", &nonce_b);
transport.set_identity(backbone_integrations::infrastructure::http::IdentityClaims {
email: Some("unverified@example.com".into()),
email_verified: Some(false),
audience: Some("client-gmail".into()),
nonce: Some(nonce_b),
sub: Some("sub-u".into()),
});
let store2 = FakeStore::new();
let svc2 = oauth_service(&pool, transport, store2.clone());
let r = svc2.complete(company, CompleteRequest { code: "c".into(), state: b.state }).await;
assert!(r.is_err(), "an unverified email passed the gauntlet");
assert!(store2.calls().is_empty());
}
#[tokio::test]
async fn ioa6_stored_credential_and_mirror_carry_honest_expiry() {
let pool = pool().await;
let company = Uuid::new_v4();
let a = run_authorize(&pool, "honest@example.com").await;
let nonce = decode_state(&a.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("honest@example.com", "client-gmail", &nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let before = Utc::now();
let out = svc.complete(company, CompleteRequest { code: "c".into(), state: a.state }).await.expect("complete");
let after = Utc::now();
let expect_low = before + Duration::seconds(86_400);
let expect_high = after + Duration::seconds(86_400);
assert!(out.expires_at > expect_low && out.expires_at < expect_high, "outcome expiry is now + expires_in");
match store.calls().first() {
Some(StoreCall::Issued { purpose, expires_at, .. }) => {
assert_eq!(purpose, PURPOSE_OAUTH_TOKEN, "the bundle is stored as an oauth_token credential");
assert!(*expires_at > expect_low && *expires_at < expect_high, "the STORED expiry is now + expires_in (got {expires_at:?})");
}
other => panic!("no issue call recorded: {other:?}"),
}
let (status, mirror, refreshed_at) = sqlx::query_as::<_, (String, Option<chrono::DateTime<Utc>>, Option<chrono::DateTime<Utc>>)>(
"SELECT status::text, expires_at, last_refreshed_at FROM integrations.integration_accounts WHERE id=$1",
)
.bind(a.account_id)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(status, "active");
let mirror = mirror.expect("an active account mirrors its expiry — never NULL on this path");
assert!(mirror > expect_low && mirror < expect_high, "the mirror carries the same honest expiry");
assert!(refreshed_at.is_some());
}
#[tokio::test]
async fn ioa6_permanent_token_refused_as_unstoreable() {
let pool = pool().await;
let company = Uuid::new_v4();
let a = run_authorize(&pool, "eternal2@example.com").await;
let nonce = decode_state(&a.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("eternal2@example.com", "client-gmail", &nonce);
transport.set_token_response(TokenResponse {
access_token: "NO-EXPIRY".into(),
refresh_token: Some("NO-EXPIRY-REFRESH".into()),
expires_in: None,
scope: None,
id_token: None,
token_type: Some("Bearer".into()),
});
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
let r = svc.complete(company, CompleteRequest { code: "c".into(), state: a.state }).await;
assert!(matches!(r, Err(backbone_integrations::application::service::integrations_oauth::OauthError::Unstoreable(_))), "a permanent token was stored: {r:?}");
assert!(store.calls().is_empty(), "nothing reached the store");
assert_eq!(account_rows(&pool, &[a.account_id]).await[0].1, "pending", "the account stays pending");
}
#[tokio::test]
async fn ioa7_all_credential_contact_rides_the_port() {
let pool = pool().await;
let company = Uuid::new_v4();
let a = run_authorize(&pool, "port@example.com").await;
let nonce = decode_state(&a.state)["nonce"].as_str().unwrap().to_string();
let transport = FakeTransport::happy("port@example.com", "client-gmail", &nonce);
let store = FakeStore::new();
let svc = oauth_service(&pool, transport, store.clone());
svc.complete(company, CompleteRequest { code: "c".into(), state: a.state }).await.expect("complete");
svc.disconnect(company, a.account_id).await.expect("disconnect");
let verbs: Vec<&'static str> = store.calls().iter().map(call_verb).collect();
assert_eq!(verbs, vec!["issue", "revoke"], "exactly an issue then a revoke — no other store contact: {verbs:?}");
let row_status = account_rows(&pool, &[a.account_id]).await[0].1.clone();
assert_eq!(row_status, "revoked", "disconnect is terminal on the account");
let secretish: Vec<String> = sqlx::query_scalar(
r#"SELECT column_name FROM information_schema.columns
WHERE table_schema='integrations' AND table_name='integration_accounts'
AND (column_name ILIKE '%token%' OR column_name ILIKE '%secret%'
OR column_name ILIKE '%cipher%' OR column_name ILIKE '%password%')"#,
)
.fetch_all(&pool)
.await
.unwrap();
assert!(secretish.is_empty(), "token-shaped columns on integration_accounts: {secretish:?}");
}
fn call_verb(c: &StoreCall) -> &'static str {
match c {
StoreCall::Issued { .. } => "issue",
StoreCall::Read { .. } => "read",
StoreCall::Rotated { .. } => "rotate",
StoreCall::Revoked { .. } => "revoke",
}
}
#[tokio::test]
async fn ioa10_callback_page_writes_nothing() {
let pool = pool().await;
let a = run_authorize(&pool, "cb@example.com").await;
let store = FakeStore::new();
let svc = oauth_service(&pool, FakeTransport::default(), store.clone());
let snapshot = account_rows(&pool, &[a.account_id]).await;
let page = svc.callback_page("AUTH-CODE", &a.state).expect("callback page");
assert_eq!(account_rows(&pool, &[a.account_id]).await, snapshot, "the GET callback changed rows");
assert!(store.calls().is_empty());
assert!(page.contains("form"), "the page carries the auto-POST form");
assert!(page.contains("complete"), "the form targets the completion action");
assert!(page.contains("AUTH-CODE") && page.contains(&a.state), "the form carries code + state");
assert!(!page.to_lowercase().contains("http-equiv=\"refresh\""), "no meta-refresh redirect");
assert!(!page.contains("Location:"), "no redirect header in the page");
let r = svc.callback_page("c", "garbage");
assert!(r.is_err());
assert_eq!(account_rows(&pool, &[a.account_id]).await, snapshot);
}
struct Authorized {
account_id: Uuid,
state: String,
}
async fn run_authorize(pool: &PgPool, account_ref: &str) -> Authorized {
let svc = oauth_service(pool, FakeTransport::default(), FakeStore::new());
let resp = svc
.authorize(AuthorizeRequest { provider: "gmail".into(), account_ref: account_ref.into(), scopes: None })
.await
.expect("authorize");
Authorized { account_id: resp.account_id, state: state_of(&resp.authorize_url) }
}
fn principal(company: Uuid, permissions: &[&str]) -> OAuthPrincipal {
OAuthPrincipal {
company_id: company,
user_id: Some(Uuid::new_v4()),
permissions: permissions.iter().map(|p| p.to_string()).collect(),
}
}
async fn body_text(response: axum::response::Response) -> String {
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.expect("response body");
String::from_utf8_lossy(&bytes).into_owned()
}
#[tokio::test]
async fn ioa2_old_flow_routes_are_absent_from_the_surface() {
let pool = pool().await;
let company = Uuid::new_v4();
let router = create_oauth_routes(Arc::new(oauth_service(&pool, FakeTransport::default(), FakeStore::new())));
let old_flow_candidates = [
"/oauth/google/authorize",
"/oauth/google/callback",
"/oauth/google/complete",
"/oauth/microsoft/authorize",
"/oauth/microsoft/callback",
"/oauth/gmail/authorize",
"/oauth/gmail/callback",
"/oauth/gmail/complete",
"/oauth/outlook/authorize",
"/oauth/outlook/callback",
"/oauth/google_calendar/authorize",
"/oauth/google_calendar/callback",
"/oauth/microsoft_calendar/authorize",
"/oauth/microsoft_calendar/callback",
"/oauth/connect",
"/oauth/token",
];
for path in old_flow_candidates {
for method in [Method::GET, Method::POST] {
let req = Request::builder()
.method(method.clone())
.uri(path)
.extension(principal(company, &["write:integrations", "delete:integrations"]))
.body(Body::empty())
.expect("request");
let status = router.clone().oneshot(req).await.expect("oneshot").status();
assert_eq!(
status,
StatusCode::NOT_FOUND,
"{method} {path}: an old-flow route must not exist on the one-generation surface"
);
}
}
let id = Uuid::new_v4();
let wrong_method_probes: Vec<(Method, String)> = vec![
(Method::GET, "/oauth/authorize".into()),
(Method::GET, "/oauth/complete".into()),
(Method::GET, format!("/oauth/{id}/disconnect")),
(Method::POST, "/oauth/callback".into()),
(Method::POST, format!("/oauth/{id}/status")),
];
for (method, path) in wrong_method_probes {
let req = Request::builder()
.method(method.clone())
.uri(&path)
.extension(principal(company, &["write:integrations", "delete:integrations"]))
.body(Body::empty())
.expect("request");
let status = router.clone().oneshot(req).await.expect("oneshot").status();
assert_eq!(
status,
StatusCode::METHOD_NOT_ALLOWED,
"{method} {path}: the new-flow verb must exist (wrong-method probe expected 405)"
);
}
}
#[tokio::test]
async fn ioa2_forged_state_on_the_http_verb_is_a_400_with_zero_writes() {
let pool = pool().await;
let company = Uuid::new_v4();
let transport = FakeTransport::default();
let store = FakeStore::new();
let router = create_oauth_routes(Arc::new(oauth_service(&pool, transport.clone(), store.clone())));
let req = Request::builder()
.method(Method::POST)
.uri("/oauth/complete")
.header("content-type", "application/x-www-form-urlencoded")
.extension(principal(company, &["write:integrations"]))
.body(Body::from("code=attacker-code&state=Zm9yZ2Vk.dmFjZQ"))
.expect("request");
let response = router.oneshot(req).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::BAD_REQUEST, "forged state must be a 400");
let body = body_text(response).await;
assert!(body.contains("OAUTH_STATE_REJECTED"), "the rejection names the state fence: {body}");
assert!(store.calls().is_empty(), "zero credential verbs");
assert_eq!(transport.call_count(), 0, "zero outbound calls");
}
#[tokio::test]
async fn ioa11_every_authed_verb_fails_closed_without_a_principal() {
let pool = pool().await;
let router = create_oauth_routes(Arc::new(oauth_service(&pool, FakeTransport::default(), FakeStore::new())));
let id = Uuid::new_v4();
let probes: Vec<(Method, String, Body)> = vec![
(
Method::POST,
"/oauth/authorize".into(),
Body::from(r#"{"provider":"gmail","account_ref":"nobody@example.com"}"#),
),
(
Method::POST,
"/oauth/complete".into(),
Body::from("code=x&state=y"),
),
(Method::POST, format!("/oauth/{id}/disconnect"), Body::empty()),
(Method::GET, format!("/oauth/{id}/status"), Body::empty()),
];
for (method, path, body) in probes {
let req = Request::builder()
.method(method.clone())
.uri(&path)
.header("content-type", "application/json")
.body(body)
.expect("request");
let response = router.clone().oneshot(req).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::UNAUTHORIZED, "{method} {path} without a principal");
let body = body_text(response).await;
assert!(body.contains("OAUTH_UNAUTHENTICATED"), "{method} {path}: {body}");
}
let req = Request::builder()
.method(Method::GET)
.uri("/oauth/callback?code=c&state=garbage")
.body(Body::empty())
.expect("request");
let response = router.oneshot(req).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::BAD_REQUEST, "public callback rejects garbage with 400");
assert!(response.headers().get("location").is_none(), "no redirect on the callback");
}
#[tokio::test]
async fn ioa11_verbs_enforce_their_own_permissions() {
let pool = pool().await;
let company = Uuid::new_v4();
let transport = FakeTransport::default();
let store = FakeStore::new();
let router = create_oauth_routes(Arc::new(oauth_service(&pool, transport.clone(), store.clone())));
let id = Uuid::new_v4();
let bare = principal(company, &[]);
let authorize = Request::builder()
.method(Method::POST)
.uri("/oauth/authorize")
.header("content-type", "application/json")
.extension(bare.clone())
.body(Body::from(r#"{"provider":"gmail","account_ref":"nobody@example.com"}"#))
.expect("request");
let response = router.clone().oneshot(authorize).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::FORBIDDEN, "authorize without write:integrations");
let disconnect = Request::builder()
.method(Method::POST)
.uri(format!("/oauth/{id}/disconnect"))
.extension(bare.clone())
.body(Body::empty())
.expect("request");
let response = router.clone().oneshot(disconnect).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::FORBIDDEN, "disconnect without delete:integrations");
let writer = principal(company, &["write:integrations"]);
let disconnect = Request::builder()
.method(Method::POST)
.uri(format!("/oauth/{id}/disconnect"))
.extension(writer.clone())
.body(Body::empty())
.expect("request");
let response = router.clone().oneshot(disconnect).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::FORBIDDEN, "write grant does not confer delete");
let full = principal(company, &["write:integrations", "delete:integrations"]);
let authorize = Request::builder()
.method(Method::POST)
.uri("/oauth/authorize")
.header("content-type", "application/json")
.extension(full.clone())
.body(Body::from(r#"{"provider":"gmail","account_ref":"route-probe@example.com"}"#))
.expect("request");
let response = router.clone().oneshot(authorize).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::OK, "authorize with the write grant");
let outcome: serde_json::Value =
serde_json::from_str(&body_text(response).await).expect("authorize response JSON");
let account_id: Uuid = outcome["data"]["account_id"]
.as_str()
.expect("account_id in the authorize response")
.parse()
.expect("account_id is a uuid");
let rows = account_rows(&pool, &[account_id]).await;
assert_eq!(rows.len(), 1, "the HTTP authorize wrote exactly one pending row");
assert_eq!(rows[0].0, account_id);
assert_eq!(rows[0].1, "pending");
let authorize_url = outcome["data"]["authorize_url"]
.as_str()
.expect("authorize_url in the authorize response")
.to_string();
let state = state_of(&authorize_url);
let nonce = decode_state(&state)["nonce"].as_str().expect("nonce in the state").to_string();
transport.set_token_response(TokenResponse {
access_token: "FAKE-ACCESS-TOKEN".into(),
refresh_token: Some("FAKE-REFRESH-TOKEN".into()),
expires_in: Some(86_400),
scope: Some("https://mail.google.com/".into()),
id_token: Some("FAKE-ID-TOKEN".into()),
token_type: Some("Bearer".into()),
});
transport.set_identity(IdentityClaims {
sub: Some("sub-route-probe@example.com".into()),
email: Some("route-probe@example.com".into()),
email_verified: Some(true),
audience: Some("client-gmail".into()),
nonce: Some(nonce.clone()),
});
let complete = Request::builder()
.method(Method::POST)
.uri("/oauth/complete")
.header("content-type", "application/x-www-form-urlencoded")
.extension(full.clone())
.body(Body::from(format!("code=provider-code&state={state}")))
.expect("request");
let response = router.clone().oneshot(complete).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::OK, "complete with the write grant");
let rows = account_rows(&pool, &[account_id]).await;
assert_eq!(rows[0].1, "active", "the HTTP complete activated the account");
assert_eq!(call_verb(&store.calls()[0]), "issue", "the credential was issued through the port");
let disconnect = Request::builder()
.method(Method::POST)
.uri(format!("/oauth/{account_id}/disconnect"))
.extension(full.clone())
.body(Body::empty())
.expect("request");
let response = router.clone().oneshot(disconnect).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::NO_CONTENT, "disconnect with the delete grant");
assert_eq!(store.calls().len(), 2, "the disconnect revoked through the port");
assert_eq!(call_verb(&store.calls()[1]), "revoke");
let status = Request::builder()
.method(Method::GET)
.uri(format!("/oauth/{account_id}/status"))
.extension(bare)
.body(Body::empty())
.expect("request");
let response = router.oneshot(status).await.expect("oneshot");
assert_eq!(response.status(), StatusCode::OK, "status needs no integration grant");
}