#![cfg(feature = "testkit")]
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use meterstore::MeterStore;
use meterstore::config::TableConfig;
use meterstore::testkit::{MeteringWorkload, TestHarness};
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
const DAYS: i64 = 6;
const METERS: usize = 8;
const EXPECTED_ROWS: u64 = (DAYS as u64) * 96 * (METERS as u64);
async fn seeded() -> (TestHarness, Arc<MeterStore>) {
let harness = TestHarness::with_config(
TableConfig::new(TestHarness::TABLE)
.settlement_lag(Duration::DAY)
.reader_grace(Duration::days(365))
.build()
.expect("config"),
)
.await
.expect("harness");
let workload = MeteringWorkload::new(START)
.seed(0x5EED)
.malo_ids(METERS)
.days(DAYS);
let (from, to) = workload.range();
harness
.ensure_partitions(from, to + Duration::days(2))
.await
.expect("partitions");
harness.seed_watermark(from).await.expect("watermark");
let store = harness.store().await.expect("store");
let series = workload.generate().expect("workload");
harness.ingest(&store, &series).await.expect("ingest");
(harness, Arc::new(store))
}
async fn resolved_count(store: &MeterStore) -> u64 {
let sql = format!(r#"SELECT count(*) AS n FROM "{}""#, store.resolved_table());
let result = store.query(&sql).await.expect("count");
use datafusion::arrow::array::AsArray;
result.batches()[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Int64Type>()
.value(0) as u64
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_count_is_exact_while_archival_moves_the_data_under_it() {
let (_h, store) = seeded().await;
assert_eq!(resolved_count(&store).await, EXPECTED_ROWS, "seeded");
let stop = Arc::new(AtomicBool::new(false));
let reads = Arc::new(AtomicU64::new(0));
let readers: Vec<_> = (0..3)
.map(|_| {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
let reads = Arc::clone(&reads);
tokio::spawn(async move {
while !stop.load(Ordering::Relaxed) {
let n = resolved_count(&store).await;
reads.fetch_add(1, Ordering::Relaxed);
assert_eq!(
n, EXPECTED_ROWS,
"a count over an unchanging range changed while archival ran: \
got {n}, expected {EXPECTED_ROWS} — a window is in neither tier"
);
tokio::task::yield_now().await;
}
})
})
.collect();
let archiver = {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
tokio::spawn(async move {
for day in 1..=DAYS {
let now = START + Duration::days(day + 1);
store.archive(now, 4).await.expect("archive");
tokio::task::yield_now().await;
}
stop.store(true, Ordering::Relaxed);
store.watermark().await.expect("watermark")
})
};
let watermark = archiver.await.expect("archiver task");
for reader in readers {
reader.await.expect("reader task");
}
assert!(
reads.load(Ordering::Relaxed) >= 3,
"the readers must actually have run"
);
assert!(
watermark.get() > START,
"archival must actually have moved the boundary: {watermark}"
);
assert_eq!(
resolved_count(&store).await,
EXPECTED_ROWS,
"and the count is still exact once everything has settled"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn ingest_archival_and_reads_together_leave_the_written_range_exact() {
let (_h, store) = seeded().await;
let counted = format!(
r#"SELECT count(*) AS n FROM "{}" WHERE "from" >= '{}' AND "from" < '{}'"#,
store.resolved_table(),
START
.format(&time::format_description::well_known::Rfc3339)
.unwrap(),
(START + Duration::days(DAYS))
.format(&time::format_description::well_known::Rfc3339)
.unwrap(),
);
let stop = Arc::new(AtomicBool::new(false));
let reads = Arc::new(AtomicU64::new(0));
let readers: Vec<_> = (0..2)
.map(|_| {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
let reads = Arc::clone(&reads);
let sql = counted.clone();
tokio::spawn(async move {
while !stop.load(Ordering::Relaxed) {
let result = store.query(&sql).await.expect("count");
use datafusion::arrow::array::AsArray;
let n = result.batches()[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Int64Type>()
.value(0) as u64;
reads.fetch_add(1, Ordering::Relaxed);
assert_eq!(
n, EXPECTED_ROWS,
"the range written up front must stay exact while ingest, \
archival and reads run together"
);
tokio::task::yield_now().await;
}
})
})
.collect();
let writer = {
let store = Arc::clone(&store);
tokio::spawn(async move {
for day in 0..3i64 {
let ahead = MeteringWorkload::new(START + Duration::days(DAYS + day))
.seed(0xA11CE)
.malo_ids(METERS)
.days(1)
.generate()
.expect("workload");
for delivery in &ahead {
store
.append(std::slice::from_ref(delivery))
.await
.expect("append");
}
tokio::task::yield_now().await;
}
})
};
let archiver = {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
tokio::spawn(async move {
for day in 1..=DAYS {
store
.archive(START + Duration::days(day + 1), 4)
.await
.expect("archive");
tokio::task::yield_now().await;
}
stop.store(true, Ordering::Relaxed);
})
};
writer.await.expect("writer task");
archiver.await.expect("archiver task");
for reader in readers {
reader.await.expect("reader task");
}
assert!(
reads.load(Ordering::Relaxed) >= 2,
"the readers must have run"
);
assert_eq!(
resolved_count(&store).await,
EXPECTED_ROWS + 3 * 96 * (METERS as u64),
"and everything written is present once it has all settled"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn corrections_landing_during_a_scan_are_never_double_counted() {
let (_h, store) = seeded().await;
store
.archive(START + Duration::days(DAYS + 2), 16)
.await
.expect("archive");
assert!(
store.watermark().await.unwrap().get() >= START + Duration::days(DAYS),
"the range must be entirely cold for this to exercise elision"
);
assert_eq!(resolved_count(&store).await, EXPECTED_ROWS, "archived");
let stop = Arc::new(AtomicBool::new(false));
let corrections = Arc::new(AtomicU64::new(0));
let readers: Vec<_> = (0..3)
.map(|_| {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
tokio::spawn(async move {
while !stop.load(Ordering::Relaxed) {
let n = resolved_count(&store).await;
assert_eq!(
n, EXPECTED_ROWS,
"a correction is a new version of an existing reading, so the \
resolved count cannot change: got {n}, expected {EXPECTED_ROWS}"
);
tokio::task::yield_now().await;
}
})
})
.collect();
let corrector = {
let store = Arc::clone(&store);
let stop = Arc::clone(&stop);
let corrections = Arc::clone(&corrections);
tokio::spawn(async move {
for round in 1..=4u64 {
let restated = MeteringWorkload::new(START)
.seed(0x5EED)
.malo_ids(METERS)
.days(1)
.version(20_260_720_000_000u128 + u128::from(round))
.generate()
.expect("restatement");
for delivery in &restated {
store
.append(std::slice::from_ref(delivery))
.await
.expect("correction");
}
corrections.fetch_add(1, Ordering::Relaxed);
tokio::task::yield_now().await;
}
stop.store(true, Ordering::Relaxed);
})
};
corrector.await.expect("corrector task");
for reader in readers {
reader.await.expect("reader task");
}
assert_eq!(corrections.load(Ordering::Relaxed), 4);
assert_eq!(
resolved_count(&store).await,
EXPECTED_ROWS,
"the resolved view still holds one row per reading"
);
let raw = format!(r#"SELECT count(*) AS n FROM "{}""#, store.raw_table());
let result = store.query(&raw).await.expect("raw count");
use datafusion::arrow::array::AsArray;
let stored = result.batches()[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Int64Type>()
.value(0) as u64;
assert!(
stored > EXPECTED_ROWS,
"the audit trail must hold every version: {stored} rows for {EXPECTED_ROWS} readings"
);
}
async fn second_replica(harness: &TestHarness) -> Arc<meterstore::PostgresHot> {
let pool = sqlx::PgPool::connect(harness.url())
.await
.expect("second replica pool");
Arc::new(meterstore::PostgresHot::new(pool))
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_lease_one_replica_holds_is_refused_to_another_and_released_to_it() {
use meterstore::tiering::HotStore;
let (harness, _store) = seeded().await;
let table = harness.config().name();
let a = harness.hot().clone();
let b = second_replica(&harness).await;
let held = a
.try_archive_lease(table)
.await
.expect("lease query")
.expect("an uncontended lease is granted");
assert!(
b.try_archive_lease(table)
.await
.expect("lease query")
.is_none(),
"a second replica took a lease the first holds"
);
assert!(
b.try_archive_lease("some_other_table")
.await
.expect("lease query")
.is_some(),
"the lease is per table, not per deployment"
);
held.release().await.expect("release");
let after = b
.try_archive_lease(table)
.await
.expect("lease query")
.expect("the lease is available once released");
after.release().await.expect("release");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn two_replicas_racing_to_archive_leave_the_range_exact() {
use meterstore::{Archiver, tiering::ColdStore};
let (harness, store) = seeded().await;
let before = resolved_count(&store).await;
assert_eq!(before, EXPECTED_ROWS, "the fixture is what it claims");
let table = harness.config().name();
let a = Archiver::new(
harness.hot().clone(),
harness.cold().clone(),
harness.config().clone(),
);
let b = Archiver::new(
second_replica(&harness).await,
harness.cold().clone(),
harness.config().clone(),
);
let mut watermark = harness.cold().watermark(table).await.expect("watermark");
let mut contended = 0usize;
let mut archived = 0usize;
for day in 1..=DAYS + 1 {
let now = START + Duration::days(day + 1);
let (ra, rb) = tokio::join!(a.run_once(now), b.run_once(now));
let (ra, rb) = (ra.expect("replica a"), rb.expect("replica b"));
for outcome in [&ra, &rb] {
if outcome.lease_contended {
contended += 1;
}
if outcome.archived_anything() {
archived += 1;
}
}
assert!(
!(ra.archived_anything() && rb.archived_anything()),
"both replicas archived in one cycle: a={ra:?} b={rb:?}"
);
let now_watermark = harness.cold().watermark(table).await.expect("watermark");
assert!(
now_watermark >= watermark,
"the watermark moved backwards: {watermark} -> {now_watermark}"
);
watermark = now_watermark;
a.verify_invariant()
.await
.expect("invariant after the race");
assert_eq!(
resolved_count(&store).await,
before,
"archival moved rows between tiers; it must not change how many there are"
);
}
assert!(archived > 0, "neither replica ever archived anything");
println!("two-replica race: {archived} archived, {contended} contended");
}