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))
}
async fn resolve(self) -> Result<(Option<ResolvedSeries>, super::QueryResult)> {
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())));
}
}
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((
merge(
self.malo_id,
self.obis_code.as_deref(),
&self.store.config().discriminator_columns(),
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);
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,
}))
}
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 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".to_owned(),
}
}
fn stored(intervals: Vec<MeterInterval>, recorded_at: OffsetDateTime) -> StoredSeries {
let scope = VersionScope::for_interval("99", 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 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 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, hours: i64| {
let mut s = one_interval("1-0:1.8.0");
s.extra.insert(
"bilanzkreis".to_string(),
ScalarValue::Utf8(Some(bk.into())),
);
s.recorded_at += time::Duration::hours(hours);
s
};
assert!(
merge(
malo(),
Some("1-0:1.8.0"),
&[],
vec![of("BK-1", 0), of("BK-2", 1)]
)
.unwrap()
.is_some()
);
}
#[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("99", 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);
}
}