use std::collections::BTreeMap;
use metering::ids::MaloId;
use metering::interval::{MeasurementUnit, MeterInterval, QualityFlag, Sparte};
use metering::measurement_series::{MeasurementSeries, MeasurementSource};
use metering::resolution::IntervalResolution;
use rust_decimal::Decimal;
use time::{Duration, OffsetDateTime};
use crate::encode::StoredSeries;
use crate::error::{Error, Result};
use crate::version::{ScopedVersion, Version, VersionScope};
pub mod harness;
pub mod postgres;
pub use harness::TestHarness;
#[derive(Debug, Clone)]
pub struct Rng(u64);
impl Rng {
pub const fn new(seed: u64) -> Self {
Self(seed)
}
pub fn next_u64(&mut self) -> u64 {
self.0 = self.0.wrapping_add(0x9E37_79B9_7F4A_7C15);
let mut z = self.0;
z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
z ^ (z >> 31)
}
pub fn below(&mut self, n: u64) -> u64 {
if n == 0 { 0 } else { self.next_u64() % n }
}
pub fn chance(&mut self, p: f64) -> bool {
let scale = 1_000_000u64;
self.below(scale) < (p.clamp(0.0, 1.0) * scale as f64) as u64
}
}
#[derive(Debug, Clone)]
pub struct MeteringWorkload {
seed: u64,
malo_ids: usize,
days: i64,
start: OffsetDateTime,
resolution: Duration,
correction_rate: f64,
gap_rate: f64,
operator: String,
malo_offset: usize,
sparte: Sparte,
unit: MeasurementUnit,
}
impl MeteringWorkload {
pub fn new(start: OffsetDateTime) -> Self {
Self {
seed: 0x5EED,
malo_ids: 10,
days: 3,
start,
resolution: Duration::minutes(15),
correction_rate: 0.0,
gap_rate: 0.0,
operator: "9900000000001".to_string(),
malo_offset: 0,
sparte: Sparte::Strom,
unit: Sparte::Strom.billing_unit(),
}
}
pub fn sparte(mut self, sparte: Sparte) -> Self {
self.sparte = sparte;
self.unit = sparte.billing_unit();
self
}
pub fn in_unit(mut self, unit: MeasurementUnit) -> Self {
self.unit = unit;
self
}
pub fn resolution(mut self, resolution: Duration) -> Result<Self> {
let seconds = resolution.whole_seconds();
if seconds <= 0 || Duration::DAY.whole_seconds() % seconds != 0 {
return Err(Error::config(format!(
"workload resolution {resolution} must be a positive divisor of 24 h"
)));
}
self.resolution = resolution;
Ok(self)
}
pub fn seed(mut self, seed: u64) -> Self {
self.seed = seed;
self
}
pub fn malo_ids(mut self, n: usize) -> Self {
self.malo_ids = n.max(1);
self
}
pub fn days(mut self, n: i64) -> Self {
self.days = n.max(1);
self
}
pub fn malo_offset(mut self, offset: usize) -> Self {
self.malo_offset = offset;
self
}
pub fn with_corrections(mut self, rate: f64) -> Self {
self.correction_rate = rate.clamp(0.0, 1.0);
self
}
pub fn with_gaps(mut self, rate: f64) -> Self {
self.gap_rate = rate.clamp(0.0, 1.0);
self
}
pub fn spanning_spring_forward(mut self) -> Self {
self.start = time::macros::datetime!(2026-03-28 00:00 UTC);
self.days = self.days.max(3);
self
}
pub fn spanning_autumn_back(mut self) -> Self {
self.start = time::macros::datetime!(2026-10-24 00:00 UTC);
self.days = self.days.max(3);
self
}
pub fn range(&self) -> (OffsetDateTime, OffsetDateTime) {
(self.start, self.start + Duration::days(self.days))
}
pub fn generate(&self) -> Result<Vec<StoredSeries>> {
let mut rng = Rng::new(self.seed);
let mut out = Vec::new();
let mut corrections: Vec<(usize, OffsetDateTime, Decimal)> = Vec::new();
for day in 0..self.days {
let day_start = self.start + Duration::days(day);
let steps = Duration::DAY.whole_seconds() / self.resolution.whole_seconds();
for malo in 0..self.malo_ids {
let mut intervals = Vec::new();
for step in 0..steps {
let from = day_start + self.resolution * step as i32;
if rng.chance(self.gap_rate) {
continue;
}
let kwh = Decimal::new(rng.below(500) as i64 + 1, 2);
if rng.chance(self.correction_rate) {
corrections.push((malo, from, kwh + Decimal::new(100, 2)));
}
intervals.push(self.interval(from, kwh));
}
for (_, group) in group_by_scope(intervals) {
let anchor = group[0].from;
out.push(self.stored(malo, group, FIRST_VERSION, anchor)?);
}
}
}
let mut by_malo: BTreeMap<usize, Vec<MeterInterval>> = BTreeMap::new();
for (malo, from, kwh) in corrections {
by_malo
.entry(malo)
.or_default()
.push(self.interval(from, kwh));
}
for (malo, mut intervals) in by_malo {
intervals.sort_by_key(|i| i.from);
for (_, group) in group_by_scope(intervals) {
let anchor = group[0].from;
out.push(self.stored(malo, group, CORRECTION_VERSION, anchor)?);
}
}
Ok(out)
}
fn interval(&self, from: OffsetDateTime, kwh: Decimal) -> MeterInterval {
MeterInterval {
from,
to: from + self.resolution,
value: kwh,
quality: QualityFlag::Measured,
obis_code: "1-0:1.8.0".parse().ok(),
}
}
fn stored(
&self,
malo: usize,
intervals: Vec<MeterInterval>,
version: u128,
scope_anchor: OffsetDateTime,
) -> Result<StoredSeries> {
let malo_id = malo_id(self.malo_offset + malo);
let recorded_at = intervals.last().map(|i| i.to).unwrap_or(scope_anchor);
let mut series = MeasurementSeries::new(
malo_id,
"1-0:1.8.0".parse().ok(),
intervals,
MeasurementSource::Mscons {
pid: 13_005,
message_ref: None,
sender_mp_id: self.operator.clone(),
},
recorded_at,
);
series.resolution = Some(self.declared_resolution());
Ok(StoredSeries::of(
self.sparte,
series,
ScopedVersion::new(
VersionScope::for_interval(&self.operator, scope_anchor)?,
Version::new(version)?,
),
recorded_at,
)
.in_unit(self.unit))
}
fn declared_resolution(&self) -> IntervalResolution {
let seconds = u32::try_from(self.resolution.whole_seconds())
.expect("the builder rejects non-positive resolutions");
IntervalResolution::from_seconds(seconds)
.expect("the builder rejects a zero-length resolution")
}
}
fn group_by_scope(intervals: Vec<MeterInterval>) -> Vec<(time::Date, Vec<MeterInterval>)> {
let mut out: Vec<(time::Date, Vec<MeterInterval>)> = Vec::new();
for interval in intervals {
let month = metering::calendar::local_month(interval.from);
match out.last_mut() {
Some((m, group)) if *m == month => group.push(interval),
_ => out.push((month, vec![interval])),
}
}
out
}
const FIRST_VERSION: u128 = 20_260_101_000_001;
const CORRECTION_VERSION: u128 = 20_260_201_000_002;
fn malo_id(n: usize) -> MaloId {
let prefix = format!("{:010}", 1_000_000_000u64 + n as u64);
let check = MaloId::compute_check_digit(&prefix).expect("ten ASCII digits");
format!("{prefix}{check}")
.parse()
.expect("a computed check digit is the one the parser recomputes")
}
type OracleKey = (String, String, OffsetDateTime, Vec<String>);
#[derive(Debug, Default, Clone)]
pub struct Oracle {
identity: Vec<String>,
rows: BTreeMap<OracleKey, (String, u128, Decimal)>,
}
impl Oracle {
pub fn new() -> Self {
Self::default()
}
pub fn for_table(config: &crate::config::ValidatedTableConfig) -> Self {
Self {
identity: config
.identity_columns()
.iter()
.map(|f| f.name().clone())
.collect(),
rows: BTreeMap::new(),
}
}
fn identity_of(&self, stored: &StoredSeries) -> Result<Vec<String>> {
self.identity
.iter()
.map(|name| match stored.extra.get(name) {
Some(datafusion::common::ScalarValue::Utf8(Some(value))) => Ok(value.clone()),
_ => Err(Error::encode(
name,
format!(
"{} declares {name:?} as an identity column, but this series carries no \
value for it — the store would have refused the write",
stored.series.malo_id
),
)),
})
.collect()
}
pub fn record(&mut self, series: &[StoredSeries]) -> Result<()> {
for stored in series {
let scope = stored.version.scope().as_str().to_string();
let version = stored.version.version().get();
let identity = self.identity_of(stored)?;
for interval in &stored.series.intervals {
let obis = interval
.obis_code
.or(stored.series.obis_code)
.ok_or_else(|| Error::encode("obis_code", "no channel on interval or series"))?
.to_string();
let key = (
stored.series.malo_id.to_string(),
obis,
interval.from,
identity.clone(),
);
match self.rows.get(&key) {
Some((existing, seen, _)) if *existing == scope => {
if version > *seen {
self.rows
.insert(key, (scope.clone(), version, interval.value));
}
}
Some((existing, _, _)) => {
return Err(Error::VersionScopeMismatch {
left: existing.clone(),
right: scope,
});
}
None => {
self.rows
.insert(key, (scope.clone(), version, interval.value));
}
}
}
}
Ok(())
}
pub fn row_count(&self, from: OffsetDateTime, to: OffsetDateTime) -> u64 {
self.rows
.keys()
.filter(|(_, _, start, _)| *start >= from && *start < to)
.count() as u64
}
pub fn sum_kwh(&self, from: OffsetDateTime, to: OffsetDateTime) -> Decimal {
self.rows
.iter()
.filter(|((_, _, start, _), _)| *start >= from && *start < to)
.map(|(_, (_, _, value))| *value)
.sum()
}
pub fn sum_kwh_for(&self, malo_id: &str, from: OffsetDateTime, to: OffsetDateTime) -> Decimal {
self.rows
.iter()
.filter(|((malo, _, start, _), _)| malo == malo_id && *start >= from && *start < to)
.map(|(_, (_, _, value))| *value)
.sum()
}
pub fn malo_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self
.rows
.keys()
.map(|(malo, _, _, _)| malo.clone())
.collect();
ids.sort();
ids.dedup();
ids
}
pub fn len(&self) -> usize {
self.rows.len()
}
pub fn is_empty(&self) -> bool {
self.rows.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use time::macros::datetime;
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
#[test]
fn an_oracle_over_a_tenant_table_keeps_the_tenants_apart() {
use crate::arrow::datatypes::{DataType, Field};
use crate::config::TableConfig;
use datafusion::common::ScalarValue;
let config = TableConfig::new("readings_versions")
.identity_column(Field::new("tenant", DataType::Utf8, false))
.build()
.expect("config");
let base = MeteringWorkload::new(START).malo_ids(1).days(1);
let for_tenant = |tenant: &str| -> Vec<StoredSeries> {
base.clone()
.generate()
.expect("workload")
.into_iter()
.map(|s| s.with_extra("tenant", ScalarValue::Utf8(Some(tenant.into()))))
.collect()
};
let mut aware = Oracle::for_table(&config);
aware.record(&for_tenant("a")).expect("tenant a");
aware.record(&for_tenant("b")).expect("tenant b");
let mut naive = Oracle::new();
naive.record(&for_tenant("a")).expect("tenant a");
naive.record(&for_tenant("b")).expect("tenant b");
assert_eq!(
aware.len(),
naive.len() * 2,
"the tenant-aware reference holds both tenants' readings"
);
}
#[test]
fn a_missing_identity_value_is_refused_rather_than_defaulted() {
use crate::arrow::datatypes::{DataType, Field};
use crate::config::TableConfig;
let config = TableConfig::new("readings_versions")
.identity_column(Field::new("tenant", DataType::Utf8, false))
.build()
.expect("config");
let series = MeteringWorkload::new(START)
.malo_ids(1)
.days(1)
.generate()
.expect("workload");
let err = Oracle::for_table(&config)
.record(&series)
.expect_err("a series with no tenant must be refused");
assert!(err.to_string().contains("tenant"), "{err}");
}
#[test]
fn a_redelivery_at_the_same_version_keeps_the_first_value() {
let mut first = MeteringWorkload::new(START)
.malo_ids(1)
.days(1)
.generate()
.expect("workload");
let mut restated = first.clone();
for s in &mut restated {
for i in &mut s.series.intervals {
i.value += Decimal::ONE;
}
}
let mut oracle = Oracle::new();
oracle.record(&first).expect("first");
let before = oracle.sum_kwh(START, START + Duration::days(1));
oracle.record(&restated).expect("restated");
assert_eq!(
oracle.sum_kwh(START, START + Duration::days(1)),
before,
"an equal version must not overwrite"
);
first.clear();
}
#[test]
fn the_generator_is_reproducible_from_its_seed() {
let workload = MeteringWorkload::new(START).seed(42).with_corrections(0.1);
let a = workload.generate().unwrap();
let b = workload.generate().unwrap();
assert_eq!(a.len(), b.len());
for (x, y) in a.iter().zip(&b) {
assert_eq!(x.series.malo_id, y.series.malo_id);
assert_eq!(x.series.intervals.len(), y.series.intervals.len());
assert_eq!(x.version, y.version);
}
}
#[test]
fn different_seeds_produce_different_workloads() {
let a = MeteringWorkload::new(START)
.seed(1)
.with_gaps(0.2)
.generate()
.unwrap();
let b = MeteringWorkload::new(START)
.seed(2)
.with_gaps(0.2)
.generate()
.unwrap();
let rows =
|v: &[StoredSeries]| -> usize { v.iter().map(|s| s.series.intervals.len()).sum() };
assert_ne!(rows(&a), rows(&b), "one seed must not stand in for another");
}
#[test]
fn a_workload_with_no_corrections_has_one_version_per_key() {
let series = MeteringWorkload::new(START)
.malo_ids(3)
.days(2)
.generate()
.unwrap();
assert!(
series
.iter()
.all(|s| s.version.version().get() == FIRST_VERSION),
"corrections must be opt-in: elision depends on them being rare"
);
}
#[test]
fn corrections_arrive_after_the_deliveries_they_correct() {
let series = MeteringWorkload::new(START)
.malo_ids(2)
.days(2)
.seed(7)
.with_corrections(0.5)
.generate()
.unwrap();
let first_correction = series
.iter()
.position(|s| s.version.version().get() == CORRECTION_VERSION)
.expect("the workload must produce corrections");
assert!(
series[..first_correction]
.iter()
.all(|s| s.version.version().get() == FIRST_VERSION)
);
}
#[test]
fn the_oracle_keeps_the_highest_version_within_a_scope() {
let workload = MeteringWorkload::new(START).malo_ids(1).days(1);
let base = workload.generate().unwrap();
let mut oracle = Oracle::new();
oracle.record(&base).unwrap();
let before = oracle.sum_kwh(START, START + Duration::DAY);
let mut corrected = base.clone();
for stored in &mut corrected {
for interval in &mut stored.series.intervals {
interval.value += Decimal::ONE;
}
stored.version = ScopedVersion::new(
stored.version.scope().clone(),
Version::new(CORRECTION_VERSION).unwrap(),
);
}
oracle.record(&corrected).unwrap();
let intervals = oracle.row_count(START, START + Duration::DAY);
assert_eq!(
oracle.sum_kwh(START, START + Duration::DAY),
before + Decimal::from(intervals),
"each interval counted once, at its corrected value"
);
}
#[test]
fn the_oracle_counts_a_corrected_interval_once() {
let workload = MeteringWorkload::new(START).malo_ids(2).days(1).seed(9);
let series = workload.generate().unwrap();
let distinct: usize = series
.iter()
.flat_map(|s| {
s.series
.intervals
.iter()
.map(|i| (s.series.malo_id, i.from))
})
.collect::<std::collections::BTreeSet<_>>()
.len();
let mut oracle = Oracle::new();
oracle.record(&series).unwrap();
assert_eq!(oracle.len(), distinct);
}
#[test]
fn a_gap_rate_actually_removes_intervals() {
let full = MeteringWorkload::new(START).malo_ids(4).days(1).seed(3);
let holey = full.clone().with_gaps(0.25);
let rows = |w: &MeteringWorkload| -> usize {
w.generate()
.unwrap()
.iter()
.map(|s| s.series.intervals.len())
.sum()
};
assert!(rows(&holey) < rows(&full));
}
#[test]
fn the_dst_workloads_span_their_transition() {
let spring = MeteringWorkload::new(START).spanning_spring_forward();
let (from, to) = spring.range();
assert!(from <= time::macros::datetime!(2026-03-29 00:00 UTC));
assert!(to > time::macros::datetime!(2026-03-29 00:00 UTC));
let autumn = MeteringWorkload::new(START).spanning_autumn_back();
let (from, to) = autumn.range();
assert!(from <= time::macros::datetime!(2026-10-25 00:00 UTC));
assert!(to > time::macros::datetime!(2026-10-25 00:00 UTC));
}
#[test]
fn generated_malo_ids_carry_a_valid_check_digit() {
for n in [0, 1, 42, 999, 123_456] {
let id = malo_id(n);
let text = id.to_string();
assert_eq!(text.len(), 11);
assert!(text.chars().all(|c| c.is_ascii_digit()));
assert_eq!(text.parse::<MaloId>().unwrap(), id);
}
let ids: std::collections::BTreeSet<_> = (0..64).map(malo_id).collect();
assert_eq!(ids.len(), 64);
}
#[test]
fn a_workload_spanning_a_month_boundary_splits_its_correction_scopes() {
let series = MeteringWorkload::new(datetime!(2026-07-30 00:00 UTC))
.malo_ids(1)
.days(4)
.seed(11)
.with_corrections(0.5)
.generate()
.unwrap();
for stored in &series {
for interval in &stored.series.intervals {
assert!(
stored.version.scope().covers(interval.from),
"scope {} does not cover {}",
stored.version.scope(),
interval.from
);
}
}
}
}