#![cfg(all(feature = "store-postgres", feature = "secret-store"))]
use std::env::var;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use ironflow_core::auth_proxy::{
AuthProxyError, AuthProxyRegistry, CredentialKind, IssuedToken, ProxyCredential,
TokenRejection, TokenRequest, token_id,
};
use ironflow_store::crypto::KeyRing;
use ironflow_store::postgres::PostgresStore;
use sqlx::postgres::PgPoolOptions;
use sqlx::{PgPool, Row, query};
use uuid::Uuid;
const CREDENTIAL: &str = "sk-ant-oat01-postgres-test";
const RAW_ROW: &str =
"SELECT id, encrypted_credential, run_id FROM ironflow.auth_proxy_grants WHERE run_id = $1";
const ROW_AS_TEXT: &str =
"SELECT g::text AS row_text FROM ironflow.auth_proxy_grants g WHERE run_id = $1";
fn database_url() -> String {
var("DATABASE_URL").expect("DATABASE_URL must be set")
}
fn hex_key(byte: u8) -> String {
format!("{byte:02x}").repeat(32)
}
fn ring(byte: u8) -> KeyRing {
KeyRing::from_spec(&format!("1:{}", hex_key(byte)), None).expect("valid ring")
}
async fn store_with_key(byte: u8) -> PostgresStore {
let mut store = PostgresStore::new(&database_url())
.await
.expect("failed to connect to PostgreSQL");
store.set_key_ring(ring(byte));
store
}
fn registry(store: PostgresStore) -> AuthProxyRegistry {
AuthProxyRegistry::with_backend(Arc::new(store))
}
async fn raw_pool() -> PgPool {
PgPoolOptions::new()
.max_connections(1)
.connect(&database_url())
.await
.expect("failed to connect to PostgreSQL")
}
fn unique_run(label: &str) -> String {
format!("test-{label}-{}", Uuid::now_v7())
}
fn now() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock after the unix epoch")
.as_secs()
}
fn request(run_id: &str, expires_at: u64) -> TokenRequest {
TokenRequest {
run_id: run_id.to_string(),
step: "review".to_string(),
expires_at,
credential: ProxyCredential::new(CredentialKind::OauthToken, CREDENTIAL.to_string()),
}
}
async fn issue(registry: &AuthProxyRegistry, run_id: &str) -> IssuedToken {
let now = now();
registry
.issue(request(run_id, now + 600), now)
.await
.expect("issue")
}
async fn count_rows(pool: &PgPool, run_id: &str) -> i64 {
query("SELECT COUNT(*) AS count FROM ironflow.auth_proxy_grants WHERE run_id = $1")
.bind(run_id)
.fetch_one(pool)
.await
.expect("count")
.get("count")
}
fn contains(haystack: &[u8], needle: &[u8]) -> bool {
haystack
.windows(needle.len())
.any(|window| window == needle)
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn grant_survives_a_restart() {
let run_id = unique_run("restart");
let issued = issue(®istry(store_with_key(0xaa).await), &run_id).await;
let restarted = registry(store_with_key(0xaa).await);
let grant = restarted
.resolve(&issued.token, now())
.await
.expect("grant survives");
assert_eq!(grant.id, issued.id);
assert_eq!(grant.run_id, run_id);
assert_eq!(grant.step, "review");
assert_eq!(grant.credential.kind(), CredentialKind::OauthToken);
assert_eq!(grant.credential.expose(), CREDENTIAL);
restarted.revoke_run(&run_id).await.expect("cleanup");
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn two_replicas_share_revocation() {
let run_id = unique_run("replicas");
let a = registry(store_with_key(0xaa).await);
let b = registry(store_with_key(0xaa).await);
let first = issue(&a, &run_id).await;
let second = issue(&a, &run_id).await;
assert!(b.resolve(&first.token, now()).await.is_ok());
assert_eq!(b.revoke_run(&run_id).await.expect("revoke run"), 2);
assert_eq!(
a.resolve(&first.token, now()).await.unwrap_err(),
TokenRejection::Unknown
);
assert_eq!(
a.resolve(&second.token, now()).await.unwrap_err(),
TokenRejection::Unknown
);
assert_eq!(a.revoke_run(&run_id).await.expect("revoke run"), 0);
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn single_revocation_is_seen_by_another_replica() {
let run_id = unique_run("revoke");
let a = registry(store_with_key(0xaa).await);
let b = registry(store_with_key(0xaa).await);
let issued = issue(&a, &run_id).await;
assert!(b.revoke(&issued.id).await.expect("revoke"));
assert!(!a.revoke(&issued.id).await.expect("revoke again"));
assert_eq!(
a.resolve(&issued.token, now()).await.unwrap_err(),
TokenRejection::Unknown
);
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn token_is_never_stored_and_credential_is_encrypted() {
let run_id = unique_run("at-rest");
let registry = registry(store_with_key(0xaa).await);
let issued = issue(®istry, &run_id).await;
let pool = raw_pool().await;
let row = query(RAW_ROW)
.bind(&run_id)
.fetch_one(&pool)
.await
.expect("row exists");
let id: String = row.get("id");
let encrypted: Vec<u8> = row.get("encrypted_credential");
let stored_run: String = row.get("run_id");
assert_eq!(id, token_id(&issued.token));
assert_ne!(id, issued.token);
assert_eq!(stored_run, run_id);
assert!(!contains(&encrypted, CREDENTIAL.as_bytes()));
assert!(!contains(&encrypted, issued.token.as_bytes()));
let text: String = query(ROW_AS_TEXT)
.bind(&run_id)
.fetch_one(&pool)
.await
.expect("row exists")
.get("row_text");
assert!(!text.contains(&issued.token), "{text}");
assert!(!text.contains(CREDENTIAL), "{text}");
let credential_hex: String = CREDENTIAL.bytes().map(|b| format!("{b:02x}")).collect();
assert!(!text.contains(&credential_hex), "{text}");
registry.revoke_run(&run_id).await.expect("cleanup");
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn purge_expired_deletes_expired_rows() {
let run_id = unique_run("purge");
let registry = registry(store_with_key(0xaa).await);
let pool = raw_pool().await;
let past_now = now() - 7_200;
registry
.issue(request(&run_id, past_now + 10), past_now)
.await
.expect("issue expired");
let live = issue(®istry, &run_id).await;
assert_eq!(count_rows(&pool, &run_id).await, 2);
assert!(registry.purge_expired(past_now + 10).await.expect("purge") >= 1);
assert_eq!(count_rows(&pool, &run_id).await, 1);
assert!(registry.resolve(&live.token, now()).await.is_ok());
registry.revoke_run(&run_id).await.expect("cleanup");
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn expired_grant_is_rejected_and_removed_on_resolve() {
let run_id = unique_run("expired");
let registry = registry(store_with_key(0xaa).await);
let pool = raw_pool().await;
let past_now = now() - 60;
let issued = registry
.issue(request(&run_id, past_now + 10), past_now)
.await
.expect("issue");
assert_eq!(
registry.resolve(&issued.token, now()).await.unwrap_err(),
TokenRejection::Expired
);
assert_eq!(count_rows(&pool, &run_id).await, 0);
assert_eq!(
registry.resolve(&issued.token, now()).await.unwrap_err(),
TokenRejection::Unknown
);
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn insert_without_key_ring_fails() {
let run_id = unique_run("no-key");
let store = PostgresStore::new(&database_url())
.await
.expect("failed to connect to PostgreSQL");
let registry = registry(store);
let now = now();
let result = registry.issue(request(&run_id, now + 600), now).await;
assert!(
matches!(result, Err(AuthProxyError::Backend(_))),
"{result:?}"
);
assert_eq!(count_rows(&raw_pool().await, &run_id).await, 0);
}
#[tokio::test]
#[ignore = "requires a live PostgreSQL database"]
async fn wrong_key_ring_cannot_read_grant() {
let run_id = unique_run("wrong-key");
let writer = registry(store_with_key(0xaa).await);
let issued = issue(&writer, &run_id).await;
let reader = registry(store_with_key(0xbb).await);
match reader.resolve(&issued.token, now()).await {
Err(TokenRejection::Unavailable(message)) => {
assert!(!message.contains(CREDENTIAL), "{message}");
}
Ok(grant) => panic!("grant {} decrypted with the wrong key", grant.id),
Err(other) => panic!("expected Unavailable, got {other:?}"),
}
writer.revoke_run(&run_id).await.expect("cleanup");
}