#![cfg(feature = "testkit")]
use metering::QualityFlag;
use metering::interval::MeterInterval;
use metering::measurement_series::{MeasurementSeries, MeasurementSource};
use meterstore::testkit::TestHarness;
use rust_decimal::Decimal;
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
const T1: OffsetDateTime = datetime!(2026-07-26 06:00 UTC);
const T_MID: OffsetDateTime = datetime!(2026-07-26 12:00 UTC);
const T2: OffsetDateTime = datetime!(2026-07-27 06:00 UTC);
fn stored(
from: OffsetDateTime,
kwh: i64,
version: u128,
recorded_at: OffsetDateTime,
) -> meterstore::encode::StoredSeries {
let series = MeasurementSeries::new(
"12345678905".parse().unwrap(),
"1-0:1.8.0".parse().ok(),
vec![MeterInterval {
from,
to: from + Duration::minutes(15),
value: Decimal::new(kwh, 0),
quality: QualityFlag::Measured,
obis_code: "1-0:1.8.0".parse().ok(),
}],
MeasurementSource::Mscons {
pid: 13_005,
message_ref: None,
sender_mp_id: "9900000000001".parse().expect("a valid Marktpartner-ID"),
},
recorded_at,
);
meterstore::encode::StoredSeries::new(
series,
meterstore::ScopedVersion::new(
meterstore::VersionScope::for_interval(
"9900000000001",
from,
metering::interval::Sparte::Strom,
)
.unwrap(),
meterstore::Version::new(version).unwrap(),
),
recorded_at,
)
}
async fn store() -> (TestHarness, meterstore::MeterStore) {
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(START, START + Duration::days(2))
.await
.expect("partitions");
harness.seed_watermark(START).await.expect("watermark");
let store = harness.store().await.expect("store");
(harness, store)
}
async fn value_at(store: &meterstore::MeterStore, from: OffsetDateTime) -> Option<Decimal> {
let series = store
.series("12345678905")
.unwrap()
.range(from, from + Duration::minutes(15))
.collect()
.await
.expect("collect");
series.and_then(|s| s.intervals.first().map(|i| i.value))
}
#[tokio::test]
async fn as_known_at_sees_the_value_in_force_at_that_time() {
let (_h, store) = store().await;
store
.append(&[stored(START, 10, 20_260_726_060_000, T1)])
.await
.expect("append original");
store
.append(&[stored(START, 42, 20_260_727_060_000, T2)])
.await
.expect("append correction");
assert_eq!(value_at(&store, START).await, Some(Decimal::new(42, 0)));
let mid = store.as_known_at(T_MID).await.expect("as_known_at");
assert_eq!(
value_at(&mid, START).await,
Some(Decimal::new(10, 0)),
"the correction was not yet known at T_MID"
);
let after = store.as_known_at(T2).await.expect("as_known_at");
assert_eq!(value_at(&after, START).await, Some(Decimal::new(42, 0)));
}
#[tokio::test]
async fn as_known_at_excludes_an_interval_first_stored_later() {
let (_h, store) = store().await;
store
.append(&[stored(START, 10, 20_260_726_060_000, T1)])
.await
.expect("append first interval");
let later = START + Duration::minutes(15);
store
.append(&[stored(later, 20, 20_260_727_060_000, T2)])
.await
.expect("append later interval");
let mid = store.as_known_at(T_MID).await.expect("as_known_at");
assert_eq!(
value_at(&mid, START).await,
Some(Decimal::new(10, 0)),
"the interval known at T_MID is present"
);
assert_eq!(
value_at(&mid, later).await,
None,
"the interval first stored at T2 did not exist at T_MID"
);
assert_eq!(value_at(&store, later).await, Some(Decimal::new(20, 0)));
}
#[tokio::test]
async fn as_known_at_reads_the_same_after_archival_to_cold() {
let (_h, store) = store().await;
store
.append(&[stored(START, 10, 20_260_726_060_000, T1)])
.await
.expect("append original");
store
.append(&[stored(START, 42, 20_260_727_060_000, T2)])
.await
.expect("append correction");
store
.archive(START + Duration::days(2), 8)
.await
.expect("archive");
assert!(
store.watermark().await.expect("watermark").get() > START,
"the interval's day must have moved into cold"
);
let mid = store.as_known_at(T_MID).await.expect("as_known_at");
assert_eq!(
value_at(&mid, START).await,
Some(Decimal::new(10, 0)),
"the as-of value is the same read from cold"
);
let after = store.as_known_at(T2).await.expect("as_known_at");
assert_eq!(value_at(&after, START).await, Some(Decimal::new(42, 0)));
}
#[tokio::test]
async fn the_ceiling_reaches_the_raw_versions_relation_too() {
let (_h, store) = store().await;
store
.append(&[stored(START, 10, 20_260_726_060_000, T1)])
.await
.expect("append original");
store
.append(&[stored(START, 42, 20_260_727_060_000, T2)])
.await
.expect("append correction");
let raw = store.raw_table();
let count = |s: meterstore::MeterStore, sql: String| async move {
let result = s.query(&sql).await.expect("query");
result.batches().iter().map(|b| b.num_rows()).sum::<usize>()
};
assert_eq!(
count(store.clone(), format!("SELECT version FROM {raw}")).await,
2
);
let mid = store.as_known_at(T_MID).await.expect("as_known_at");
assert_eq!(
count(mid, format!("SELECT version FROM {raw}")).await,
1,
"only the version recorded by T_MID was known then"
);
}
#[tokio::test]
async fn a_ceiling_before_any_delivery_is_empty() {
let (_h, store) = store().await;
store
.append(&[stored(START, 10, 20_260_726_060_000, T1)])
.await
.expect("append");
let before = store
.as_known_at(datetime!(2026-07-01 00:00 UTC))
.await
.expect("as_known_at");
assert_eq!(
value_at(&before, START).await,
None,
"nothing had been recorded yet"
);
}