#![cfg(feature = "testkit")]
use datafusion::prelude::SessionContext;
use meterstore::testkit::{MeteringWorkload, Oracle, TestHarness};
use rust_decimal::Decimal;
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
async fn archived(workload: MeteringWorkload) -> (TestHarness, meterstore::MeterStore, Oracle) {
let harness = TestHarness::start().await.expect("harness");
let (from, to) = workload.range();
harness
.ensure_partitions(from, to + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(from).await.expect("watermark");
let store = harness.store().await.expect("store");
let series = workload.generate().expect("workload");
let mut oracle = Oracle::new();
oracle.record(&series).expect("oracle");
harness.ingest(&store, &series).await.expect("ingest");
store
.archive(to + Duration::days(2), 64)
.await
.expect("archive");
assert!(
store.watermark().await.unwrap().get() >= to,
"the suite is meaningless unless everything reached the cold tier"
);
(harness, store, oracle)
}
async fn external_reader(harness: &TestHarness) -> SessionContext {
use datafusion::datasource::file_format::parquet::ParquetFormat;
use datafusion::datasource::listing::{
ListingOptions, ListingTable, ListingTableConfig, ListingTableUrl,
};
use std::sync::Arc;
let files = harness.parquet_files();
assert!(
!files.is_empty(),
"archival wrote no Parquet files, so there is nothing to read"
);
let ctx = SessionContext::new();
let options =
ListingOptions::new(Arc::new(ParquetFormat::default())).with_file_extension(".parquet");
let urls: Vec<ListingTableUrl> = files
.iter()
.map(|f| ListingTableUrl::parse(format!("file://{}", f.display())).expect("file url"))
.collect();
let resolved = options
.infer_schema(&ctx.state(), &urls[0])
.await
.expect("infer schema from the files alone");
let config = ListingTableConfig::new_with_multi_paths(urls)
.with_listing_options(options)
.with_schema(resolved);
ctx.register_table(
TestHarness::TABLE,
Arc::new(ListingTable::try_new(config).expect("listing table")),
)
.expect("register");
ctx
}
async fn count(ctx: &SessionContext, sql: &str) -> i64 {
use datafusion::arrow::array::AsArray;
let batches = ctx
.sql(sql)
.await
.expect("plan")
.collect()
.await
.expect("run");
batches[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Int64Type>()
.value(0)
}
async fn total(ctx: &SessionContext, sql: &str) -> Decimal {
use datafusion::arrow::array::AsArray;
let batches = ctx
.sql(sql)
.await
.expect("plan")
.collect()
.await
.expect("run");
let array = batches[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Decimal128Type>();
Decimal::try_from_i128_with_scale(array.value(0), u32::from(array.scale() as u8))
.expect("decimal in range")
.normalize()
}
#[tokio::test]
async fn the_files_are_readable_without_meterstore_or_iceberg() {
let workload = MeteringWorkload::new(START)
.seed(0x09E4)
.malo_ids(4)
.days(3);
let (from, to) = workload.range();
let (harness, _store, oracle) = archived(workload).await;
let ctx = external_reader(&harness).await;
let rows = count(&ctx, "SELECT COUNT(*) FROM readings_versions").await as u64;
assert_eq!(
rows,
oracle.row_count(from, to),
"a plain Parquet reader must see every row"
);
}
#[tokio::test]
async fn the_published_resolution_sql_gives_an_external_engine_the_right_answer() {
let workload = MeteringWorkload::new(START)
.seed(0xC0FFEE)
.malo_ids(4)
.days(3)
.with_corrections(0.2);
let (from, to) = workload.range();
let (harness, store, oracle) = archived(workload).await;
let ctx = external_reader(&harness).await;
let resolution = store.resolution_sql();
let resolved = total(
&ctx,
&format!("SELECT COALESCE(SUM(value), 0) FROM ({resolution}) AS resolved"),
)
.await;
assert_eq!(
resolved,
oracle.sum_kwh(from, to).normalize(),
"the published SQL must reproduce MeterStore's own answer"
);
}
#[tokio::test]
async fn the_naive_query_really_does_double_count() {
let workload = MeteringWorkload::new(START)
.seed(0xBAD)
.malo_ids(3)
.days(2)
.with_corrections(0.3);
let (from, to) = workload.range();
let (harness, _store, oracle) = archived(workload).await;
let ctx = external_reader(&harness).await;
let naive = total(
&ctx,
"SELECT COALESCE(SUM(value), 0) FROM readings_versions",
)
.await;
assert!(
naive > oracle.sum_kwh(from, to).normalize(),
"a raw sum must overstate; if it does not, the fixture produced no corrections"
);
}
#[tokio::test]
async fn stored_values_are_self_describing() {
let workload = MeteringWorkload::new(START)
.seed(0x5E1F)
.malo_ids(2)
.days(1);
let (harness, _store, _oracle) = archived(workload).await;
let ctx = external_reader(&harness).await;
let measured = count(
&ctx,
"SELECT COUNT(*) FROM readings_versions WHERE quality = 'MEASURED'",
)
.await;
assert!(measured > 0, "quality must be readable as its own name");
let resolution = count(
&ctx,
"SELECT COUNT(*) FROM readings_versions WHERE resolution = 'PT15M'",
)
.await;
assert!(
resolution > 0,
"resolution must be an ISO 8601 duration, not a local code"
);
let obis = count(
&ctx,
"SELECT COUNT(DISTINCT obis_code) FROM readings_versions WHERE obis_code LIKE '1-0:%'",
)
.await;
assert_eq!(obis, 1);
}
#[tokio::test]
async fn decimals_survive_at_full_precision_outside_this_crate() {
let workload = MeteringWorkload::new(START).seed(0xDEC).malo_ids(3).days(2);
let (from, to) = workload.range();
let (harness, _store, oracle) = archived(workload).await;
let ctx = external_reader(&harness).await;
let sum = total(
&ctx,
"SELECT COALESCE(SUM(value), 0) FROM readings_versions",
)
.await;
assert_eq!(sum, oracle.sum_kwh(from, to).normalize());
let batches = ctx
.sql("SELECT value FROM readings_versions LIMIT 1")
.await
.unwrap()
.collect()
.await
.unwrap();
assert!(matches!(
batches[0].schema().field(0).data_type(),
datafusion::arrow::datatypes::DataType::Decimal128(18, 6)
));
}
#[tokio::test]
async fn the_footer_carries_the_tuning_the_design_claims() {
use parquet::file::reader::{FileReader, SerializedFileReader};
use parquet::schema::types::ColumnPath;
let workload = MeteringWorkload::new(START)
.seed(0xF007)
.malo_ids(4)
.days(2);
let (harness, _store, _oracle) = archived(workload).await;
let files = harness.parquet_files();
let file = std::fs::File::open(&files[0]).expect("open a data file");
let reader = SerializedFileReader::new(file).expect("parse the footer");
let metadata = reader.metadata();
let row_group = metadata.row_group(0);
let sorting = row_group
.sorting_columns()
.expect("the footer must declare a sort order");
assert_eq!(sorting.len(), 2, "malo_id then from");
assert!(sorting.iter().all(|c| !c.descending));
for name in ["malo_id", "obis_code"] {
let column = row_group
.columns()
.iter()
.find(|c| c.column_path() == &ColumnPath::from(name))
.unwrap_or_else(|| panic!("{name} must be a column"));
assert!(
column.bloom_filter_offset().is_some(),
"{name} must carry a bloom filter"
);
}
let from_column = row_group
.columns()
.iter()
.find(|c| c.column_path() == &ColumnPath::from("from"))
.expect("from column");
assert!(
from_column.offset_index_offset().is_some(),
"page index must be written, or layer 4 of the pruning stack is absent"
);
}
#[tokio::test]
async fn the_files_really_are_sorted_the_way_the_footer_says() {
let workload = MeteringWorkload::new(START)
.seed(0x5017)
.malo_ids(5)
.days(2);
let (harness, _store, _oracle) = archived(workload).await;
for path in harness.parquet_files() {
let ctx = SessionContext::new();
ctx.register_parquet(
"one_file",
path.to_str().expect("utf-8 path"),
datafusion::prelude::ParquetReadOptions::default(),
)
.await
.expect("register a single file");
let batches = ctx
.sql(r#"SELECT malo_id, "from" FROM one_file"#)
.await
.expect("plan")
.collect()
.await
.expect("run");
let mut previous: Option<(String, i64)> = None;
for batch in &batches {
let malo = datafusion::arrow::compute::cast(
batch.column(0),
&datafusion::arrow::datatypes::DataType::Utf8,
)
.expect("malo_id is a string column");
use datafusion::arrow::array::AsArray;
let malo = malo.as_string::<i32>();
let from = batch
.column(1)
.as_primitive::<datafusion::arrow::datatypes::TimestampMicrosecondType>();
for i in 0..batch.num_rows() {
let current = (malo.value(i).to_string(), from.value(i));
if let Some(prev) = &previous {
assert!(
prev <= ¤t,
"{} breaks its declared (malo_id, from) order: {prev:?} then {current:?}",
path.display()
);
}
previous = Some(current);
}
}
}
}