use metering::QualityFlag;
use metering::ids::{MaloId, MeloId};
use metering::measurement_series::{MeasurementSource, ProvenanceEntry};
use metering::obis::ObisCode;
use metering::reading::MeterReading;
use metering::resolution::IntervalResolution;
use time::OffsetDateTime;
use std::collections::BTreeMap;
use datafusion::common::ScalarValue;
use crate::encode::StoredReadings;
use crate::encode::schema::col;
use crate::error::{Error, Result};
#[derive(Debug, Clone)]
pub struct ReadingsQuery<'a> {
store: &'a crate::session::MeterStore,
malo_id: MaloId,
melo_id: Option<MeloId>,
obis_code: Option<String>,
from: Option<OffsetDateTime>,
to: Option<OffsetDateTime>,
filters: Vec<(String, ScalarValue)>,
quality: Vec<QualityFlag>,
latest_only: bool,
}
impl<'a> ReadingsQuery<'a> {
pub(crate) fn new(store: &'a crate::session::MeterStore, malo_id: MaloId) -> Self {
Self {
store,
malo_id,
melo_id: None,
obis_code: None,
from: None,
to: None,
filters: Vec::new(),
quality: Vec::new(),
latest_only: false,
}
}
pub fn melo(mut self, melo_id: &str) -> Result<Self> {
self.melo_id = Some(
melo_id
.parse::<MeloId>()
.map_err(|e| Error::encode(col::MELO_ID, format!("{melo_id:?}: {e}")))?,
);
Ok(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 column_eq(mut self, name: &str, value: ScalarValue) -> Result<Self> {
let accepted = filterable_columns(self.store);
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
}
#[must_use]
pub fn range(mut self, from: OffsetDateTime, to: OffsetDateTime) -> Self {
self.from = Some(from);
self.to = Some(to);
self
}
#[must_use]
pub fn since(mut self, from: OffsetDateTime) -> Self {
self.from = Some(from);
self
}
#[must_use]
pub fn until(mut self, to: OffsetDateTime) -> Self {
self.to = Some(to);
self
}
pub async fn collect(self) -> Result<Option<StoredReadings>> {
Ok(self.resolve().await?.0)
}
pub async fn values(self) -> Result<Vec<MeterReading>> {
Ok(self
.collect()
.await?
.map(|r| r.readings)
.unwrap_or_default())
}
pub async fn latest(mut self) -> Result<Option<MeterReading>> {
self.latest_only = true;
Ok(self
.resolve()
.await?
.0
.and_then(|r| r.readings.into_iter().next_back()))
}
pub async fn collect_with_provenance(
self,
) -> Result<(Option<StoredReadings>, super::QueryResult)> {
self.resolve().await
}
pub async fn deliveries(self) -> Result<Vec<StoredReadings>> {
Ok(self.scan().await?.0)
}
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().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, StoredReadings>> {
Ok(self.collect_by_channel_with_provenance().await?.0)
}
pub async fn collect_by_channel_with_provenance(
self,
) -> Result<(BTreeMap<ObisCode, StoredReadings>, super::QueryResult)> {
let malo_id = self.malo_id;
let discriminators = self.store.config().discriminator_columns();
let (stored, result) = self.scan().await?;
Ok((split_by_register(malo_id, &discriminators, stored)?, result))
}
async fn resolve(self) -> Result<(Option<StoredReadings>, super::QueryResult)> {
let malo_id = self.malo_id;
let obis = self.obis_code.clone();
let discriminators = self.store.config().discriminator_columns();
let (stored, result) = self.scan().await?;
Ok((
merge(malo_id, obis.as_deref(), &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()))];
let mut bind = |sql: &str, value: ScalarValue| {
params.push(value);
conditions.push(sql.replace("$?", &format!("${}", params.len())));
};
if let Some(melo) = &self.melo_id {
bind(
&format!(r#""{}" = $?"#, col::MELO_ID),
ScalarValue::Utf8(Some(melo.to_string())),
);
}
if let Some(obis) = &self.obis_code {
bind(
&format!(r#""{}" = $?"#, col::OBIS_CODE),
ScalarValue::Utf8(Some(obis.clone())),
);
}
if let Some(from) = self.from {
bind(
&format!(r#""{}" >= $?"#, col::FROM),
crate::encode::schema::timestamp_scalar(from),
);
}
if let Some(to) = self.to {
bind(
&format!(r#""{}" < $?"#, col::FROM),
crate::encode::schema::timestamp_scalar(to),
);
}
for (name, value) in &self.filters {
bind(&format!(r#""{name}" = $?"#), value.clone());
}
if !self.quality.is_empty() {
let start = params.len() + 1;
let placeholders = (0..self.quality.len())
.map(|i| format!("${}", start + 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<StoredReadings>, super::QueryResult)> {
let (conditions, params) = self.predicate();
let merge_key = self.store.config().merge_key();
let tail = match self.latest_only {
true => format!(
r#"ORDER BY "{at}" DESC, {rest} LIMIT 1"#,
at = col::FROM,
rest = super::series::key_order_without_start(&merge_key),
),
false => format!("ORDER BY {}", super::series::merge_key_order(&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::readings_from_record_batch(batch)?);
}
Ok((stored, result))
}
}
fn filterable_columns(store: &crate::session::MeterStore) -> Vec<String> {
let mut accepted: Vec<String> = store
.config()
.extra_columns()
.iter()
.map(|f| f.name().clone())
.collect();
for column in store.config().discriminator_columns() {
if column != col::MELO_ID && !accepted.contains(&column) {
accepted.push(column);
}
}
accepted
}
fn split_by_register(
malo_id: MaloId,
discriminators: &[String],
stored: Vec<StoredReadings>,
) -> Result<BTreeMap<ObisCode, StoredReadings>> {
let mut grouped: BTreeMap<ObisCode, Vec<StoredReadings>> = BTreeMap::new();
for delivery in stored {
grouped
.entry(delivery.obis_code)
.or_default()
.push(delivery);
}
let mut out = BTreeMap::new();
for (register, deliveries) in grouped {
if let Some(folded) = merge(
malo_id,
Some(®ister.to_string()),
discriminators,
deliveries,
)? {
out.insert(register, folded);
}
}
Ok(out)
}
fn merge(
malo_id: MaloId,
obis_code: Option<&str>,
discriminators: &[String],
mut stored: Vec<StoredReadings>,
) -> Result<Option<StoredReadings>> {
if stored.is_empty() {
return Ok(None);
}
refuse_mixed_registers(&stored, discriminators)?;
stored.sort_by_key(|s| s.recorded_at);
let mut readings: Vec<MeterReading> = Vec::new();
let mut melo_id: Option<MeloId> = None;
let mut cadence: Option<IntervalResolution> = None;
let mut provenance: Vec<ProvenanceEntry> = Vec::new();
let mut source: Option<MeasurementSource> = None;
let mut recorded_at: Option<OffsetDateTime> = None;
let mut sparte = None;
let mut unit = None;
let mut version = None;
let mut extra = BTreeMap::new();
let mut obis = None;
for delivery in stored {
readings.extend(delivery.readings);
if delivery.melo_id.is_some() {
melo_id = delivery.melo_id;
}
if delivery.cadence.is_some() {
cadence = delivery.cadence;
}
provenance.extend(delivery.provenance);
obis = Some(delivery.obis_code);
source = Some(delivery.source);
recorded_at = Some(delivery.recorded_at);
sparte = Some(delivery.sparte);
unit = Some(delivery.unit);
version = Some(delivery.version);
extra = delivery.extra;
}
readings.sort_by_key(|r| r.at);
refuse_repeated_instants(malo_id, &readings)?;
let (Some(source), Some(recorded_at), Some(sparte), Some(unit), Some(obis), Some(version)) =
(source, recorded_at, sparte, unit, obis, version)
else {
return Ok(None);
};
let obis = match obis_code {
Some(code) => code.parse().map_err(|e| {
Error::decode(col::OBIS_CODE, format!("{code:?} is not an OBIS code: {e}"))
})?,
None => obis,
};
let mut out = StoredReadings::new(
malo_id,
obis,
sparte,
readings,
source,
version,
recorded_at,
)
.in_unit(unit);
out.melo_id = melo_id;
out.cadence = cadence;
if !provenance.is_empty() {
out.provenance = provenance;
}
out.extra = extra;
Ok(Some(out))
}
fn refuse_repeated_instants(
malo_id: MaloId,
readings: &[metering::reading::MeterReading],
) -> Result<()> {
for pair in readings.windows(2) {
if pair[0].at != pair[1].at {
continue;
}
return Err(Error::InvariantViolated {
table: malo_id.to_string(),
detail: format!(
"two register readings survived resolution for one meter 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. Differencing them \
would produce a consumption figure that means nothing",
at = pair[0].at,
a = pair[0].value,
b = pair[1].value,
),
});
}
Ok(())
}
fn refuse_mixed_registers(stored: &[StoredReadings], discriminators: &[String]) -> Result<()> {
let key_of = |s: &StoredReadings| -> Result<(String, Vec<(String, String)>)> {
Ok((
s.obis_code.to_string(),
crate::session::store::discriminator_values(
s.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 differs = match first.0 == key.0 {
false => format!("two registers ({} and {})", first.0, key.0),
true => format!("two meters ({} and {})", render(&first.1), render(&key.1)),
};
return Err(Error::config(format!(
"{malo} spans {differs} over this range, and a Zählerstandsgang describes one — \
folding them interleaves two cumulative sequences, so differencing the result \
produces advances that belong to neither meter. Narrow the read with \
.obis(..), .melo(..) or .column_eq(..), or scope the session",
malo = stored[0].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(", "),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::version::{ScopedVersion, Version, VersionScope};
use metering::interval::Sparte;
use metering::reading::MeterReading;
use rust_decimal::Decimal;
use time::macros::datetime;
fn malo() -> MaloId {
"12345678905".parse().expect("a valid MaLo-ID")
}
fn source() -> MeasurementSource {
MeasurementSource::Mscons {
pid: 13_005,
message_ref: None,
sender_mp_id: "9900000000001".parse().expect("a valid Marktpartner-ID"),
}
}
fn one_reading(obis: &str) -> StoredReadings {
let at = datetime!(2026-07-20 00:00 UTC);
let code = obis.parse().expect("a valid OBIS code");
StoredReadings::new(
malo(),
code,
Sparte::Strom,
vec![MeterReading {
at,
value: Decimal::new(1_234, 0),
quality: metering::QualityFlag::Measured,
obis_code: Some(code),
}],
source(),
ScopedVersion::new(
VersionScope::for_interval("9900000000001", at, Sparte::Strom).unwrap(),
Version::new(20_260_720_000_001).unwrap(),
),
at,
)
}
#[test]
fn two_register_values_at_one_instant_are_refused_rather_than_differenced() {
let mut second = one_reading("1-0:1.8.0");
second.readings[0].value = Decimal::new(9_999, 0);
second.recorded_at += time::Duration::hours(1);
let err = merge(malo(), None, &[], vec![one_reading("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("1234") && msg.contains("9999"), "both: {msg}");
}
#[test]
fn a_register_read_at_successive_instants_still_folds() {
let first = one_reading("1-0:1.8.0");
let mut later = one_reading("1-0:1.8.0");
later.readings[0].at += time::Duration::minutes(15);
later.readings[0].value = Decimal::new(1_240, 0);
later.recorded_at += time::Duration::hours(1);
let folded = merge(malo(), None, &[], vec![first, later])
.unwrap()
.expect("one history");
assert_eq!(folded.readings.len(), 2);
}
#[test]
fn a_per_register_read_splits_what_the_fold_refuses() {
let split = split_by_register(
malo(),
&[],
vec![one_reading("1-0:1.8.1"), one_reading("1-0:1.8.2")],
)
.expect("two registers are two histories, not an error");
assert_eq!(split.len(), 2);
for code in ["1-0:1.8.1", "1-0:1.8.2"] {
let register: ObisCode = code.parse().unwrap();
assert_eq!(split[®ister].obis_code, register);
assert_eq!(split[®ister].readings.len(), 1);
}
}
#[test]
fn a_per_register_read_still_refuses_two_meters_on_one_register() {
let a = one_reading("1-0:1.8.0")
.with_melo_id("DE0001112223334445556667778889990".parse().unwrap());
let b = one_reading("1-0:1.8.0")
.with_melo_id("DE0009998887776665554443332221110".parse().unwrap());
let err = split_by_register(malo(), &[col::MELO_ID.to_string()], vec![a, b])
.expect_err("two meters on one register are two readings");
let msg = err.to_string();
assert!(msg.contains("two meters"), "{msg}");
assert!(msg.contains(".melo(..)"), "the fix has to be named: {msg}");
}
#[test]
fn a_per_register_read_of_nothing_is_an_empty_map() {
assert!(
split_by_register(malo(), &[], Vec::new())
.unwrap()
.is_empty()
);
}
}