#![cfg(feature = "testkit")]
use metering::interval::{MeasurementUnit, Sparte};
use meterstore::encode::StoredSeries;
use meterstore::testkit::{MeteringWorkload, Oracle, TestHarness};
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
async fn split_store(workload: &MeteringWorkload) -> (TestHarness, meterstore::MeterStore, Oracle) {
let (from, to) = workload.range();
let harness = TestHarness::start().await.expect("harness");
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(from + Duration::days(2), 1)
.await
.expect("archive");
(harness, store, oracle)
}
async fn scalar(store: &meterstore::MeterStore, sql: &str) -> String {
let result = store.query(sql).await.expect("query");
let batch = result
.batches()
.iter()
.find(|b| b.num_rows() > 0)
.expect("a row");
meterstore::arrow::util::display::array_value_to_string(batch.column(0), 0).expect("render")
}
#[tokio::test]
async fn water_keeps_its_cubic_metres_across_both_tiers() {
let workload = MeteringWorkload::new(START)
.seed(0xA7E_0000)
.sparte(Sparte::Wasser)
.malo_ids(3)
.days(3);
let (_h, store, oracle) = split_store(&workload).await;
let (from, to) = workload.range();
assert_eq!(
scalar(&store, "SELECT COUNT(*) FROM readings").await,
oracle.row_count(from, to).to_string(),
"every water reading must survive both tiers"
);
let units = store
.query("SELECT DISTINCT sparte, unit FROM readings")
.await
.expect("query");
let rendered = meterstore::arrow::util::pretty::pretty_format_batches(units.batches())
.expect("render")
.to_string();
assert!(rendered.contains("WASSER"), "{rendered}");
assert!(rendered.contains("M3"), "{rendered}");
assert!(
!rendered.contains("KWH"),
"water has no calorific value, so no row may claim kWh: {rendered}"
);
}
#[tokio::test]
async fn gas_may_be_stored_on_either_side_of_the_brennwert_conversion() {
let converted = MeteringWorkload::new(START)
.seed(0x6A5_0001)
.sparte(Sparte::Gas)
.malo_ids(2)
.days(2);
assert_eq!(Sparte::Gas.billing_unit(), MeasurementUnit::KiloWattHour);
let (_h, store, _oracle) = split_store(&converted).await;
assert_eq!(
scalar(&store, "SELECT DISTINCT unit FROM readings").await,
"KWH",
"the default for gas is the settled unit"
);
let raw = MeteringWorkload::new(START)
.seed(0x6A5_0002)
.sparte(Sparte::Gas)
.in_unit(MeasurementUnit::CubicMetre)
.malo_ids(2)
.days(2);
let (_h2, store2, _o2) = split_store(&raw).await;
assert_eq!(
scalar(&store2, "SELECT DISTINCT unit FROM readings").await,
"M3",
"unconverted gas is m³ and must say so"
);
}
#[test]
fn a_unit_the_commodity_cannot_have_is_refused_at_the_write() {
let workload = MeteringWorkload::new(START)
.sparte(Sparte::Wasser)
.in_unit(MeasurementUnit::KiloWattHour)
.malo_ids(1)
.days(1);
let series = workload.generate().expect("generating is not what fails");
let err =
meterstore::encode::to_record_batch(&series).expect_err("water cannot be expressed in kWh");
let msg = err.to_string();
assert!(msg.contains("WASSER"), "{msg}");
assert!(msg.contains("KWH"), "{msg}");
assert!(
msg.contains("M3"),
"the message must name the unit that would have been right: {msg}"
);
}
#[test]
fn every_sparte_admits_its_own_measured_and_billing_units() {
for sparte in Sparte::ALL {
for unit in MeasurementUnit::ALL {
let admissible = unit == sparte.measured_unit() || unit == sparte.billing_unit();
let series = MeteringWorkload::new(START)
.sparte(sparte)
.in_unit(unit)
.malo_ids(1)
.days(1)
.generate()
.expect("workload");
assert_eq!(
meterstore::encode::to_record_batch(&series).is_ok(),
admissible,
"{sparte} in {unit} should be {}",
if admissible { "accepted" } else { "refused" }
);
}
}
}
#[tokio::test]
async fn a_mixed_portfolio_sums_per_unit_and_never_across_them() {
let harness = TestHarness::start().await.expect("harness");
let to = START + Duration::days(2);
harness
.ensure_partitions(START, to + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(START).await.expect("watermark");
let store = harness.store().await.expect("store");
let mut all: Vec<StoredSeries> = Vec::new();
for (i, sparte) in [Sparte::Strom, Sparte::Gas, Sparte::Wasser]
.into_iter()
.enumerate()
{
all.extend(
MeteringWorkload::new(START)
.seed(0x1_0000 + i as u64)
.malo_offset(i * 100)
.sparte(sparte)
.malo_ids(2)
.days(2)
.generate()
.expect("workload"),
);
}
harness.ingest(&store, &all).await.expect("ingest");
let batches = store
.query("SELECT unit, COUNT(*) FROM readings GROUP BY unit ORDER BY unit")
.await
.expect("query");
let rendered = meterstore::arrow::util::pretty::pretty_format_batches(batches.batches())
.expect("render")
.to_string();
assert!(rendered.contains("KWH"), "{rendered}");
assert!(rendered.contains("M3"), "{rendered}");
}
fn gas_quarter(from: OffsetDateTime, value: i64) -> metering::interval::MeterInterval {
metering::interval::MeterInterval {
from,
to: from + Duration::minutes(15),
value: rust_decimal::Decimal::new(value, 0),
quality: metering::QualityFlag::Measured,
obis_code: "7-1:99.33.0".parse().ok(),
}
}
fn gas_series(malo: &str, from: OffsetDateTime, to: OffsetDateTime, value: i64) -> StoredSeries {
let intervals: Vec<_> = std::iter::successors(Some(from), |t| Some(*t + Duration::minutes(15)))
.take_while(|t| *t < to)
.map(|t| gas_quarter(t, value))
.collect();
let mut series = metering::measurement_series::MeasurementSeries::new(
malo.parse().expect("a valid MaLo-ID"),
"7-1:99.33.0".parse().ok(),
intervals,
metering::measurement_series::MeasurementSource::Mscons {
pid: 13_005,
message_ref: None,
sender_mp_id: "9900000000001".parse().expect("a valid Marktpartner-ID"),
},
to,
);
series.resolution = Some(metering::IntervalResolution::QuarterHour);
StoredSeries::of(
Sparte::Gas,
series,
meterstore::ScopedVersion::new(
meterstore::VersionScope::for_interval("9900000000001", from, Sparte::Gas)
.expect("scope"),
meterstore::Version::new(20_261_001_000_001).expect("version"),
),
to,
)
}
#[tokio::test]
async fn a_gas_delivery_scoped_to_its_real_bilanzierungsmonat_is_accepted() {
let straddling = datetime!(2026-03-01 1:00 UTC);
let series = gas_series(
"10000000009",
straddling,
straddling + Duration::hours(1),
7,
);
assert_eq!(
series.version.scope().period(),
"2026-02",
"the derived scope must be February's, not March's"
);
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(
straddling - Duration::days(1),
straddling + Duration::days(1),
)
.await
.expect("partitions");
let store = harness.store().await.expect("store");
let outcome = store
.append(&[series])
.await
.expect("a correctly scoped gas delivery");
assert_eq!(outcome.total(), 4, "four quarter-hours");
let mut wrong = gas_series(
"10000000009",
straddling,
straddling + Duration::hours(1),
7,
);
wrong.version = meterstore::ScopedVersion::new(
meterstore::VersionScope::new("9900000000001", 2026, 3).expect("scope"),
meterstore::Version::new(20_261_001_000_002).expect("version"),
);
let err = store
.append(&[wrong])
.await
.expect_err("the calendar month is not this interval's gas Bilanzierungsmonat")
.to_string();
assert!(err.contains("Bilanzierungsmonat"), "{err}");
assert!(err.contains("GAS"), "{err}");
}
#[tokio::test]
async fn a_gas_lastgang_groups_onto_the_gastag_not_the_calendar_day() {
let first = datetime!(2026-07-15 4:00 UTC); let last = first + Duration::days(2);
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(first, last + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(first).await.expect("watermark");
let store = harness.store().await.expect("store");
harness
.ingest(&store, &[gas_series("10000000009", first, last, 4)])
.await
.expect("ingest");
let rendered = |batches: &[meterstore::arrow::array::RecordBatch]| {
meterstore::arrow::util::pretty::pretty_format_batches(batches)
.expect("render")
.to_string()
};
let gas = store
.query(
r#"SELECT meter_gas_day("from") AS d, COUNT(*) AS n
FROM readings GROUP BY 1 ORDER BY 1"#,
)
.await
.expect("query");
let gas = rendered(gas.batches());
assert!(gas.contains("2026-07-15"), "{gas}");
assert!(gas.contains("2026-07-16"), "{gas}");
assert!(
!gas.contains("2026-07-17"),
"two whole Gastage must be two buckets: {gas}"
);
let balanced = store
.query(
r#"SELECT meter_balancing_day("from", sparte) AS d, COUNT(*) AS n
FROM readings GROUP BY 1 ORDER BY 1"#,
)
.await
.expect("query");
assert_eq!(
rendered(balanced.batches()),
gas,
"dispatch must not differ"
);
let calendar = store
.query(
r#"SELECT meter_local_day("from") AS d, COUNT(*) AS n
FROM readings GROUP BY 1 ORDER BY 1"#,
)
.await
.expect("query");
let calendar = rendered(calendar.batches());
assert!(
calendar.contains("2026-07-17"),
"the calendar day spills into a third bucket, which is the bug: {calendar}"
);
}
#[tokio::test]
async fn gas_completeness_puts_the_long_day_on_the_saturday() {
let first = datetime!(2026-10-24 4:00 UTC); let last = datetime!(2026-10-26 5:00 UTC);
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(first, last + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(first).await.expect("watermark");
let store = harness.store().await.expect("store");
harness
.ingest(&store, &[gas_series("10000000009", first, last, 4)])
.await
.expect("ingest");
let report = store.completeness(first, last).await.expect("completeness");
assert_eq!(report.len(), 1, "one channel: {report:?}");
let row = &report[0];
assert_eq!(row.sparte, Sparte::Gas);
assert_eq!(row.actual, 196);
assert_eq!(row.expected, 196);
assert!(
row.is_complete(),
"a complete pair of Gastage must report complete: {row:?}"
);
assert_eq!(row.missing, 0);
assert_eq!(row.surplus, 0, "the Saturday is 25 hours long, not 24");
assert_eq!(row.first_gap, None);
}
#[tokio::test]
async fn a_gap_in_a_gas_day_is_still_found() {
let first = datetime!(2026-10-24 4:00 UTC);
let last = datetime!(2026-10-26 5:00 UTC);
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(first, last + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(first).await.expect("watermark");
let store = harness.store().await.expect("store");
let cut = datetime!(2026-10-25 4:00 UTC);
harness
.ingest(
&store,
&[
gas_series("10000000009", first, cut, 4),
gas_series("10000000009", datetime!(2026-10-25 5:00 UTC), last, 4),
],
)
.await
.expect("ingest");
let report = store.completeness(first, last).await.expect("completeness");
let row = &report[0];
assert_eq!(row.missing, 4, "the missing hour must be reported: {row:?}");
assert_eq!(row.surplus, 0);
assert_eq!(
row.first_gap,
Some(time::macros::date!(2026 - 10 - 24)),
"and attributed to the Gastag that began on the Saturday"
);
}