#![allow(clippy::unwrap_used)] #![allow(clippy::missing_panics_doc)] #![allow(missing_docs)]
use std::sync::Arc;
use fraiseql_server::usage::{
aggregator::{PostgresBackend, UsageAggregator},
events::MutationAuditEvent,
};
use sqlx::PgPool;
async fn setup_pg() -> (PgPool, fraiseql_test_support::Service) {
let svc = fraiseql_test_support::postgres()
.await
.expect("DATABASE_URL must be set (or enable fraiseql-test-support/local-testcontainers)");
let pool = PgPool::connect(svc.url()).await.unwrap();
sqlx::query("DROP TABLE IF EXISTS fraiseql_usage_counters CASCADE")
.execute(&pool)
.await
.unwrap();
(pool, svc)
}
fn event(tenant: &str, period: &str, entity: &str) -> MutationAuditEvent {
MutationAuditEvent::new(format!("create_{entity}"), entity, "create", tenant, period)
}
#[tokio::test]
async fn test_postgres_backend_flush_and_load_round_trip() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.unwrap();
agg.record(&event("acme", "2026-05", "User"));
agg.record(&event("acme", "2026-05", "User"));
agg.record(&event("acme", "2026-05", "Order"));
agg.flush_to_backend().await.unwrap();
let new_agg = UsageAggregator::new_with_backend(backend.clone());
new_agg.load_from_backend().await.unwrap();
assert_eq!(new_agg.query("acme", "2026-05").mutations["User"], 2);
assert_eq!(new_agg.query("acme", "2026-05").mutations["Order"], 1);
}
#[tokio::test]
async fn test_postgres_backend_flush_is_idempotent() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.unwrap();
agg.record(&event("t1", "2026-05", "Widget"));
agg.record(&event("t1", "2026-05", "Widget"));
agg.flush_to_backend().await.unwrap();
agg.flush_to_backend().await.unwrap();
let new_agg = UsageAggregator::new_with_backend(backend.clone());
new_agg.load_from_backend().await.unwrap();
assert_eq!(new_agg.query("t1", "2026-05").mutations["Widget"], 2);
}
#[tokio::test]
async fn test_postgres_backend_load_merges_with_inflight() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.unwrap();
for _ in 0..3 {
agg.record(&event("tenant", "2026-05", "Thing"));
}
agg.flush_to_backend().await.unwrap();
let new_agg = UsageAggregator::new_with_backend(backend.clone());
new_agg.record(&event("tenant", "2026-05", "Thing"));
new_agg.record(&event("tenant", "2026-05", "Thing"));
new_agg.load_from_backend().await.unwrap();
assert_eq!(new_agg.query("tenant", "2026-05").mutations["Thing"], 5);
}
#[tokio::test]
async fn test_postgres_backend_tenant_isolation() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.unwrap();
agg.record(&event("tenant_a", "2026-05", "User"));
agg.record(&event("tenant_b", "2026-05", "User"));
agg.record(&event("tenant_b", "2026-05", "User"));
agg.flush_to_backend().await.unwrap();
let new_agg = UsageAggregator::new_with_backend(backend.clone());
new_agg.load_from_backend().await.unwrap();
assert_eq!(new_agg.query("tenant_a", "2026-05").mutations["User"], 1);
assert_eq!(new_agg.query("tenant_b", "2026-05").mutations["User"], 2);
}
#[tokio::test]
async fn test_postgres_backend_empty_load_is_noop() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let agg = UsageAggregator::new_with_backend(backend);
agg.load_from_backend().await.unwrap();
assert_eq!(agg.entry_count(), 0);
}
struct UnreadablePostgres(PostgresBackend);
#[async_trait::async_trait]
impl fraiseql_server::usage::aggregator::UsageBackend for UnreadablePostgres {
async fn flush_deltas(
&self,
deltas: &std::collections::HashMap<(String, String, String), u64>,
) -> Result<(), String> {
self.0.flush_deltas(deltas).await
}
async fn load(
&self,
) -> Result<std::collections::HashMap<(String, String, String), u64>, String> {
Err("statement timeout".to_string())
}
}
async fn persisted_count(pool: &PgPool, tenant: &str, period: &str, entity: &str) -> i64 {
sqlx::query_scalar::<_, i64>(
"SELECT count FROM fraiseql_usage_counters
WHERE tenant_id = $1 AND period = $2 AND entity_type = $3",
)
.bind(tenant)
.bind(period)
.bind(entity)
.fetch_one(pool)
.await
.unwrap()
}
#[tokio::test]
async fn a_failed_startup_load_cannot_destroy_the_persisted_month() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let seed = UsageAggregator::new_with_backend(backend.clone());
seed.load_from_backend().await.unwrap();
for _ in 0..413 {
seed.record(&event("acme", "2026-07", "Order"));
}
seed.flush_to_backend().await.unwrap();
assert_eq!(persisted_count(&pool, "acme", "2026-07", "Order").await, 413);
let restarted = UsageAggregator::new_with_backend(Arc::new(UnreadablePostgres(
PostgresBackend::new(pool.clone()).await.unwrap(),
)));
restarted.load_from_backend().await.unwrap_err();
for _ in 0..12 {
restarted.record(&event("acme", "2026-07", "Order"));
}
restarted.flush_to_backend().await.unwrap_err();
assert_eq!(
persisted_count(&pool, "acme", "2026-07", "Order").await,
413,
"#861: a process that could not read the counters must not overwrite them"
);
}
#[tokio::test]
async fn replicas_sum_their_intervals_in_postgres() {
let (pool, _container) = setup_pg().await;
let backend = Arc::new(PostgresBackend::new(pool.clone()).await.unwrap());
let seed = UsageAggregator::new_with_backend(backend.clone());
seed.load_from_backend().await.unwrap();
for _ in 0..1000 {
seed.record(&event("acme", "2026-07", "Order"));
}
seed.flush_to_backend().await.unwrap();
for own in [7_u32, 5, 3] {
let replica = UsageAggregator::new_with_backend(backend.clone());
replica.load_from_backend().await.unwrap();
for _ in 0..own {
replica.record(&event("acme", "2026-07", "Order"));
}
replica.flush_to_backend().await.unwrap();
}
assert_eq!(
persisted_count(&pool, "acme", "2026-07", "Order").await,
1015,
"#861: an absolute UPSERT kept only the last replica's total (1003), losing 12"
);
}