use crate::convenience::ZeitpunktError;
use crate::generated::v202607::{Lastgang, Zeitreihe, Zeitreihenwert};
#[cfg(feature = "decimal")]
use crate::generated::v202607::{Messwert, Zaehlwerk};
use std::ops::Range;
use time::{Duration, OffsetDateTime};
#[derive(Debug, Clone, PartialEq)]
pub struct PlacedValue<'a> {
pub index: usize,
pub range: Range<OffsetDateTime>,
pub value: &'a Zeitreihenwert,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct UnplacedValue {
pub index: usize,
pub reason: UnplacedReason,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum UnplacedReason {
MissingZeitraum,
NotAnInstantRange,
Malformed(ZeitpunktError),
Reversed,
}
impl std::fmt::Display for UnplacedReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
UnplacedReason::MissingZeitraum => f.write_str("no zeitraum"),
UnplacedReason::NotAnInstantRange => {
f.write_str("zeitraum does not state both a date and a time of day with an offset")
}
UnplacedReason::Malformed(e) => write!(f, "{e}"),
UnplacedReason::Reversed => f.write_str("interval ends before it starts"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
#[non_exhaustive]
pub struct CoverageReport {
pub reference: Option<Range<OffsetDateTime>>,
pub covered: Duration,
pub gaps: Vec<Range<OffsetDateTime>>,
pub overlaps: Vec<Range<OffsetDateTime>>,
pub unplaced: Vec<UnplacedValue>,
pub wrong_length: Vec<usize>,
pub unusable: Vec<usize>,
pub out_of_order: bool,
}
impl CoverageReport {
#[must_use]
pub fn is_complete(&self) -> bool {
self.reference.is_some()
&& self.gaps.is_empty()
&& self.overlaps.is_empty()
&& self.unplaced.is_empty()
&& self.wrong_length.is_empty()
}
#[must_use]
pub fn is_usable(&self) -> bool {
self.is_complete() && self.unusable.is_empty()
}
#[allow(clippy::cast_precision_loss)]
#[must_use]
pub fn coverage_ratio(&self) -> Option<f64> {
let reference = self.reference.as_ref()?;
let span = (reference.end - reference.start).whole_nanoseconds();
if span <= 0 {
return None;
}
Some(self.covered.whole_nanoseconds() as f64 / span as f64)
}
#[must_use]
pub fn missing(&self) -> Duration {
self.gaps
.iter()
.map(|g| g.end - g.start)
.fold(Duration::ZERO, |a, b| a + b)
}
}
pub trait Bo4eTimeSeries {
fn werte(&self) -> &[Zeitreihenwert];
fn expected_interval(&self) -> Option<Duration>;
fn einheit(&self) -> Option<crate::generated::v202607::Mengeneinheit>;
fn placed(&self) -> impl Iterator<Item = PlacedValue<'_>> {
self.werte()
.iter()
.enumerate()
.filter_map(|(index, value)| {
let range = place(value).ok()?;
Some(PlacedValue {
index,
range,
value,
})
})
}
fn span(&self) -> Option<Range<OffsetDateTime>> {
let mut iter = self.placed();
let first = iter.next()?.range;
Some(iter.fold(first, |acc, p| {
acc.start.min(p.range.start)..acc.end.max(p.range.end)
}))
}
fn audit(&self) -> CoverageReport {
audit_inner(self.werte(), self.expected_interval(), None)
}
fn audit_over(&self, reference: Range<OffsetDateTime>) -> CoverageReport {
audit_inner(self.werte(), self.expected_interval(), Some(reference))
}
fn all_values_usable(&self) -> bool {
self.werte()
.iter()
.all(|v| v.status.is_none_or(|s| s.is_usable()))
}
#[cfg(feature = "decimal")]
#[cfg_attr(docsrs, doc(cfg(feature = "decimal")))]
fn sum(&self) -> Option<rust_decimal::Decimal> {
if self.einheit().is_some_and(|u| !u.is_extensive()) {
return None;
}
self.werte()
.iter()
.try_fold(rust_decimal::Decimal::ZERO, |acc, v| {
if v.status.is_some_and(|s| !s.is_usable()) {
return None;
}
acc.checked_add(v.wert?)
})
}
#[cfg(feature = "decimal")]
#[cfg_attr(docsrs, doc(cfg(feature = "decimal")))]
fn integrate(&self) -> Option<rust_decimal::Decimal> {
if self.einheit().is_some_and(|u| u.energy_unit().is_none()) {
return None;
}
self.placed()
.try_fold(rust_decimal::Decimal::ZERO, |acc, p| {
if p.value.status.is_some_and(|s| !s.is_usable()) {
return None;
}
let hours = crate::units::duration_to_hours(p.range.end - p.range.start)?;
acc.checked_add(p.value.wert?.checked_mul(hours)?)
})
}
fn integrated_unit(&self) -> Option<crate::generated::v202607::Mengeneinheit> {
self.einheit()?.energy_unit()
}
}
fn audit_inner(
werte: &[Zeitreihenwert],
expected_interval: Option<Duration>,
reference: Option<Range<OffsetDateTime>>,
) -> CoverageReport {
let mut report = CoverageReport::default();
let mut ranges: Vec<Range<OffsetDateTime>> = Vec::with_capacity(werte.len());
let mut last_start: Option<OffsetDateTime> = None;
for (index, value) in werte.iter().enumerate() {
match place(value) {
Err(reason) => report.unplaced.push(UnplacedValue { index, reason }),
Ok(range) => {
if value.status.is_some_and(|s| !s.is_usable()) {
report.unusable.push(index);
}
if last_start.is_some_and(|prev| range.start < prev) {
report.out_of_order = true;
}
last_start = Some(range.start);
if let Some(expected) = expected_interval {
if range.end - range.start != expected {
report.wrong_length.push(index);
}
}
ranges.push(range);
}
}
}
let reference = reference.or_else(|| {
let first = ranges.first()?.clone();
Some(
ranges
.iter()
.fold(first, |acc, r| acc.start.min(r.start)..acc.end.max(r.end)),
)
});
let Some(reference) = reference else {
return report;
};
report.reference = Some(reference.clone());
ranges.retain(|r| !r.is_empty() && r.start < reference.end && r.end > reference.start);
for r in &mut ranges {
r.start = r.start.max(reference.start);
r.end = r.end.min(reference.end);
}
ranges.sort_unstable_by_key(|r| (r.start, r.end));
let mut cursor = reference.start;
for r in &ranges {
if r.start > cursor {
report.gaps.push(cursor..r.start);
report.covered += r.end - r.start;
} else if r.end > cursor {
if r.start < cursor {
report.overlaps.push(r.start..cursor);
}
report.covered += r.end - cursor;
} else {
report.overlaps.push(r.start..r.end);
}
cursor = cursor.max(r.end);
}
if cursor < reference.end {
report.gaps.push(cursor..reference.end);
}
merge_adjacent(&mut report.overlaps);
report
}
fn place(value: &Zeitreihenwert) -> Result<Range<OffsetDateTime>, UnplacedReason> {
let zeitraum = value
.zeitraum
.as_ref()
.ok_or(UnplacedReason::MissingZeitraum)?;
let range = zeitraum
.as_instant_range()
.ok_or(UnplacedReason::NotAnInstantRange)?
.map_err(UnplacedReason::Malformed)?;
if range.end < range.start {
return Err(UnplacedReason::Reversed);
}
Ok(range)
}
fn merge_adjacent(ranges: &mut Vec<Range<OffsetDateTime>>) {
if ranges.len() < 2 {
return;
}
ranges.sort_unstable_by_key(|r| (r.start, r.end));
let mut merged: Vec<Range<OffsetDateTime>> = Vec::with_capacity(ranges.len());
for r in ranges.iter().cloned() {
match merged.last_mut() {
Some(last) if r.start <= last.end => last.end = last.end.max(r.end),
_ => merged.push(r),
}
}
*ranges = merged;
}
impl Bo4eTimeSeries for Lastgang {
fn werte(&self) -> &[Zeitreihenwert] {
self.werte.as_deref().unwrap_or_default()
}
fn expected_interval(&self) -> Option<Duration> {
#[cfg(feature = "decimal")]
{
self.zeit_intervall_laenge.as_duration()
}
#[cfg(not(feature = "decimal"))]
{
None
}
}
fn einheit(&self) -> Option<crate::generated::v202607::Mengeneinheit> {
self.messgroesse
}
}
impl Bo4eTimeSeries for Zeitreihe {
fn werte(&self) -> &[Zeitreihenwert] {
self.werte.as_deref().unwrap_or_default()
}
fn expected_interval(&self) -> Option<Duration> {
None
}
fn einheit(&self) -> Option<crate::generated::v202607::Mengeneinheit> {
self.einheit
}
}
#[cfg(feature = "decimal")]
#[cfg_attr(docsrs, doc(cfg(feature = "decimal")))]
#[derive(Debug, Clone, PartialEq)]
pub struct Reading<'a> {
pub index: usize,
pub at: OffsetDateTime,
pub value: rust_decimal::Decimal,
pub source: &'a Messwert,
}
#[cfg(feature = "decimal")]
#[cfg_attr(docsrs, doc(cfg(feature = "decimal")))]
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum ConsumptionError {
#[error("a consumption needs two readings, and this register has {count}")]
TooFewReadings {
count: usize,
},
#[error(
"reading {index} fell from {from} to {to}, and the register states no \
vorkommastelle, so a wrap-around cannot be told from a fault"
)]
DecreasedWithoutRegisterWidth {
index: usize,
from: rust_decimal::Decimal,
to: rust_decimal::Decimal,
},
#[error("reading {index} is marked Z78_GERAETEWECHSEL; split the series there")]
MeterExchange {
index: usize,
},
#[error("reading {index} is in a unit that does not convert to the register's")]
IncompatibleUnit {
index: usize,
},
#[error("the consumption arithmetic overflowed")]
Overflow,
}
#[cfg(feature = "decimal")]
impl Zaehlwerk {
#[must_use]
pub fn readings(&self) -> Vec<Reading<'_>> {
let target = self.einheit;
let mut out: Vec<Reading<'_>> = self
.messwerte
.as_deref()
.unwrap_or_default()
.iter()
.enumerate()
.filter_map(|(index, m)| {
if m.messwertstatus.is_some_and(|s| !s.is_usable()) {
return None;
}
let value = reading_value(m, target)?;
Some(Reading {
index,
at: m.zeitpunkt?,
value,
source: m,
})
})
.collect();
out.sort_by_key(|r| r.at);
out
}
#[must_use]
pub fn register_capacity(&self) -> Option<rust_decimal::Decimal> {
let digits = u32::try_from(self.vorkommastelle?).ok()?;
if digits > 28 {
return None;
}
#[allow(clippy::unnecessary_fallible_conversions)]
rust_decimal::Decimal::try_from(10u128.checked_pow(digits)?).ok()
}
pub fn consumption_between(
&self,
from: rust_decimal::Decimal,
to: rust_decimal::Decimal,
) -> Result<rust_decimal::Decimal, ConsumptionError> {
self.consumption_at(0, from, to)
}
fn consumption_at(
&self,
index: usize,
from: rust_decimal::Decimal,
to: rust_decimal::Decimal,
) -> Result<rust_decimal::Decimal, ConsumptionError> {
use rust_decimal::Decimal;
let mut delta = to.checked_sub(from).ok_or(ConsumptionError::Overflow)?;
if delta < Decimal::ZERO {
let capacity = self
.register_capacity()
.ok_or(ConsumptionError::DecreasedWithoutRegisterWidth { index, from, to })?;
delta = delta
.checked_add(capacity)
.ok_or(ConsumptionError::Overflow)?;
if delta < Decimal::ZERO {
return Err(ConsumptionError::DecreasedWithoutRegisterWidth { index, from, to });
}
}
delta
.checked_mul(self.wandlerfaktor.unwrap_or(Decimal::ONE))
.ok_or(ConsumptionError::Overflow)
}
pub fn total_consumption(&self) -> Result<rust_decimal::Decimal, ConsumptionError> {
if let Some(index) = self.unconvertible_reading() {
return Err(ConsumptionError::IncompatibleUnit { index });
}
let readings = self.readings();
if readings.len() < 2 {
return Err(ConsumptionError::TooFewReadings {
count: readings.len(),
});
}
if let Some(r) = readings.iter().find(|r| {
r.source.messwertstatuszusatz
== Some(crate::generated::v202607::Messwertstatuszusatz::Z78Geraetewechsel)
}) {
return Err(ConsumptionError::MeterExchange { index: r.index });
}
readings
.windows(2)
.try_fold(rust_decimal::Decimal::ZERO, |acc, pair| {
let step = self.consumption_at(pair[1].index, pair[0].value, pair[1].value)?;
acc.checked_add(step).ok_or(ConsumptionError::Overflow)
})
}
fn unconvertible_reading(&self) -> Option<usize> {
let target = self.einheit?;
self.messwerte
.as_deref()
.unwrap_or_default()
.iter()
.enumerate()
.find(|(_, m)| {
m.messwertstatus.is_none_or(|s| s.is_usable())
&& m.zeitpunkt.is_some()
&& m.wert.is_some()
&& reading_value(m, Some(target)).is_none()
})
.map(|(index, _)| index)
}
}
#[cfg(feature = "decimal")]
fn reading_value(
m: &Messwert,
target: Option<crate::generated::v202607::Mengeneinheit>,
) -> Option<rust_decimal::Decimal> {
let menge = m.wert.as_ref()?;
match (target, menge.einheit) {
(Some(t), Some(u)) if u != t => menge.convert_to(t)?.wert,
_ => menge.wert,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::generated::v202607::{Menge, Mengeneinheit, Zeitraum};
use time::macros::datetime;
#[cfg(feature = "decimal")]
fn value(from: OffsetDateTime, minutes: i64, wert: i64) -> Zeitreihenwert {
Zeitreihenwert {
wert: Some(rust_decimal::Decimal::from(wert)),
zeitraum: Some(Zeitraum::from_instants(
from,
from + Duration::minutes(minutes),
)),
..Default::default()
}
}
#[cfg(not(feature = "decimal"))]
fn value(from: OffsetDateTime, minutes: i64, wert: i64) -> Zeitreihenwert {
Zeitreihenwert {
wert: Some(wert.to_string()),
zeitraum: Some(Zeitraum::from_instants(
from,
from + Duration::minutes(minutes),
)),
..Default::default()
}
}
fn lastgang(werte: Vec<Zeitreihenwert>) -> Lastgang {
Lastgang {
messgroesse: Some(Mengeneinheit::Kw),
werte: Some(werte),
..Lastgang::new(quarter_hour_interval())
}
}
#[cfg(feature = "decimal")]
fn quarter_hour_interval() -> Menge {
Menge {
wert: Some(rust_decimal::Decimal::from(15)),
einheit: Some(Mengeneinheit::Minute),
..Default::default()
}
}
#[cfg(not(feature = "decimal"))]
fn quarter_hour_interval() -> Menge {
Menge {
wert: Some("15".to_owned()),
einheit: Some(Mengeneinheit::Minute),
..Default::default()
}
}
const T0: OffsetDateTime = datetime!(2026-01-01 00:00 +01:00);
#[test]
fn a_contiguous_series_is_complete() {
let lg = lastgang(
(0..4)
.map(|i| value(T0 + Duration::minutes(15 * i), 15, 100))
.collect(),
);
let report = lg.audit();
assert!(report.is_complete(), "{report:?}");
assert_eq!(report.covered, Duration::HOUR);
assert_eq!(report.coverage_ratio(), Some(1.0));
assert!(!report.out_of_order);
assert_eq!(report.span_or_none(), lg.span());
}
#[test]
fn a_hole_becomes_a_gap() {
let lg = lastgang(vec![
value(T0, 15, 100),
value(T0 + Duration::minutes(30), 15, 100),
]);
let report = lg.audit();
assert!(!report.is_complete());
assert_eq!(
report.gaps,
[T0 + Duration::minutes(15)..T0 + Duration::minutes(30)]
);
assert_eq!(report.missing(), Duration::minutes(15));
assert_eq!(report.covered, Duration::minutes(30));
}
#[test]
fn a_duplicate_reading_becomes_an_overlap_and_is_counted_once() {
let lg = lastgang(vec![value(T0, 15, 100), value(T0, 15, 100)]);
let report = lg.audit();
assert!(!report.is_complete());
assert_eq!(report.overlaps, [T0..T0 + Duration::minutes(15)]);
assert_eq!(report.covered, Duration::minutes(15));
assert!(report.gaps.is_empty());
}
#[test]
fn a_partial_overlap_reports_only_the_shared_stretch() {
let lg = lastgang(vec![
value(T0, 15, 100),
value(T0 + Duration::minutes(10), 15, 100),
]);
let report = lg.audit();
assert_eq!(
report.overlaps,
[T0 + Duration::minutes(10)..T0 + Duration::minutes(15)]
);
assert_eq!(report.covered, Duration::minutes(25));
}
#[test]
fn an_entry_of_the_wrong_length_is_named_by_index() {
let lg = lastgang(vec![
value(T0, 15, 100),
value(T0 + Duration::minutes(15), 30, 100),
]);
let report = lg.audit();
assert_eq!(report.wrong_length, [1]);
assert!(report.gaps.is_empty());
}
#[test]
fn a_zeitreihe_declares_no_interval_so_none_is_wrong() {
let zr = Zeitreihe {
einheit: Some(Mengeneinheit::Kwh),
werte: Some(vec![
value(T0, 15, 100),
value(T0 + Duration::minutes(15), 45, 100),
]),
..Default::default()
};
let report = zr.audit();
assert!(report.wrong_length.is_empty());
assert!(report.is_complete(), "{report:?}");
}
#[test]
fn an_unplaceable_entry_is_reported_with_its_reason() {
let lg = lastgang(vec![
value(T0, 15, 100),
Zeitreihenwert::default(),
Zeitreihenwert {
zeitraum: Some(Zeitraum {
startdatum: Some(time::macros::date!(2026 - 01 - 01)),
..Default::default()
}),
..Default::default()
},
]);
let report = lg.audit();
assert_eq!(
report.unplaced,
[
UnplacedValue {
index: 1,
reason: UnplacedReason::MissingZeitraum
},
UnplacedValue {
index: 2,
reason: UnplacedReason::NotAnInstantRange
},
]
);
assert!(!report.is_complete());
}
#[test]
fn unsorted_entries_are_flagged_but_still_measured() {
let lg = lastgang(vec![
value(T0 + Duration::minutes(15), 15, 100),
value(T0, 15, 100),
]);
let report = lg.audit();
assert!(report.out_of_order);
assert!(report.gaps.is_empty());
assert!(report.overlaps.is_empty());
assert_eq!(report.covered, Duration::minutes(30));
}
#[test]
fn auditing_over_a_reference_finds_a_missing_tail() {
let lg = lastgang(vec![value(T0, 15, 100)]);
assert!(lg.audit().is_complete());
let report = lg.audit_over(T0..T0 + Duration::HOUR);
assert_eq!(
report.gaps,
[T0 + Duration::minutes(15)..T0 + Duration::HOUR]
);
assert_eq!(report.coverage_ratio(), Some(0.25));
}
#[test]
fn entries_outside_the_reference_are_clipped_not_counted() {
let lg = lastgang(vec![
value(T0 - Duration::HOUR, 15, 100),
value(T0, 15, 100),
]);
let report = lg.audit_over(T0..T0 + Duration::minutes(15));
assert!(report.is_complete(), "{report:?}");
assert_eq!(report.covered, Duration::minutes(15));
}
#[test]
fn an_empty_series_has_no_reference() {
let report = lastgang(vec![]).audit();
assert_eq!(report.reference, None);
assert_eq!(report.coverage_ratio(), None);
assert!(!report.is_complete());
}
#[cfg(feature = "decimal")]
#[test]
fn power_integrates_but_does_not_sum() {
use rust_decimal::Decimal;
let lg = lastgang(vec![
value(T0, 15, 400),
value(T0 + Duration::minutes(15), 15, 480),
]);
assert_eq!(lg.integrate(), Some(Decimal::from(220)));
assert_eq!(lg.sum(), None, "kW is not an extensive quantity");
}
#[cfg(feature = "decimal")]
#[test]
fn a_stated_unit_admits_exactly_one_aggregate() {
let power = lastgang(vec![value(T0, 15, 400)]);
assert!(power.sum().is_none() && power.integrate().is_some());
assert_eq!(power.integrated_unit(), Some(Mengeneinheit::Kwh));
let energy = Zeitreihe {
einheit: Some(Mengeneinheit::Kwh),
werte: Some(vec![value(T0, 15, 100)]),
..Default::default()
};
assert!(energy.sum().is_some() && energy.integrate().is_none());
assert_eq!(energy.integrated_unit(), None);
let unstated = Zeitreihe {
werte: Some(vec![value(T0, 15, 100)]),
..Default::default()
};
assert!(unstated.sum().is_some() && unstated.integrate().is_some());
assert_eq!(unstated.integrated_unit(), None);
}
#[cfg(feature = "decimal")]
#[test]
fn energy_sums() {
use rust_decimal::Decimal;
let zr = Zeitreihe {
einheit: Some(Mengeneinheit::Kwh),
werte: Some(vec![
value(T0, 15, 100),
value(T0 + Duration::minutes(15), 15, 120),
]),
..Default::default()
};
assert_eq!(zr.sum(), Some(Decimal::from(220)));
}
#[cfg(feature = "decimal")]
#[test]
fn a_missing_value_makes_the_sum_unavailable() {
let zr = Zeitreihe {
einheit: Some(Mengeneinheit::Kwh),
werte: Some(vec![value(T0, 15, 100), Zeitreihenwert::default()]),
..Default::default()
};
assert_eq!(zr.sum(), None);
}
#[test]
fn a_declared_absence_keeps_the_timeline_but_not_the_data() {
use crate::generated::v202607::Messwertstatus;
let mut werte = vec![value(T0, 15, 100), value(T0 + Duration::minutes(15), 15, 0)];
werte[0].status = Some(Messwertstatus::Abgelesen);
werte[1].status = Some(Messwertstatus::Fehlt);
let lg = lastgang(werte);
let report = lg.audit();
assert!(report.is_complete(), "the timeline has no hole: {report:?}");
assert_eq!(report.unusable, [1]);
assert!(!report.is_usable());
#[cfg(feature = "decimal")]
assert_eq!(lg.integrate(), None);
}
#[test]
fn a_substitute_reading_is_usable() {
use crate::generated::v202607::Messwertstatus;
let mut werte = vec![value(T0, 15, 400)];
werte[0].status = Some(Messwertstatus::Ersatzwert);
let lg = lastgang(werte);
let report = lg.audit();
assert!(report.unusable.is_empty());
assert!(report.is_usable());
assert!(!Messwertstatus::Ersatzwert.is_measured());
#[cfg(feature = "decimal")]
assert_eq!(lg.integrate(), Some(rust_decimal::Decimal::from(100)));
}
#[test]
fn all_values_usable_agrees_with_the_audit() {
use crate::generated::v202607::Messwertstatus;
let clean = lastgang(vec![value(T0, 15, 100)]);
assert!(clean.all_values_usable());
assert!(clean.audit().unusable.is_empty());
let mut werte = vec![value(T0, 15, 100)];
werte[0].status = Some(Messwertstatus::NichtVerwendbar);
let dirty = lastgang(werte);
assert!(!dirty.all_values_usable());
assert_eq!(dirty.audit().unusable, [0]);
let mut werte = vec![Zeitreihenwert::default()];
werte[0].status = Some(Messwertstatus::Fehlt);
assert!(!lastgang(werte).all_values_usable());
}
#[test]
fn expected_interval_reads_the_declared_menge() {
let lg = lastgang(vec![]);
#[cfg(feature = "decimal")]
assert_eq!(lg.expected_interval(), Some(Duration::minutes(15)));
#[cfg(not(feature = "decimal"))]
assert_eq!(lg.expected_interval(), None);
}
impl CoverageReport {
fn span_or_none(&self) -> Option<Range<OffsetDateTime>> {
self.reference.clone()
}
}
}