use std::collections::BTreeMap;
use datafusion::common::ScalarValue;
use metering::ids::{MaloId, MeloId};
use metering::interval::QualityFlag;
use metering::interval::Sparte;
use metering::measurement_series::{MeasurementSeries, MeasurementSource, ProvenanceEntry};
use metering::obis::ObisCode;
use time::OffsetDateTime;
use crate::encode::schema::col;
use crate::error::{Error, Result};
#[derive(Debug, Clone)]
pub struct ResolvedSeries {
pub sparte: Sparte,
pub extra: BTreeMap<String, ScalarValue>,
pub series: MeasurementSeries,
}
#[derive(Debug, Clone)]
pub struct SeriesQuery<'a> {
store: &'a crate::session::MeterStore,
malo_id: MaloId,
obis_code: Option<String>,
from: Option<OffsetDateTime>,
to: Option<OffsetDateTime>,
filters: Vec<(String, ScalarValue)>,
quality: Vec<QualityFlag>,
latest_only: bool,
}
impl<'a> SeriesQuery<'a> {
pub(crate) fn new(store: &'a crate::session::MeterStore, malo_id: MaloId) -> Self {
Self {
store,
malo_id,
obis_code: None,
from: None,
to: None,
filters: Vec::new(),
quality: Vec::new(),
latest_only: false,
}
}
pub fn column_eq(mut self, name: &str, value: ScalarValue) -> Result<Self> {
let mut accepted: Vec<String> = self
.store
.config()
.extra_columns()
.iter()
.map(|f| f.name().clone())
.collect();
for column in self.store.config().discriminator_columns() {
if !accepted.contains(&column) {
accepted.push(column);
}
}
if !accepted.iter().any(|c| c == name) {
return Err(Error::config(format!(
"{name:?} is not a filterable column of {}: this store accepts [{}]. \
Column names are written into SQL as identifiers, which cannot be \
parameterised, so only declared ones are accepted",
self.store.table(),
accepted.join(", "),
)));
}
self.filters.push((name.to_string(), value));
Ok(self)
}
#[must_use]
pub fn quality_in(mut self, flags: &[QualityFlag]) -> Self {
self.quality.extend_from_slice(flags);
self
}
pub fn obis(mut self, obis_code: &str) -> Result<Self> {
self.obis_code = Some(crate::encode::canonical_obis(obis_code)?);
Ok(self)
}
pub fn range(mut self, from: OffsetDateTime, to: OffsetDateTime) -> Self {
self.from = Some(from);
self.to = Some(to);
self
}
pub fn since(mut self, from: OffsetDateTime) -> Self {
self.from = Some(from);
self
}
pub fn until(mut self, to: OffsetDateTime) -> Self {
self.to = Some(to);
self
}
pub async fn collect(self) -> Result<Option<MeasurementSeries>> {
Ok(self.collect_with_provenance().await?.0)
}
pub async fn collect_with_sparte(self) -> Result<Option<(Sparte, MeasurementSeries)>> {
Ok(self.resolve().await?.0.map(|r| (r.sparte, r.series)))
}
pub async fn collect_resolved(self) -> Result<Option<ResolvedSeries>> {
Ok(self.resolve().await?.0)
}
pub async fn latest(self) -> Result<Option<metering::interval::MeterInterval>> {
Ok(self
.latest_resolved()
.await?
.and_then(|r| r.series.intervals.into_iter().next_back()))
}
pub async fn latest_resolved(mut self) -> Result<Option<ResolvedSeries>> {
self.latest_only = true;
Ok(self.resolve().await?.0)
}
pub async fn intervals(self) -> Result<Vec<metering::interval::MeterInterval>> {
Ok(self
.collect()
.await?
.map(|s| s.intervals)
.unwrap_or_default())
}
pub async fn collect_with_provenance(
self,
) -> Result<(Option<MeasurementSeries>, super::QueryResult)> {
let (resolved, result) = self.resolve().await?;
Ok((resolved.map(|r| r.series), result))
}
pub async fn channels(&self) -> Result<Vec<ObisCode>> {
let (conditions, params) = self.predicate();
let sql = format!(
r#"SELECT DISTINCT "{obis}" FROM "{table}" WHERE {conditions} ORDER BY 1"#,
obis = col::OBIS_CODE,
table = self.store.resolved_table(),
conditions = conditions.join(" AND "),
);
let result = self.store.query_with_params(&sql, params).await?;
let mut out = Vec::new();
for batch in result.batches() {
let codes =
crate::encode::column::<crate::arrow::array::StringArray>(batch, col::OBIS_CODE)?;
for i in 0..batch.num_rows() {
out.push(codes.value(i).parse::<ObisCode>().map_err(|e| {
Error::decode(
col::OBIS_CODE,
format!("{:?} is not an OBIS code: {e}", codes.value(i)),
)
})?);
}
}
out.sort_unstable();
out.dedup();
Ok(out)
}
pub async fn collect_by_channel(self) -> Result<BTreeMap<ObisCode, ResolvedSeries>> {
Ok(self.collect_by_channel_with_provenance().await?.0)
}
pub async fn collect_by_channel_with_provenance(
self,
) -> Result<(BTreeMap<ObisCode, ResolvedSeries>, super::QueryResult)> {
let malo_id = self.malo_id;
let discriminators = self.store.config().discriminator_columns();
let (stored, result) = self.scan().await?;
Ok((split_by_channel(malo_id, &discriminators, stored)?, result))
}
fn predicate(&self) -> (Vec<String>, Vec<ScalarValue>) {
let mut conditions = vec![format!(r#""{}" = $1"#, col::MALO_ID)];
let mut params: Vec<ScalarValue> = vec![ScalarValue::Utf8(Some(self.malo_id.to_string()))];
if let Some(obis) = &self.obis_code {
conditions.push(format!(r#""{}" = ${}"#, col::OBIS_CODE, params.len() + 1));
params.push(ScalarValue::Utf8(Some(obis.clone())));
}
if let Some(from) = self.from {
conditions.push(format!(r#""{}" >= ${}"#, col::FROM, params.len() + 1));
params.push(crate::encode::schema::timestamp_scalar(from));
}
if let Some(to) = self.to {
conditions.push(format!(r#""{}" < ${}"#, col::FROM, params.len() + 1));
params.push(crate::encode::schema::timestamp_scalar(to));
}
for (name, value) in &self.filters {
conditions.push(format!(r#""{}" = ${}"#, name, params.len() + 1));
params.push(value.clone());
}
if !self.quality.is_empty() {
let placeholders = (0..self.quality.len())
.map(|i| format!("${}", params.len() + 1 + i))
.collect::<Vec<_>>()
.join(", ");
conditions.push(format!(r#""{}" IN ({placeholders})"#, col::QUALITY));
for q in &self.quality {
params.push(ScalarValue::Utf8(Some(q.as_str().to_owned())));
}
}
(conditions, params)
}
async fn scan(&self) -> Result<(Vec<crate::encode::StoredSeries>, super::QueryResult)> {
let (conditions, params) = self.predicate();
let tail = if self.latest_only {
format!(
r#"ORDER BY "{from}" DESC, {rest} LIMIT 1"#,
from = col::FROM,
rest = key_order_without_start(&self.store.config().merge_key()),
)
} else {
format!(
"ORDER BY {}",
merge_key_order(&self.store.config().merge_key()),
)
};
let sql = format!(
r#"SELECT * FROM "{table}" WHERE {conditions} {tail}"#,
table = self.store.resolved_table(),
conditions = conditions.join(" AND "),
);
let result = self.store.query_with_params(&sql, params).await?;
let mut stored = Vec::new();
for batch in result.batches() {
stored.extend(crate::encode::from_record_batch(batch)?);
}
Ok((stored, result))
}
async fn resolve(self) -> Result<(Option<ResolvedSeries>, super::QueryResult)> {
let malo_id = self.malo_id;
let obis_code = self.obis_code.clone();
let discriminators = self.store.config().discriminator_columns();
let (stored, result) = self.scan().await?;
Ok((
merge(malo_id, obis_code.as_deref(), &discriminators, stored)?,
result,
))
}
}
pub(crate) fn merge_key_order(merge_key: &[String]) -> String {
merge_key
.iter()
.filter(|c| c.as_str() != col::FROM)
.map(|c| format!("\"{c}\""))
.chain(std::iter::once(format!("\"{}\"", col::FROM)))
.collect::<Vec<_>>()
.join(", ")
}
pub(crate) fn key_order_without_start(merge_key: &[String]) -> String {
merge_key
.iter()
.filter(|c| c.as_str() != col::FROM)
.map(|c| format!("\"{c}\""))
.collect::<Vec<_>>()
.join(", ")
}
fn merge(
malo_id: MaloId,
obis_code: Option<&str>,
discriminators: &[String],
mut stored: Vec<crate::encode::StoredSeries>,
) -> Result<Option<ResolvedSeries>> {
if stored.is_empty() {
return Ok(None);
}
refuse_mixed_readings(&stored, discriminators)?;
stored.sort_by_key(|s| s.recorded_at);
let mut intervals = Vec::new();
let mut source: Option<MeasurementSource> = None;
let mut provenance: Vec<ProvenanceEntry> = Vec::new();
let mut melo_id: Option<MeloId> = None;
let mut resolution = None;
let mut recorded_at: Option<OffsetDateTime> = None;
let mut sparte: Option<Sparte> = None;
let mut extra: BTreeMap<String, ScalarValue> = BTreeMap::new();
for series in stored {
intervals.extend(series.series.intervals);
source = Some(series.series.source);
provenance.extend(series.series.provenance);
if series.series.melo_id.is_some() {
melo_id = series.series.melo_id;
}
if series.series.resolution.is_some() {
resolution = series.series.resolution;
}
recorded_at = Some(series.recorded_at);
sparte = Some(series.sparte);
extra = series.extra;
}
intervals.sort_by_key(|i| i.from);
refuse_repeated_instants(malo_id, &intervals)?;
let obis = match obis_code {
Some(code) => Some(code.parse::<ObisCode>().map_err(|e| {
Error::decode(col::OBIS_CODE, format!("{code:?} is not an OBIS code: {e}"))
})?),
None => single_channel(&intervals),
};
let (Some(source), Some(recorded_at), Some(sparte)) = (source, recorded_at, sparte) else {
return Ok(None);
};
let mut series = MeasurementSeries::new(malo_id, obis, intervals, source, recorded_at);
series.melo_id = melo_id;
series.resolution = resolution;
if !provenance.is_empty() {
series.provenance = provenance;
}
Ok(Some(ResolvedSeries {
sparte,
extra,
series,
}))
}
fn refuse_repeated_instants(
malo_id: MaloId,
intervals: &[metering::interval::MeterInterval],
) -> Result<()> {
for pair in intervals.windows(2) {
if pair[0].from != pair[1].from {
continue;
}
return Err(Error::InvariantViolated {
table: malo_id.to_string(),
detail: format!(
"two values survived resolution for one reading at {at}: {a} and {b}. \
Resolution partitions by the merge key and version_scope, so the cause \
is almost always two network operators for one reading — which both \
write paths refuse, unless integrity constraints are off or something \
other than meterstore wrote these rows. Folding them into one series \
would sum to twice the truth",
at = pair[0].from,
a = pair[0].value,
b = pair[1].value,
),
});
}
Ok(())
}
type ReadingKey = (Option<String>, Vec<(String, String)>);
fn refuse_mixed_readings(
stored: &[crate::encode::StoredSeries],
discriminators: &[String],
) -> Result<()> {
let key_of = |s: &crate::encode::StoredSeries| -> Result<ReadingKey> {
Ok((
s.series.obis_code.map(|c| c.to_string()),
crate::session::store::discriminator_values(
s.series.melo_id.as_ref(),
&s.extra,
discriminators,
)?,
))
};
let first = key_of(&stored[0])?;
for other in &stored[1..] {
let key = key_of(other)?;
if key == first {
continue;
}
let (channel, identity) = (&first.0, &first.1);
let differs = if channel != &key.0 {
format!(
"two channels ({} and {})",
channel.as_deref().unwrap_or("<none>"),
key.0.as_deref().unwrap_or("<none>"),
)
} else {
format!("two readings ({} and {})", render(identity), render(&key.1))
};
return Err(Error::config(format!(
"{malo} spans {differs} over this range, and a MeasurementSeries can only \
describe one — folding them puts two values at the same instant into one \
series, which sums to twice the truth with nothing to notice it. Narrow the \
read with .obis(..) or .column_eq(..), or scope the session",
malo = stored[0].series.malo_id,
)));
}
Ok(())
}
fn split_by_channel(
malo_id: MaloId,
discriminators: &[String],
stored: Vec<crate::encode::StoredSeries>,
) -> Result<BTreeMap<ObisCode, ResolvedSeries>> {
let mut grouped: BTreeMap<ObisCode, Vec<crate::encode::StoredSeries>> = BTreeMap::new();
for series in stored {
let channel = series.series.obis_code.ok_or_else(|| {
Error::decode(
col::OBIS_CODE,
format!(
"a decoded series for {malo_id} carries no channel, so it cannot be placed \
in a per-channel read"
),
)
})?;
grouped.entry(channel).or_default().push(series);
}
let mut out = BTreeMap::new();
for (channel, series) in grouped {
if let Some(resolved) = merge(malo_id, Some(&channel.to_string()), discriminators, series)?
{
out.insert(channel, resolved);
}
}
Ok(out)
}
fn render(values: &[(String, String)]) -> String {
match values.is_empty() {
true => "<none>".to_string(),
false => values
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect::<Vec<_>>()
.join(", "),
}
}
fn single_channel(intervals: &[metering::interval::MeterInterval]) -> Option<ObisCode> {
let mut seen: Option<ObisCode> = None;
for interval in intervals {
match (seen, interval.obis_code) {
(_, None) => return None,
(None, Some(code)) => seen = Some(code),
(Some(a), Some(b)) if a == b => {}
_ => return None,
}
}
seen
}
#[cfg(test)]
mod tests {
use super::*;
use crate::encode::StoredSeries;
use crate::version::{ScopedVersion, Version, VersionScope};
use metering::interval::{MeterInterval, QualityFlag};
use rust_decimal::Decimal;
use time::macros::datetime;
fn interval(from: OffsetDateTime, kwh: i64) -> MeterInterval {
MeterInterval {
from,
to: from + time::Duration::minutes(15),
value: Decimal::new(kwh, 0),
quality: QualityFlag::Measured,
obis_code: Some("1-0:1.8.0".parse().unwrap()),
}
}
fn malo() -> MaloId {
"12345678905".parse().unwrap()
}
fn source() -> MeasurementSource {
MeasurementSource::Mscons {
pid: 13_005,
message_ref: Some("MSG-1".to_owned()),
sender_mp_id: "9900000000001".parse().expect("a valid Marktpartner-ID"),
}
}
fn stored(intervals: Vec<MeterInterval>, recorded_at: OffsetDateTime) -> StoredSeries {
let scope =
VersionScope::for_interval("9900000000001", intervals[0].from, Sparte::Strom).unwrap();
StoredSeries::new(
MeasurementSeries::new(
malo(),
Some("1-0:1.8.0".parse().unwrap()),
intervals,
source(),
recorded_at,
),
ScopedVersion::new(scope, Version::new(20_260_701_000_001).unwrap()),
recorded_at,
)
}
fn one_interval(obis: &str) -> StoredSeries {
let from = datetime!(2026-07-20 00:00 UTC);
let mut s = stored(
vec![MeterInterval {
from,
to: from + time::Duration::minutes(15),
value: rust_decimal::Decimal::new(15, 1),
quality: metering::QualityFlag::Measured,
obis_code: obis.parse().ok(),
}],
from,
);
s.series.obis_code = obis.parse().ok();
s
}
#[test]
fn the_latest_read_orders_totally() {
let key = vec![
col::MALO_ID.to_string(),
col::MELO_ID.to_string(),
col::OBIS_CODE.to_string(),
col::FROM.to_string(),
"tenant".to_string(),
];
let rest = key_order_without_start(&key);
assert_eq!(
rest, r#""malo_id", "melo_id", "obis_code", "tenant""#,
"every merge-key column but the start, in key order"
);
assert!(!rest.contains(col::FROM), "the read already ordered on it");
assert!(
!key_order_without_start(&[col::MALO_ID.to_string(), col::FROM.to_string(),])
.is_empty()
);
}
#[test]
fn a_per_channel_read_splits_what_the_fold_refuses() {
let split = split_by_channel(
malo(),
&[],
vec![one_interval("1-0:1.8.0"), one_interval("1-0:2.8.0")],
)
.expect("two channels are two series, not an error");
assert_eq!(split.len(), 2);
let import: ObisCode = "1-0:1.8.0".parse().unwrap();
let export: ObisCode = "1-0:2.8.0".parse().unwrap();
assert_eq!(split[&import].series.obis_code, Some(import));
assert_eq!(split[&export].series.obis_code, Some(export));
assert_eq!(split[&import].series.intervals.len(), 1);
assert_eq!(split[&export].series.intervals.len(), 1);
assert_eq!(
split[&import].series.intervals[0].from,
split[&export].series.intervals[0].from
);
}
#[test]
fn a_per_channel_read_still_refuses_two_readings_on_one_channel() {
let mut a = one_interval("1-0:1.8.0");
let mut b = one_interval("1-0:1.8.0");
a.extra.insert(
"tenant".to_string(),
ScalarValue::Utf8(Some("alpha".to_string())),
);
b.extra.insert(
"tenant".to_string(),
ScalarValue::Utf8(Some("beta".to_string())),
);
let err = split_by_channel(malo(), &["tenant".to_string()], vec![a, b])
.expect_err("two tenants on one channel are two readings");
let msg = err.to_string();
assert!(msg.contains("two readings"), "{msg}");
assert!(msg.contains("tenant=alpha"), "{msg}");
let mut scoped = one_interval("1-0:1.8.0");
scoped.extra.insert(
"tenant".to_string(),
ScalarValue::Utf8(Some("alpha".to_string())),
);
let split = split_by_channel(malo(), &["tenant".to_string()], vec![scoped]).unwrap();
assert_eq!(split.len(), 1);
}
#[test]
fn channels_come_back_in_obis_order_not_alphabetical_order() {
let mut codes: Vec<ObisCode> = ["1-0:2.8.0", "1-0:10.8.0", "1-0:1.8.0"]
.iter()
.map(|c| c.parse().unwrap())
.collect();
codes.sort_unstable();
assert_eq!(
codes.iter().map(ToString::to_string).collect::<Vec<_>>(),
["1-0:1.8.0", "1-0:2.8.0", "1-0:10.8.0"],
);
let mut text: Vec<&str> = vec!["1-0:2.8.0", "1-0:10.8.0", "1-0:1.8.0"];
text.sort_unstable();
assert_ne!(
text,
codes.iter().map(ToString::to_string).collect::<Vec<_>>(),
"if the two agreed, sorting in Rust would prove nothing"
);
}
#[test]
fn a_per_channel_read_of_nothing_is_an_empty_map() {
assert!(
split_by_channel(malo(), &[], Vec::new())
.unwrap()
.is_empty()
);
}
#[test]
fn folding_two_channels_into_one_series_is_refused() {
let mut export = one_interval("1-0:2.8.0");
export.recorded_at += time::Duration::hours(1);
let err = merge(malo(), None, &[], vec![one_interval("1-0:1.8.0"), export])
.unwrap_err()
.to_string();
assert!(err.contains("two channels"), "{err}");
assert!(
err.contains("1-0:1.8.0") && err.contains("1-0:2.8.0"),
"{err}"
);
assert!(err.contains(".obis("), "the message names the fix: {err}");
}
#[test]
fn two_values_at_one_instant_are_refused_rather_than_summed() {
let mut second = one_interval("1-0:1.8.0");
second.series.intervals[0].value = rust_decimal::Decimal::new(99, 1);
second.recorded_at += time::Duration::hours(1);
let err = merge(malo(), None, &[], vec![one_interval("1-0:1.8.0"), second]).unwrap_err();
assert!(
matches!(err, Error::InvariantViolated { .. }),
"stored data that should not exist, not a delivery being refused: {err:?}"
);
let msg = err.to_string();
assert!(msg.contains("version_scope"), "{msg}");
assert!(
msg.contains("1.5") && msg.contains("9.9"),
"both values: {msg}"
);
}
#[test]
fn ordinary_deliveries_of_one_series_still_fold() {
let first = one_interval("1-0:1.8.0");
let mut later = one_interval("1-0:1.8.0");
later.series.intervals[0].from += time::Duration::minutes(15);
later.series.intervals[0].to += time::Duration::minutes(15);
later.recorded_at += time::Duration::hours(1);
let folded = merge(malo(), None, &[], vec![first, later])
.unwrap()
.expect("one series");
assert_eq!(folded.series.intervals.len(), 2);
}
#[test]
fn folding_two_tenants_into_one_series_is_refused() {
let of = |tenant: &str| {
let mut s = one_interval("1-0:1.8.0");
s.extra.insert(
"tenant".to_string(),
ScalarValue::Utf8(Some(tenant.to_string())),
);
s
};
let err = merge(
malo(),
Some("1-0:1.8.0"),
&["tenant".to_string()],
vec![of("a"), of("b")],
)
.unwrap_err()
.to_string();
assert!(err.contains("two readings"), "{err}");
assert!(
err.contains("tenant=a") && err.contains("tenant=b"),
"{err}"
);
assert!(err.contains(".column_eq("), "{err}");
}
#[test]
fn a_column_that_is_not_in_the_merge_key_does_not_split_a_series() {
let of = |bk: &str, quarters: i64| {
let mut s = one_interval("1-0:1.8.0");
let shift = time::Duration::minutes(15 * quarters);
s.series.intervals[0].from += shift;
s.series.intervals[0].to += shift;
s.extra.insert(
"bilanzkreis".to_string(),
ScalarValue::Utf8(Some(bk.into())),
);
s.recorded_at += time::Duration::hours(quarters);
s
};
let folded = merge(
malo(),
Some("1-0:1.8.0"),
&[],
vec![of("BK-1", 0), of("BK-2", 1)],
)
.unwrap()
.expect("one series");
assert_eq!(folded.series.intervals.len(), 2, "both groups are kept");
}
#[test]
fn decoded_groups_fold_back_into_one_series() {
let a = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 1)],
datetime!(2026-07-21 00:00 UTC),
);
let b = stored(
vec![interval(datetime!(2026-07-20 00:15 UTC), 2)],
datetime!(2026-07-22 00:00 UTC),
);
let ResolvedSeries { series, .. } = merge(malo(), Some("1-0:1.8.0"), &[], vec![b, a])
.unwrap()
.expect("rows were supplied");
assert_eq!(series.intervals.len(), 2);
assert_eq!(series.malo_id, malo());
}
#[test]
fn intervals_come_back_in_time_order() {
let later = stored(
vec![interval(datetime!(2026-07-20 12:00 UTC), 1)],
datetime!(2026-07-21 00:00 UTC),
);
let earlier = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 2)],
datetime!(2026-07-21 00:00 UTC),
);
let ResolvedSeries { series, .. } =
merge(malo(), Some("1-0:1.8.0"), &[], vec![later, earlier])
.unwrap()
.expect("rows were supplied");
assert_eq!(series.intervals[0].from, datetime!(2026-07-20 00:00 UTC));
assert_eq!(series.intervals[1].from, datetime!(2026-07-20 12:00 UTC));
}
#[test]
fn an_empty_read_is_absence_rather_than_an_empty_series() {
assert!(
merge(malo(), Some("1-0:1.8.0"), &[], Vec::new())
.unwrap()
.is_none()
);
}
#[test]
fn the_newest_delivery_supplies_the_series_level_fields() {
let mut old = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 1)],
datetime!(2026-07-21 00:00 UTC),
);
old.series.resolution = None;
let mut new = stored(
vec![interval(datetime!(2026-07-20 00:15 UTC), 2)],
datetime!(2026-07-25 00:00 UTC),
);
new.series.resolution = Some(metering::IntervalResolution::QuarterHour);
let ResolvedSeries { series, .. } = merge(malo(), Some("1-0:1.8.0"), &[], vec![new, old])
.unwrap()
.expect("rows were supplied");
assert_eq!(
series.resolution,
Some(metering::IntervalResolution::QuarterHour)
);
}
#[test]
fn a_mixed_channel_read_carries_no_series_level_obis() {
let mut other = interval(datetime!(2026-07-20 00:15 UTC), 2);
other.obis_code = Some("1-0:2.8.0".parse().unwrap());
let mixed = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 1), other],
datetime!(2026-07-21 00:00 UTC),
);
let ResolvedSeries { series, .. } = merge(malo(), None, &[], vec![mixed])
.unwrap()
.expect("rows were supplied");
assert_eq!(series.obis_code, None);
}
#[test]
fn a_single_channel_read_recovers_its_obis_without_being_told() {
let one = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 1)],
datetime!(2026-07-21 00:00 UTC),
);
let ResolvedSeries { series, .. } = merge(malo(), None, &[], vec![one])
.unwrap()
.expect("rows were supplied");
assert_eq!(series.obis_code, Some("1-0:1.8.0".parse().unwrap()));
}
#[test]
fn the_stored_provenance_trail_survives_the_fold() {
let mut s = stored(
vec![interval(datetime!(2026-07-20 00:00 UTC), 1)],
datetime!(2026-07-21 00:00 UTC),
);
let trail = s.series.provenance.clone();
assert!(!trail.is_empty(), "the fixture must carry a trail");
s.series.provenance = trail.clone();
let ResolvedSeries { series, .. } = merge(malo(), Some("1-0:1.8.0"), &[], vec![s])
.unwrap()
.expect("rows were supplied");
assert_eq!(series.provenance, trail);
}
#[test]
fn the_commodity_survives_the_fold() {
let scope = VersionScope::for_interval(
"9900000000001",
datetime!(2026-07-20 00:00 UTC),
Sparte::Strom,
)
.unwrap();
let gas = StoredSeries::of(
Sparte::Gas,
MeasurementSeries::new(
malo(),
Some("7-1:3.0.0".parse().unwrap()),
vec![interval(datetime!(2026-07-20 00:00 UTC), 1)],
source(),
datetime!(2026-07-21 00:00 UTC),
),
ScopedVersion::new(scope, Version::new(20_260_701_000_001).unwrap()),
datetime!(2026-07-21 00:00 UTC),
);
let ResolvedSeries { sparte, .. } = merge(malo(), None, &[], vec![gas])
.unwrap()
.expect("rows were supplied");
assert_eq!(sparte, Sparte::Gas);
}
}