use rust_decimal::Decimal;
use std::collections::HashMap;
use time::{Duration, OffsetDateTime};
pub const MABIS_SLOT: Duration = Duration::minutes(15);
pub use crate::ids::BilanzierungsgebietId;
pub use crate::ids::BilanzkreisId;
pub use crate::ids::MabisZaehlpunktId;
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct SumInterval {
#[serde(with = "time::serde::rfc3339")]
pub from: OffsetDateTime,
#[serde(with = "time::serde::rfc3339")]
pub to: OffsetDateTime,
pub quantity_kwh: Decimal,
pub malo_count: u32,
pub substituted_count: u32,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Summenzeitreihe {
pub bilanzierungsgebiet_id: BilanzierungsgebietId,
pub mabis_zp_id: MabisZaehlpunktId,
#[serde(with = "time::serde::rfc3339")]
pub period_from: OffsetDateTime,
#[serde(with = "time::serde::rfc3339")]
pub period_to: OffsetDateTime,
#[serde(with = "time::serde::rfc3339")]
pub version: OffsetDateTime,
pub intervals: Vec<SumInterval>,
pub sender_mp_id: String,
pub receiver_mp_id: String,
pub slot_length_minutes: i64,
}
impl Summenzeitreihe {
#[must_use]
pub fn total_kwh(&self) -> Decimal {
self.intervals.iter().map(|i| i.quantity_kwh).sum()
}
#[must_use]
pub fn interval_count(&self) -> usize {
self.intervals.len()
}
#[must_use]
pub fn has_substituted_values(&self) -> bool {
self.intervals.iter().any(|i| i.substituted_count > 0)
}
#[must_use]
pub fn expected_slot_count(&self) -> usize {
let secs = (self.period_to - self.period_from).whole_seconds();
let slot = (self.slot_length_minutes * 60).max(1);
usize::try_from(secs / slot).unwrap_or(0)
}
#[must_use]
pub fn missing_slot_count(&self) -> usize {
self.expected_slot_count()
.saturating_sub(self.intervals.len())
}
#[must_use]
pub fn is_complete(&self) -> bool {
self.missing_slot_count() == 0
}
#[must_use]
pub fn monthly_totals(&self) -> Vec<metering::ResampledBucket> {
use metering::{MeterInterval, QualityFlag, ResampleConfig, resample};
let intervals: Vec<MeterInterval> = self
.intervals
.iter()
.map(|iv| {
let quality = if iv.substituted_count > 0 {
QualityFlag::Substituted
} else {
QualityFlag::Measured
};
MeterInterval {
from: iv.from,
to: iv.to,
value: iv.quantity_kwh,
quality,
obis_code: None,
}
})
.collect();
resample(&intervals, &ResampleConfig::to_monthly())
}
}
pub struct SummenzeitreiheBuilder {
bilanzierungsgebiet_id: BilanzierungsgebietId,
mabis_zp_id: MabisZaehlpunktId,
period_from: OffsetDateTime,
period_to: OffsetDateTime,
version: OffsetDateTime,
sender_mp_id: String,
receiver_mp_id: String,
slot_length: Duration,
slots: HashMap<(i128, i128), (Decimal, u32, u32)>,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error(
"interval {from}..{to} spans {actual_minutes} min, but this Bilanzierungsgebiet settles on a {expected_minutes} min grid"
)]
pub struct SlotResolutionError {
pub from: OffsetDateTime,
pub to: OffsetDateTime,
pub actual_minutes: i64,
pub expected_minutes: i64,
}
impl SummenzeitreiheBuilder {
#[must_use]
#[allow(clippy::too_many_arguments)]
pub fn new(
bilanzierungsgebiet_id: BilanzierungsgebietId,
mabis_zp_id: MabisZaehlpunktId,
period_from: OffsetDateTime,
period_to: OffsetDateTime,
version: OffsetDateTime,
sender_mp_id: impl Into<String>,
receiver_mp_id: impl Into<String>,
slot_length: Duration,
) -> Self {
Self {
bilanzierungsgebiet_id,
mabis_zp_id,
period_from,
period_to,
version,
sender_mp_id: sender_mp_id.into(),
receiver_mp_id: receiver_mp_id.into(),
slot_length,
slots: HashMap::new(),
}
}
#[must_use]
pub fn expected_slot_count(&self) -> usize {
let span = self.period_to - self.period_from;
usize::try_from(span.whole_seconds() / self.slot_length.whole_seconds().max(1)).unwrap_or(0)
}
pub fn add_malo(
&mut self,
intervals: &[metering::MeterInterval],
) -> Result<(), SlotResolutionError> {
for iv in intervals {
let actual = iv.to - iv.from;
if actual != self.slot_length {
return Err(SlotResolutionError {
from: iv.from,
to: iv.to,
actual_minutes: actual.whole_minutes(),
expected_minutes: self.slot_length.whole_minutes(),
});
}
}
for iv in intervals {
let key = (iv.from.unix_timestamp_nanos(), iv.to.unix_timestamp_nanos());
let is_substituted = !matches!(
iv.quality,
metering::QualityFlag::Measured | metering::QualityFlag::Calculated
);
let entry = self.slots.entry(key).or_insert((Decimal::ZERO, 0, 0));
entry.0 += iv.value;
entry.1 += 1;
if is_substituted {
entry.2 += 1;
}
}
Ok(())
}
#[must_use]
pub fn build(self) -> Summenzeitreihe {
let mut intervals: Vec<SumInterval> = self
.slots
.into_iter()
.filter_map(|((from_ns, to_ns), (kwh, count, sub))| {
let from = OffsetDateTime::from_unix_timestamp_nanos(from_ns).ok()?;
let to = OffsetDateTime::from_unix_timestamp_nanos(to_ns).ok()?;
if from >= self.period_to || to <= self.period_from {
return None;
}
Some(SumInterval {
from,
to,
quantity_kwh: kwh,
malo_count: count,
substituted_count: sub,
})
})
.collect();
intervals.sort_by_key(|i| i.from);
Summenzeitreihe {
mabis_zp_id: self.mabis_zp_id,
bilanzierungsgebiet_id: self.bilanzierungsgebiet_id,
period_from: self.period_from,
period_to: self.period_to,
version: self.version,
intervals,
sender_mp_id: self.sender_mp_id,
receiver_mp_id: self.receiver_mp_id,
slot_length_minutes: self.slot_length.whole_minutes(),
}
}
}
#[cfg(test)]
mod tests {
fn bg(code: &str) -> BilanzierungsgebietId {
BilanzierungsgebietId::new(code).expect("a valid Y-type Bilanzierungsgebiet EIC")
}
use super::*;
use metering::{MeterInterval, QualityFlag};
use rust_decimal::dec;
use time::Duration;
use time::macros::datetime;
fn period() -> (OffsetDateTime, OffsetDateTime) {
(
datetime!(2026-06-01 0:00 UTC),
datetime!(2026-07-01 0:00 UTC),
)
}
fn make_iv(
from: OffsetDateTime,
kwh: rust_decimal::Decimal,
quality: QualityFlag,
) -> MeterInterval {
MeterInterval {
from,
to: from + Duration::minutes(15),
value: kwh,
quality,
obis_code: None,
}
}
#[test]
fn empty_builder_produces_empty_series() {
let (from, to) = period();
let builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-TEST-01U"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-03 05:00 UTC),
"9900357000004",
"9900077000006",
MABIS_SLOT,
);
let series = builder.build();
assert_eq!(series.interval_count(), 0);
assert_eq!(series.total_kwh(), Decimal::ZERO);
}
#[test]
fn two_malos_aggregated_correctly() {
let (from, to) = period();
let iv_start = datetime!(2026-06-01 0:00 UTC);
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-TEST-01U"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-03 05:00 UTC),
"9900357000004",
"9900077000006",
MABIS_SLOT,
);
builder
.add_malo(&[make_iv(iv_start, dec!(2.5), QualityFlag::Measured)])
.unwrap();
builder
.add_malo(&[make_iv(iv_start, dec!(3.0), QualityFlag::Substituted)])
.unwrap();
let series = builder.build();
assert_eq!(series.interval_count(), 1);
assert_eq!(series.intervals[0].quantity_kwh, dec!(5.5)); assert_eq!(series.intervals[0].malo_count, 2);
assert_eq!(series.intervals[0].substituted_count, 1);
assert!(series.has_substituted_values());
}
#[test]
fn intervals_outside_period_are_excluded() {
let (from, to) = period();
let before_period = datetime!(2026-05-31 23:45 UTC);
let in_period = datetime!(2026-06-01 0:00 UTC);
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-TEST-01U"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-08 05:00 UTC),
"SENDER",
"BIKO",
MABIS_SLOT,
);
builder
.add_malo(&[
MeterInterval {
from: before_period,
to: from, value: dec!(10.0),
quality: QualityFlag::Measured,
obis_code: None,
},
make_iv(in_period, dec!(5.0), QualityFlag::Measured),
])
.unwrap();
let series = builder.build();
assert_eq!(
series.interval_count(),
1,
"outside interval must be excluded"
);
assert_eq!(series.total_kwh(), dec!(5.0));
}
#[test]
fn estimated_quality_counts_as_substituted() {
let (from, to) = period();
let iv_start = datetime!(2026-06-15 12:00 UTC);
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-QUALIT-5"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-03 05:00 UTC),
"SENDER",
"BIKO",
MABIS_SLOT,
);
builder
.add_malo(&[make_iv(iv_start, dec!(1.0), QualityFlag::Estimated)])
.unwrap();
let series = builder.build();
assert_eq!(
series.intervals[0].substituted_count, 1,
"Estimated must be counted as substituted"
);
}
#[test]
fn monthly_totals_uses_resample() {
let (from, to) = period();
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-MONTHL-Q"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-08 05:00 UTC),
"SENDER",
"BIKO",
MABIS_SLOT,
);
builder
.add_malo(&[make_iv(
datetime!(2026-06-15 10:00 UTC),
dec!(2.0),
QualityFlag::Measured,
)])
.unwrap();
let series = builder.build();
let monthly = series.monthly_totals();
assert_eq!(
monthly.len(),
1,
"all intervals in June should produce one bucket"
);
assert_eq!(monthly[0].total, dec!(2.0));
}
#[test]
fn a_monthly_bucket_is_rejected_rather_than_settled_as_a_slot() {
let (from, to) = period();
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-TEST-01U"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-03 05:00 UTC),
"SENDER",
"BIKO",
MABIS_SLOT,
);
let err = builder
.add_malo(&[MeterInterval {
from,
to,
value: dec!(1234.5),
quality: QualityFlag::Measured,
obis_code: None,
}])
.expect_err("a month-long bucket is not a quarter-hourly slot");
assert_eq!(err.expected_minutes, 15);
assert_eq!(err.actual_minutes, 30 * 24 * 60);
assert_eq!(
builder.build().total_kwh(),
dec!(0),
"a rejected MaLo must contribute nothing"
);
}
#[test]
fn a_partially_covered_period_reports_its_missing_slots() {
let (from, to) = period();
let mut builder = SummenzeitreiheBuilder::new(
bg("11YMAKO-TEST-01U"),
MabisZaehlpunktId::new("DE0004030099000000000000000012345").unwrap(),
from,
to,
datetime!(2026-07-03 05:00 UTC),
"SENDER",
"BIKO",
MABIS_SLOT,
);
builder
.add_malo(&[
make_iv(from, dec!(1.0), QualityFlag::Measured),
make_iv(
from + Duration::minutes(15),
dec!(1.0),
QualityFlag::Measured,
),
])
.unwrap();
let series = builder.build();
assert_eq!(series.expected_slot_count(), 30 * 96);
assert_eq!(series.missing_slot_count(), 30 * 96 - 2);
assert!(!series.is_complete());
}
}
#[cfg(test)]
mod identifier_tests {
use super::*;
const ZP: &str = "DE0004030099000000000000000012345";
fn bg(code: &str) -> BilanzierungsgebietId {
BilanzierungsgebietId::new(code).expect("a valid Y-type Bilanzierungsgebiet EIC")
}
fn series(mabis_zp_id: MabisZaehlpunktId, gebiet: &str) -> Summenzeitreihe {
SummenzeitreiheBuilder::new(
bg(gebiet),
mabis_zp_id,
OffsetDateTime::UNIX_EPOCH,
OffsetDateTime::UNIX_EPOCH + Duration::days(1),
OffsetDateTime::UNIX_EPOCH,
"9900000000001",
"9900000000002",
Duration::minutes(15),
)
.build()
}
#[test]
fn a_malformed_meldepunkt_cannot_reach_a_summenzeitreihe() {
assert!(MabisZaehlpunktId::new("11XSWISSGRIDBGX8").is_err());
assert!(MabisZaehlpunktId::new("").is_err());
}
#[test]
fn the_meldepunkt_territory_swap_is_unrepresentable() {
assert!(BilanzierungsgebietId::new(ZP).is_err());
assert!(MabisZaehlpunktId::new("11YMAKO-TEST-01U").is_err());
}
#[test]
fn a_bilanzkreis_is_refused_where_a_territory_belongs() {
assert!(BilanzierungsgebietId::new("11XSUEDWESTSTRO8").is_err());
assert!(BilanzkreisId::new("11YMAKO-TEST-01U").is_err());
}
#[test]
fn a_well_formed_pair_is_accepted() {
let s = series(MabisZaehlpunktId::new(ZP).unwrap(), "11YMAKO-TEST-01U");
assert_eq!(s.mabis_zp_id.as_str(), ZP);
assert_eq!(s.bilanzierungsgebiet_id.as_ref(), "11YMAKO-TEST-01U");
}
}