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,
}
}
#[must_use]
pub fn column_eq(mut self, name: &str, value: ScalarValue) -> Self {
self.filters.push((name.to_string(), value));
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(timestamp(from));
}
if let Some(to) = self.to {
conditions.push(format!(r#""{}" < ${}"#, col::FROM, params.len() + 1));
params.push(timestamp(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 LIMIT 1"#, from = col::FROM)
} else {
format!(
r#"ORDER BY "{malo}", "{obis}", "{from}""#,
malo = col::MALO_ID,
obis = col::OBIS_CODE,
from = col::FROM,
)
};
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(), stored)?,
result,
))
}
}
fn timestamp(t: OffsetDateTime) -> ScalarValue {
ScalarValue::TimestampMicrosecond(
Some((t.unix_timestamp_nanos() / 1_000) as i64),
Some("UTC".into()),
)
}
fn merge(
malo_id: MaloId,
obis_code: Option<&str>,
mut stored: Vec<crate::encode::StoredSeries>,
) -> Result<Option<ResolvedSeries>> {
if stored.is_empty() {
return Ok(None);
}
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,
}))
}
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).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,
)
}
#[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)).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);
}
}