mod aggregator_tests {
use std::collections::HashMap;
use super::super::{aggregator::*, events::MutationAuditEvent};
fn event(tenant: &str, period: &str, entity: &str) -> MutationAuditEvent {
MutationAuditEvent {
mutation_name: format!("create_{entity}"),
entity_type: entity.to_owned(),
operation: "create".to_owned(),
tenant_id: tenant.to_owned(),
period: period.to_owned(),
}
}
#[test]
fn test_record_and_query_single_tenant() {
let agg = UsageAggregator::new();
for _ in 0..4 {
agg.record(&event("tenant_a", "2026-05", "User"));
}
for _ in 0..3 {
agg.record(&event("tenant_a", "2026-05", "Order"));
}
let summary = agg.query("tenant_a", "2026-05");
assert_eq!(summary.mutations.get("User"), Some(&4));
assert_eq!(summary.mutations.get("Order"), Some(&3));
}
#[test]
fn test_record_and_query_two_tenants() {
let agg = UsageAggregator::new();
for _ in 0..5 {
agg.record(&event("tenant_a", "2026-05", "User"));
}
for _ in 0..2 {
agg.record(&event("tenant_b", "2026-05", "User"));
}
for _ in 0..3 {
agg.record(&event("tenant_b", "2026-05", "Product"));
}
let a = agg.query("tenant_a", "2026-05");
assert_eq!(a.mutations.get("User"), Some(&5));
assert_eq!(a.mutations.get("Product"), None);
let b = agg.query("tenant_b", "2026-05");
assert_eq!(b.mutations.get("User"), Some(&2));
assert_eq!(b.mutations.get("Product"), Some(&3));
}
#[test]
fn test_record_across_periods_does_not_bleed() {
let agg = UsageAggregator::new();
for _ in 0..10 {
agg.record(&event("t1", "2026-04", "Widget"));
}
for _ in 0..3 {
agg.record(&event("t1", "2026-05", "Widget"));
}
assert_eq!(agg.query("t1", "2026-04").mutations.get("Widget"), Some(&10));
assert_eq!(agg.query("t1", "2026-05").mutations.get("Widget"), Some(&3));
}
#[test]
fn test_record_10_events_across_2_tenants_3_entities() {
let agg = UsageAggregator::new();
let events = [
("tenant_a", "Alpha"),
("tenant_a", "Beta"),
("tenant_a", "Alpha"),
("tenant_b", "Gamma"),
("tenant_a", "Alpha"),
("tenant_b", "Gamma"),
("tenant_a", "Beta"),
("tenant_b", "Gamma"),
("tenant_a", "Alpha"),
("tenant_a", "Beta"),
];
for (tenant, entity) in events {
agg.record(&event(tenant, "2026-05", entity));
}
let a = agg.query("tenant_a", "2026-05");
assert_eq!(a.mutations.get("Alpha"), Some(&4));
assert_eq!(a.mutations.get("Beta"), Some(&3));
assert_eq!(a.mutations.get("Gamma"), None);
let b = agg.query("tenant_b", "2026-05");
assert_eq!(b.mutations.get("Gamma"), Some(&3));
assert_eq!(b.mutations.len(), 1);
}
#[test]
fn test_empty_result_for_unknown_tenant() {
let agg = UsageAggregator::new();
let summary = agg.query("nobody", "2026-05");
assert!(summary.mutations.is_empty());
}
#[test]
fn test_empty_result_for_unknown_period() {
let agg = UsageAggregator::new();
agg.record(&event("tenant_a", "2026-05", "User"));
let summary = agg.query("tenant_a", "2026-06");
assert!(summary.mutations.is_empty());
}
#[test]
fn test_validate_period_valid() {
assert!(validate_period("2026-04"));
assert!(validate_period("2026-01"));
assert!(validate_period("2026-12"));
assert!(validate_period("1000-06"));
assert!(validate_period("9999-11"));
}
#[test]
fn test_validate_period_invalid_month() {
assert!(!validate_period("2026-00")); assert!(!validate_period("2026-13")); assert!(!validate_period("2026-99"));
}
#[test]
fn test_validate_period_invalid_format() {
assert!(!validate_period("2026")); assert!(!validate_period("26-04")); assert!(!validate_period("2026/04")); assert!(!validate_period("2026-4")); assert!(!validate_period("2026-04-01")); assert!(!validate_period("")); }
#[test]
fn test_counters_reset_on_new_aggregator_without_persistence() {
let agg = UsageAggregator::new();
agg.record(&event("tenant_a", "2026-05", "User"));
assert_eq!(agg.query("tenant_a", "2026-05").mutations["User"], 1);
let new_agg = UsageAggregator::new();
assert_eq!(new_agg.query("tenant_a", "2026-05").mutations.get("User"), None);
}
struct InMemoryPersistenceBackend {
store: std::sync::Mutex<HashMap<(String, String, String), u64>>,
}
impl InMemoryPersistenceBackend {
fn new() -> Self {
Self {
store: std::sync::Mutex::new(HashMap::new()),
}
}
}
#[async_trait::async_trait]
impl UsageBackend for InMemoryPersistenceBackend {
async fn flush_deltas(
&self,
deltas: &HashMap<(String, String, String), u64>,
) -> Result<(), String> {
let mut store = self.store.lock().map_err(|e| e.to_string())?;
for (key, &delta) in deltas {
*store.entry(key.clone()).or_insert(0) += delta;
}
Ok(())
}
async fn load(&self) -> Result<HashMap<(String, String, String), u64>, String> {
let store = self.store.lock().map_err(|e| e.to_string())?;
Ok(store.clone())
}
}
struct UnreadableBackend {
flushes: std::sync::atomic::AtomicUsize,
}
#[async_trait::async_trait]
impl UsageBackend for UnreadableBackend {
async fn flush_deltas(
&self,
_deltas: &HashMap<(String, String, String), u64>,
) -> Result<(), String> {
self.flushes.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
Ok(())
}
async fn load(&self) -> Result<HashMap<(String, String, String), u64>, String> {
Err("connection reset by peer".to_string())
}
}
#[tokio::test]
async fn a_failed_startup_load_blocks_the_flush() {
let backend = std::sync::Arc::new(UnreadableBackend {
flushes: std::sync::atomic::AtomicUsize::new(0),
});
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.expect_err("the fixture backend cannot be read");
assert!(!agg.is_loaded(), "a failed load must leave the aggregator disarmed");
agg.record(&event("acme", "2026-07", "Order"));
let err = agg
.flush_to_backend()
.await
.expect_err("#861: a process that could not read the counters must not write them");
assert!(err.contains("startup load"), "the refusal must say why: {err}");
assert_eq!(
backend.flushes.load(std::sync::atomic::Ordering::Relaxed),
0,
"#861: the backend must not be written to at all"
);
}
#[tokio::test]
async fn concurrent_replicas_sum_rather_than_overwrite() {
let backend = std::sync::Arc::new(InMemoryPersistenceBackend::new());
let seed = UsageAggregator::new_with_backend(backend.clone());
seed.load_from_backend().await.expect("load");
for _ in 0..1000 {
seed.record(&event("acme", "2026-07", "Order"));
}
seed.flush_to_backend().await.expect("flush");
let replicas = [7_u32, 5, 3];
for own in replicas {
let r = UsageAggregator::new_with_backend(backend.clone());
r.load_from_backend().await.expect("load");
for _ in 0..own {
r.record(&event("acme", "2026-07", "Order"));
}
r.flush_to_backend().await.expect("flush");
}
let reader = UsageAggregator::new_with_backend(backend.clone());
reader.load_from_backend().await.expect("load");
assert_eq!(
reader.query("acme", "2026-07").mutations["Order"],
1000 + 7 + 5 + 3,
"#861: every replica's interval must be added; an absolute write kept only the last"
);
}
#[tokio::test]
async fn repeated_flushes_do_not_double_count() {
let backend = std::sync::Arc::new(InMemoryPersistenceBackend::new());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.expect("load");
agg.record(&event("t1", "2026-07", "User"));
agg.record(&event("t1", "2026-07", "User"));
agg.flush_to_backend().await.expect("flush");
agg.flush_to_backend().await.expect("flush again");
agg.flush_to_backend().await.expect("and again");
let reader = UsageAggregator::new_with_backend(backend.clone());
reader.load_from_backend().await.expect("load");
assert_eq!(
reader.query("t1", "2026-07").mutations["User"],
2,
"three flushes of two events must persist two, not six"
);
}
#[tokio::test]
async fn test_flush_and_load_round_trip() {
let backend = std::sync::Arc::new(InMemoryPersistenceBackend::new());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.expect("load before flush (#861)");
agg.record(&event("tenant_a", "2026-05", "User"));
agg.record(&event("tenant_a", "2026-05", "User"));
agg.record(&event("tenant_b", "2026-05", "Order"));
agg.flush_to_backend().await.expect("flush");
let new_agg = UsageAggregator::new_with_backend(backend.clone());
assert_eq!(new_agg.query("tenant_a", "2026-05").mutations.get("User"), None);
new_agg.load_from_backend().await.expect("load");
assert_eq!(new_agg.query("tenant_a", "2026-05").mutations["User"], 2);
assert_eq!(new_agg.query("tenant_b", "2026-05").mutations["Order"], 1);
}
#[tokio::test]
async fn test_load_merges_with_inflight_events() {
let backend = std::sync::Arc::new(InMemoryPersistenceBackend::new());
let agg = UsageAggregator::new_with_backend(backend.clone());
agg.load_from_backend().await.expect("load before flush (#861)");
agg.record(&event("t1", "2026-05", "User"));
agg.flush_to_backend().await.expect("flush");
let new_agg = UsageAggregator::new_with_backend(backend.clone());
new_agg.record(&event("t1", "2026-05", "User"));
new_agg.record(&event("t1", "2026-05", "User"));
new_agg.load_from_backend().await.expect("load");
assert_eq!(new_agg.query("t1", "2026-05").mutations["User"], 3);
}
#[tokio::test]
async fn test_noop_backend_flush_and_load_are_harmless() {
let agg = UsageAggregator::new(); agg.record(&event("t1", "2026-05", "User"));
agg.flush_to_backend().await.expect("flush ok");
let new_agg = UsageAggregator::new();
new_agg.load_from_backend().await.expect("load ok");
assert_eq!(new_agg.query("t1", "2026-05").mutations.get("User"), None); }
}
mod layer_tests {
use std::sync::Arc;
use chrono::Utc;
use tracing_subscriber::{Registry, layer::SubscriberExt as _};
use super::super::layer::*;
use crate::usage::aggregator::UsageAggregator;
fn current_period() -> String {
Utc::now().format("%Y-%m").to_string()
}
#[test]
fn test_layer_captures_mutation_audit_event() {
let aggregator = Arc::new(UsageAggregator::new());
let layer = MutationAuditLayer::new(Arc::clone(&aggregator));
let subscriber = Registry::default().with(layer);
let _guard = tracing::subscriber::set_default(subscriber);
tracing::info!(
target: "fraiseql::mutation_audit",
mutation_name = "create_user",
entity_type = %"User",
operation = %"create",
tenant_id = %"acme",
"mutation.executed"
);
let period = current_period();
let summary = aggregator.query("acme", &period);
assert_eq!(summary.mutations.get("User"), Some(&1));
}
#[test]
fn test_layer_ignores_other_targets() {
let aggregator = Arc::new(UsageAggregator::new());
let layer = MutationAuditLayer::new(Arc::clone(&aggregator));
let subscriber = Registry::default().with(layer);
let _guard = tracing::subscriber::set_default(subscriber);
tracing::info!(
target: "fraiseql::other",
mutation_name = "create_user",
entity_type = %"User",
operation = %"create",
tenant_id = %"acme",
"not an audit event"
);
let summary = aggregator.query("acme", ¤t_period());
assert!(summary.mutations.is_empty());
}
#[test]
fn test_layer_aggregates_multiple_events_across_tenants() {
let aggregator = Arc::new(UsageAggregator::new());
let layer = MutationAuditLayer::new(Arc::clone(&aggregator));
let subscriber = Registry::default().with(layer);
let _guard = tracing::subscriber::set_default(subscriber);
let period = current_period();
for _ in 0..3 {
tracing::info!(
target: "fraiseql::mutation_audit",
mutation_name = "create_user",
entity_type = %"User",
operation = %"create",
tenant_id = %"tenant_x",
"mutation.executed"
);
}
for _ in 0..2 {
tracing::info!(
target: "fraiseql::mutation_audit",
mutation_name = "delete_order",
entity_type = %"Order",
operation = %"delete",
tenant_id = %"tenant_y",
"mutation.executed"
);
}
let x = aggregator.query("tenant_x", &period);
assert_eq!(x.mutations.get("User"), Some(&3));
assert_eq!(x.mutations.get("Order"), None);
let y = aggregator.query("tenant_y", &period);
assert_eq!(y.mutations.get("Order"), Some(&2));
assert_eq!(y.mutations.get("User"), None);
}
#[test]
fn test_layer_handles_empty_tenant_id() {
let aggregator = Arc::new(UsageAggregator::new());
let layer = MutationAuditLayer::new(Arc::clone(&aggregator));
let subscriber = Registry::default().with(layer);
let _guard = tracing::subscriber::set_default(subscriber);
tracing::info!(
target: "fraiseql::mutation_audit",
mutation_name = "update_product",
entity_type = %"Product",
operation = %"update",
tenant_id = %"",
"mutation.executed"
);
let summary = aggregator.query("", ¤t_period());
assert_eq!(summary.mutations.get("Product"), Some(&1));
}
#[test]
fn test_aggregator_accessor() {
let aggregator = Arc::new(UsageAggregator::new());
let layer = MutationAuditLayer::new(Arc::clone(&aggregator));
assert!(Arc::ptr_eq(&aggregator, layer.aggregator()));
}
}