use std::collections::HashMap;
use mako_engine::types::Pruefidentifikator;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
projection::Projection,
types::{DeviceId, MarktpartnerCode, MeLo, MessageRef, Sparte},
workflow::{CommandPayload, EventPayload, PendingDeadline, Workflow, WorkflowOutput},
};
use mako_fristen::{
APERAK_STROM_WINDOW_LABEL, HolidayCalendar, aperak_strom_due_at, deadline_at_werktage,
};
use time::OffsetDateTime;
pub const WORKFLOW_NAME: &str = "wim-device-change";
pub const ANTWORT_FRIST_WINDOW_LABEL: &str = "wim-device-change-antwort-frist";
pub const DEVICE_CHANGE_PIDS: &[u32] = &[
55_039, 55_042, 55_051, 55_168, 44_039, 44_042, 44_051, 44_168, ];
#[must_use]
pub fn wim_sparte(pid: u32) -> Option<Sparte> {
match pid {
55_039 | 55_042 | 55_051 | 55_168 => Some(Sparte::Strom),
44_039 | 44_042 | 44_051 | 44_168 => Some(Sparte::Gas),
_ => None,
}
}
#[must_use]
pub fn transaktionsgruende(pid: u32) -> &'static [&'static str] {
let request = if wim_sparte(pid).is_some() {
pid
} else {
match antwort_pid_meaning(pid) {
Some((r, _)) => r,
None => return &[],
}
};
match request {
55_039 => &["E03", "ZR9"],
44_039 => &["E03", "ZR9"],
55_042 => &["E01", "E02", "E03", "ZJ4"],
44_042 => &["E01", "E02", "E03"],
55_051 => &["E01", "E03", "Z33", "ZZB"],
44_051 => &["E01", "E03", "Z33"],
55_168 | 44_168 => &["E01", "E02", "E03"],
_ => &[],
}
}
pub const TRANSAKTIONSGRUND_WECHSEL: &str = "E03";
pub const ENDE_MSB_VOM_NB_PID: u32 = 44_183;
#[must_use]
pub const fn zuordnungs_stunde(sparte: Sparte) -> u8 {
match sparte {
Sparte::Strom => 0,
Sparte::Gas => 6,
}
}
#[must_use]
pub fn antwort_frist_werktage(request_pid: u32) -> Option<u32> {
use mako_fristen::antwort::FristShape;
if !DEVICE_CHANGE_PIDS.contains(&request_pid) {
return None;
}
match mako_fristen::antwort::antwort_obligation(request_pid)?.frist {
FristShape::WerktageAtCutoff(n) => Some(n),
_ => None,
}
}
pub const AUFTRAG_ANTWORT_WINDOW_LABEL: &str = "wim-device-change-auftrag-antwort";
pub const DEVICE_CHANGE_ANTWORT_PIDS: &[(u32, u32, bool)] = &[
(55_040, 55_039, true),
(55_041, 55_039, false),
(55_043, 55_042, true),
(55_044, 55_042, false),
(55_052, 55_051, true),
(55_053, 55_051, false),
(55_169, 55_168, true),
(55_170, 55_168, false),
(44_040, 44_039, true),
(44_041, 44_039, false),
(44_043, 44_042, true),
(44_044, 44_042, false),
(44_052, 44_051, true),
(44_053, 44_051, false),
(44_169, 44_168, true),
];
#[must_use]
pub fn antwort_pid_meaning(pid: u32) -> Option<(u32, bool)> {
DEVICE_CHANGE_ANTWORT_PIDS
.iter()
.find(|(antwort, _, _)| *antwort == pid)
.map(|(_, request, confirmed)| (*request, *confirmed))
}
#[must_use]
pub fn wim_ebd(request_pid: u32) -> Option<&'static str> {
use mako_pruefung::codes as c;
match request_pid {
55_039 => Some(c::EBD_KUENDIGUNG_MSB),
55_042 => Some(c::EBD_ANMELDUNG_MSB),
55_051 => Some(c::EBD_ABMELDUNG_MSB),
55_168 => Some(c::EBD_VERPFLICHTUNGSANFRAGE),
44_039 => Some(c::EBD_KUENDIGUNG_MSB_GAS),
44_042 => Some(c::EBD_ANMELDUNG_MSB_GAS),
44_051 => Some(c::EBD_ABMELDUNG_MSB_GAS),
44_168 => Some(c::EBD_VERPFLICHTUNGSANFRAGE_GAS),
_ => None,
}
}
#[must_use]
pub fn antwort_pid_for(request_pid: u32, bestaetigt: bool) -> Option<u32> {
DEVICE_CHANGE_ANTWORT_PIDS
.iter()
.find(|(_, request, confirmed)| *request == request_pid && *confirmed == bestaetigt)
.map(|(antwort, _, _)| *antwort)
}
pub const IFTSTA_PIDS: &[u32] = &[
21_007, 21_009, 21_010, 21_011, 21_012, 21_013, 21_015, 21_018, 21_036,
];
pub const GESAMTVORGANG_PIDS: &[u32] = &[21_009, 21_010, 21_011, 21_012, 21_013];
pub const GESAMTVORGANG_ERFOLG_PID: u32 = 21_010;
pub const GESAMTVORGANG_SCHEITERN_PID: u32 = 21_009;
pub const ZUORDNUNG_ERFOLG_PID: u32 = 21_012;
pub const ZUORDNUNG_SCHEITERN_PID: u32 = 21_011;
pub const GESAMTVORGANG_AUSGEBLIEBEN_PID: u32 = 21_013;
pub const GESAMTVORGANG_MELDUNG_WINDOW_LABEL: &str = "wim-gesamtvorgang-meldung";
pub const GESAMTVORGANG_AUSBLEIBEN_WINDOW_LABEL: &str = "wim-gesamtvorgang-ausbleiben";
pub const ZUORDNUNG_ANTWORT_WINDOW_LABEL: &str = "wim-gesamtvorgang-antwort";
pub const GESAMTVORGANG_MELDUNG_WT: u32 = 10;
pub const GESAMTVORGANG_AUSBLEIBEN_WT: u32 = 11;
fn parse_yyyymmdd(raw: &str) -> Option<time::Date> {
let digits: String = raw.chars().filter(char::is_ascii_digit).collect();
if digits.len() != 8 {
return None;
}
let year: i32 = digits[0..4].parse().ok()?;
let month = time::Month::try_from(digits[4..6].parse::<u8>().ok()?).ok()?;
let day: u8 = digits[6..8].parse().ok()?;
time::Date::from_calendar_date(year, month, day).ok()
}
fn aperak_deadline(sparte: Sparte, pid: u32, received_at: OffsetDateTime) -> PendingDeadline {
match sparte {
Sparte::Strom => {
PendingDeadline::new(APERAK_STROM_WINDOW_LABEL, aperak_strom_due_at(received_at))
}
Sparte::Gas => {
let (label, due) = mako_fristen::aperak_gas_due_at(pid, received_at);
PendingDeadline::new(label, due)
}
}
}
fn berlin_cutoff(date: time::Date) -> OffsetDateTime {
mako_fristen::berlin_at(
date,
time::Time::from_hms(17, 0, 0).expect("17:00 is valid"),
)
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum DeviceChangeEvent {
Initiated {
melo_id: MeLo,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
device_id: DeviceId,
document_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
#[serde(default)]
vorgangsnummer: Option<String>,
#[serde(default)]
process_date: Option<String>,
#[serde(default)]
transaktionsgrund: Option<String>,
},
ValidationPassed {
message_ref: MessageRef,
},
AperakDispatched {
positive: bool,
reason: Option<String>,
},
AntwortGesendet {
pruefidentifikator: Pruefidentifikator,
bestaetigt: bool,
antwort_code: String,
antwort_ebd: String,
bemerkung: Option<String>,
abweichender_termin: Option<String>,
},
GesamtvorgangGemeldet {
erfolgreich: bool,
zuordnungsbeginn: Option<String>,
outbound: bool,
message_ref: MessageRef,
},
ZuordnungEntschieden {
pruefidentifikator: Pruefidentifikator,
zugeordnet: bool,
zuordnungsbeginn: Option<String>,
outbound: bool,
},
Completed {
device_id: DeviceId,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
InformationEmpfangen {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
},
AuftragGesendet {
melo_id: MeLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
process_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
},
AntwortEmpfangen {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
message_ref: MessageRef,
is_confirmed: bool,
reason: Option<String>,
#[serde(default)]
bestaetigter_termin: Option<String>,
},
}
impl EventPayload for DeviceChangeEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AuftragGesendet { .. } => "WimDeviceChangeAuftragGesendet",
Self::AntwortEmpfangen { .. } => "WimDeviceChangeAntwortEmpfangen",
Self::Initiated { .. } => "WimDeviceChangeInitiated",
Self::ValidationPassed { .. } => "WimDeviceChangeValidationPassed",
Self::AperakDispatched { .. } => "WimDeviceChangeAperakDispatched",
Self::AntwortGesendet { .. } => "WimDeviceChangeAntwortGesendet",
Self::GesamtvorgangGemeldet { .. } => "WimDeviceChangeGesamtvorgangGemeldet",
Self::ZuordnungEntschieden { .. } => "WimDeviceChangeZuordnungEntschieden",
Self::Completed { .. } => "WimDeviceChangeCompleted",
Self::Rejected { .. } => "WimDeviceChangeRejected",
Self::DeadlineExpired { .. } => "WimDeviceChangeDeadlineExpired",
Self::InformationEmpfangen { .. } => "WimDeviceChangeInformationEmpfangen",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DeviceChangeData {
pub melo_id: MeLo,
pub incoming_msb: MarktpartnerCode,
pub grid_operator: MarktpartnerCode,
pub device_id: DeviceId,
pub document_date: String,
pub pruefidentifikator: Pruefidentifikator,
pub message_ref: MessageRef,
#[serde(default)]
pub vorgangsnummer: Option<String>,
#[serde(default)]
pub process_date: Option<String>,
#[serde(default)]
pub bestaetigter_zuordnungsbeginn: Option<String>,
#[serde(default)]
pub transaktionsgrund: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum DeviceChangeState {
#[default]
New,
AuftragGesendet(DeviceChangeData),
AuftragBestaetigt(DeviceChangeData),
Initiated(DeviceChangeData),
ValidationPassed(DeviceChangeData),
AperakSent(DeviceChangeData),
AntwortGesendet(DeviceChangeData),
GesamtvorgangGemeldet(DeviceChangeData),
Completed(DeviceChangeData),
Rejected {
reason: String,
},
}
impl mako_engine::workflow::OccupiesBusinessKey for DeviceChangeState {
fn occupies_business_key(&self) -> bool {
match self {
Self::AuftragGesendet(_)
| Self::AuftragBestaetigt(_)
| Self::Initiated(_)
| Self::ValidationPassed(_)
| Self::AperakSent(_)
| Self::AntwortGesendet(_)
| Self::GesamtvorgangGemeldet(_) => true,
Self::New | Self::Completed(_) | Self::Rejected { .. } => false,
}
}
}
impl DeviceChangeState {
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::AuftragGesendet(_) => "AuftragGesendet",
Self::AuftragBestaetigt(_) => "AuftragBestaetigt",
Self::Initiated(_) => "Initiated",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AperakSent(_) => "AperakSent",
Self::AntwortGesendet(_) => "AntwortGesendet",
Self::GesamtvorgangGemeldet(_) => "GesamtvorgangGemeldet",
Self::Completed(_) => "Completed",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum DeviceChangeCommand {
InitiateDeviceChange {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
process_date: String,
message_ref: MessageRef,
},
ReceiveAntwort {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
message_ref: MessageRef,
reason: Option<String>,
bestaetigter_termin: Option<String>,
},
ReceiveUtilmd {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
device_id: DeviceId,
document_date: String,
message_ref: MessageRef,
vorgangsnummer: Option<String>,
process_date: Option<String>,
transaktionsgrund: Option<String>,
validation_passed: bool,
validation_errors: Vec<String>,
received_at: OffsetDateTime,
},
ReceiveRestOrder {
tx_id: String,
sender_mp_id: MarktpartnerCode,
melo_id: MeLo,
device_category: String,
process_date: String,
},
DispatchAperak {
positive: bool,
reason: Option<String>,
},
DispatchAntwort {
bestaetigt: bool,
antwort_code: String,
bemerkung: Option<String>,
abweichender_termin: Option<String>,
},
MeldeGesamtvorgang {
erfolgreich: bool,
zuordnungsbeginn: Option<String>,
},
ReceiveGesamtvorgang {
pid: Pruefidentifikator,
zuordnungsbeginn: Option<String>,
message_ref: MessageRef,
},
DispatchZuordnung {
zugeordnet: bool,
},
ReceiveZuordnungsantwort {
pid: Pruefidentifikator,
zuordnungsbeginn: Option<String>,
},
Complete {
device_id: DeviceId,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
ReceiveInformation {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
}
impl CommandPayload for DeviceChangeCommand {}
pub struct WimDeviceChangeWorkflow;
impl Workflow for WimDeviceChangeWorkflow {
type State = DeviceChangeState;
type Event = DeviceChangeEvent;
type Command = DeviceChangeCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
ANTWORT_FRIST_WINDOW_LABEL,
DeviceChangeState::Initiated(_) | DeviceChangeState::ValidationPassed(_),
) => Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
(AUFTRAG_ANTWORT_WINDOW_LABEL, DeviceChangeState::AuftragGesendet(_)) => {
Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
(GESAMTVORGANG_MELDUNG_WINDOW_LABEL, DeviceChangeState::AuftragBestaetigt(_)) => {
Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
(GESAMTVORGANG_AUSBLEIBEN_WINDOW_LABEL, DeviceChangeState::AntwortGesendet(_)) => {
Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
(ZUORDNUNG_ANTWORT_WINDOW_LABEL, DeviceChangeState::GesamtvorgangGemeldet(_)) => {
Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
DeviceChangeEvent::AuftragGesendet {
melo_id,
sender,
receiver,
process_date,
message_ref,
pruefidentifikator,
} => DeviceChangeState::AuftragGesendet(DeviceChangeData {
transaktionsgrund: None,
melo_id: melo_id.clone(),
incoming_msb: sender.clone(),
grid_operator: receiver.clone(),
device_id: DeviceId::new(""),
document_date: process_date.clone(),
pruefidentifikator: *pruefidentifikator,
message_ref: message_ref.clone(),
vorgangsnummer: None,
process_date: Some(process_date.clone()),
bestaetigter_zuordnungsbeginn: None,
}),
DeviceChangeEvent::AntwortEmpfangen {
is_confirmed,
reason,
bestaetigter_termin,
..
} => match state {
DeviceChangeState::AuftragGesendet(mut data) => {
if *is_confirmed {
data.bestaetigter_zuordnungsbeginn = bestaetigter_termin
.clone()
.or_else(|| data.process_date.clone());
DeviceChangeState::AuftragBestaetigt(data)
} else {
DeviceChangeState::Rejected {
reason: reason
.clone()
.unwrap_or_else(|| "Auftrag vom Marktpartner abgelehnt".to_owned()),
}
}
}
other => other,
},
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb,
grid_operator,
device_id,
document_date,
message_ref,
pruefidentifikator,
vorgangsnummer,
process_date,
transaktionsgrund,
} => DeviceChangeState::Initiated(DeviceChangeData {
transaktionsgrund: transaktionsgrund.clone(),
melo_id: melo_id.clone(),
incoming_msb: incoming_msb.clone(),
grid_operator: grid_operator.clone(),
device_id: device_id.clone(),
document_date: document_date.clone(),
pruefidentifikator: *pruefidentifikator,
message_ref: message_ref.clone(),
vorgangsnummer: vorgangsnummer.clone(),
process_date: process_date.clone(),
bestaetigter_zuordnungsbeginn: None,
}),
DeviceChangeEvent::ValidationPassed { .. } => {
if let DeviceChangeState::Initiated(data) = state {
DeviceChangeState::ValidationPassed(data)
} else {
state
}
}
DeviceChangeEvent::AperakDispatched { positive, .. } => match state {
DeviceChangeState::ValidationPassed(data) => {
if *positive {
DeviceChangeState::AperakSent(data)
} else {
DeviceChangeState::Rejected {
reason: "negative APERAK".to_owned(),
}
}
}
_ => state,
},
DeviceChangeEvent::AntwortGesendet {
bestaetigt,
abweichender_termin,
..
} => match state {
DeviceChangeState::ValidationPassed(mut data)
| DeviceChangeState::AperakSent(mut data) => {
if *bestaetigt {
data.bestaetigter_zuordnungsbeginn = abweichender_termin
.clone()
.or_else(|| data.process_date.clone());
DeviceChangeState::AntwortGesendet(data)
} else {
DeviceChangeState::Rejected {
reason: "Ablehnung versendet".to_owned(),
}
}
}
other => other,
},
DeviceChangeEvent::GesamtvorgangGemeldet {
erfolgreich,
zuordnungsbeginn,
..
} => match state {
DeviceChangeState::AntwortGesendet(mut data)
| DeviceChangeState::AuftragBestaetigt(mut data) => {
if *erfolgreich {
if let Some(d) = zuordnungsbeginn {
data.bestaetigter_zuordnungsbeginn = Some(d.clone());
}
DeviceChangeState::GesamtvorgangGemeldet(data)
} else {
DeviceChangeState::Rejected {
reason: "Gesamtvorgang gescheitert — der MSBA bleibt zugeordnet"
.to_owned(),
}
}
}
other => other,
},
DeviceChangeEvent::ZuordnungEntschieden {
zugeordnet,
zuordnungsbeginn,
..
} => match state {
DeviceChangeState::GesamtvorgangGemeldet(mut data) => {
if *zugeordnet {
if let Some(d) = zuordnungsbeginn {
data.bestaetigter_zuordnungsbeginn = Some(d.clone());
}
DeviceChangeState::Completed(data)
} else {
DeviceChangeState::Rejected {
reason: "Zuordnung nicht erfolgt — der MSBA bleibt zugeordnet"
.to_owned(),
}
}
}
other => other,
},
DeviceChangeEvent::Completed { device_id } => match state {
DeviceChangeState::AperakSent(mut data)
| DeviceChangeState::AntwortGesendet(mut data)
| DeviceChangeState::GesamtvorgangGemeldet(mut data)
| DeviceChangeState::AuftragBestaetigt(mut data) => {
data.device_id = device_id.clone();
DeviceChangeState::Completed(data)
}
other => other,
},
DeviceChangeEvent::Rejected { reason } => DeviceChangeState::Rejected {
reason: reason.clone(),
},
DeviceChangeEvent::DeadlineExpired { label, .. } => match state {
DeviceChangeState::Completed(_) | DeviceChangeState::Rejected { .. } => state,
_ => DeviceChangeState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
DeviceChangeEvent::InformationEmpfangen { .. } => state,
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
DeviceChangeCommand::InitiateDeviceChange {
pid,
sender,
receiver,
melo_id,
process_date,
message_ref,
} => {
if !matches!(state, DeviceChangeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if !DEVICE_CHANGE_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected a WiM MSB-Wechsel PID (55039, 55042, 55051, 55168), got {pid}",
)));
}
let outbox = PendingOutbox::new(
"UTILMD",
receiver.as_str(),
serde_json::json!({
"direction": "outbound",
"pid": pid.as_u32(),
"sender": sender.as_str(),
"receiver": receiver.as_str(),
"melo": melo_id.as_str(),
"process_date": process_date,
"message_ref": message_ref.as_str(),
}),
)
.caused_by(0);
let event = DeviceChangeEvent::AuftragGesendet {
melo_id,
sender,
receiver,
process_date,
message_ref,
pruefidentifikator: pid,
};
Ok(WorkflowOutput::with_outbox(vec![event], vec![outbox]))
}
DeviceChangeCommand::ReceiveAntwort {
pid,
sender,
message_ref,
reason,
bestaetigter_termin,
} => {
let DeviceChangeState::AuftragGesendet(data) = state else {
return Err(WorkflowError::invalid_state(
"AuftragGesendet",
state.status_str(),
));
};
let Some((request_pid, is_confirmed)) = antwort_pid_meaning(pid.as_u32()) else {
return Err(WorkflowError::rejected(format!(
"PID {pid} is not a WiM MSB-Wechsel Antwort (expected one of \
55040, 55041, 55043, 55044, 55052, 55053, 55169, 55170)",
)));
};
if request_pid != data.pruefidentifikator.as_u32() {
return Err(WorkflowError::rejected(format!(
"Antwort PID {pid} answers request {request_pid}, but this process \
sent {}",
data.pruefidentifikator,
)));
}
let events = vec![DeviceChangeEvent::AntwortEmpfangen {
pruefidentifikator: pid,
sender,
message_ref,
is_confirmed,
reason,
bestaetigter_termin: bestaetigter_termin.clone(),
}];
let beginn = bestaetigter_termin.or_else(|| data.process_date.clone());
match (
is_confirmed,
request_pid,
beginn.as_deref().and_then(parse_yyyymmdd),
) {
(true, 55_042, Some(d)) => Ok(WorkflowOutput::with_outbox_and_deadlines(
events,
vec![],
vec![PendingDeadline::new(
GESAMTVORGANG_MELDUNG_WINDOW_LABEL,
berlin_cutoff(mako_fristen::add_werktage(
d,
GESAMTVORGANG_MELDUNG_WT,
HolidayCalendar::BdewMaKo,
)),
)],
)),
_ => Ok(events.into()),
}
}
DeviceChangeCommand::ReceiveRestOrder {
tx_id,
sender_mp_id,
melo_id,
device_category,
process_date,
} => {
if !matches!(state, DeviceChangeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
let pid = Pruefidentifikator::new(55_042).map_err(|e| {
WorkflowError::rejected(format!(
"constant PID 55042 (WiM Anmeldung MSB) invalid: {e}"
))
})?;
let device_id = DeviceId::new(&*tx_id);
let message_ref = MessageRef::new(&*tx_id);
Ok(vec![
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb: sender_mp_id,
grid_operator: MarktpartnerCode::new(""),
device_id,
document_date: format!("{process_date}|category={device_category}"),
message_ref: message_ref.clone(),
pruefidentifikator: pid,
transaktionsgrund: None,
vorgangsnummer: Some(tx_id.clone()),
process_date: Some(process_date.clone()),
},
DeviceChangeEvent::ValidationPassed { message_ref },
]
.into())
}
DeviceChangeCommand::ReceiveUtilmd {
pid,
transaktionsgrund,
sender,
receiver,
melo_id,
device_id,
document_date,
message_ref,
vorgangsnummer,
process_date,
validation_passed,
validation_errors,
received_at,
} => {
if !matches!(state, DeviceChangeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
let Some(sparte) = wim_sparte(pid.as_u32()) else {
return Err(WorkflowError::rejected(format!(
"PID {} is not a WiM Messstellenbetrieb PID (expected one of {DEVICE_CHANGE_PIDS:?})",
pid.as_u32()
)));
};
let sender_mp_id = sender.clone();
let receiver_gln = receiver.clone();
let mut events = vec![DeviceChangeEvent::Initiated {
melo_id,
incoming_msb: sender,
grid_operator: receiver,
device_id,
document_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
vorgangsnummer,
process_date,
transaktionsgrund,
}];
if validation_passed {
events.push(DeviceChangeEvent::ValidationPassed { message_ref });
let aperak_send_dl = aperak_deadline(sparte, pid.as_u32(), received_at);
let frist_wt = antwort_frist_werktage(pid.as_u32())
.expect("PID guard above restricts this to the MSB-Wechsel family");
let process_dl = PendingDeadline::new(
ANTWORT_FRIST_WINDOW_LABEL,
deadline_at_werktage(received_at, frist_wt, HolidayCalendar::BdewMaKo),
);
Ok(WorkflowOutput::with_outbox_and_deadlines(
events,
vec![],
vec![aperak_send_dl, process_dl],
))
} else {
let reason = validation_errors.join("; ");
events.push(DeviceChangeEvent::Rejected {
reason: reason.clone(),
});
let aperak_send_dl = aperak_deadline(sparte, pid.as_u32(), received_at);
let outbox = vec![
PendingOutbox::new(
"APERAK",
sender_mp_id.as_str(),
serde_json::json!({
"sender": receiver_gln.as_str(),
"receiver": sender_mp_id.as_str(),
"pid": 29001_u32,
"positive": false,
"error_code": mako_engine::erc::codes::Z29,
"reason": reason,
}),
)
.caused_by(0),
];
Ok(WorkflowOutput::with_outbox_and_deadlines(
events,
outbox,
vec![aperak_send_dl],
))
}
}
DeviceChangeCommand::DispatchAperak { positive, reason } => {
let data = match state {
DeviceChangeState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
};
let sparte = wim_sparte(data.pruefidentifikator.as_u32()).ok_or_else(|| {
WorkflowError::rejected(format!(
"PID {} is not a WiM Messstellenbetrieb PID",
data.pruefidentifikator
))
})?;
let suppress_wire = positive
&& !mako_fristen::aperak_hat_anerkennungsmeldung(sparte == Sparte::Gas);
let mut aperak_payload = serde_json::json!({
"sender": data.grid_operator.as_str(),
"pid": data.pruefidentifikator.as_u32(),
"melo": data.melo_id.as_str(),
"positive": positive,
});
if suppress_wire {
aperak_payload["suppress_wire"] = serde_json::Value::Bool(true);
}
aperak_payload["orig_message_ref"] =
serde_json::Value::String(data.message_ref.as_str().to_owned());
if let Some(ref r) = reason {
aperak_payload["reason"] = serde_json::Value::String(r.clone());
}
let outbox_entry =
PendingOutbox::new("APERAK", data.incoming_msb.as_str(), aperak_payload)
.caused_by(0);
Ok(WorkflowOutput::with_outbox(
vec![DeviceChangeEvent::AperakDispatched { positive, reason }],
vec![outbox_entry],
))
}
DeviceChangeCommand::DispatchAntwort {
bestaetigt,
antwort_code,
bemerkung,
abweichender_termin,
} => {
let data = match state {
DeviceChangeState::ValidationPassed(d) | DeviceChangeState::AperakSent(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed or AperakSent",
state.status_str(),
));
}
};
let request_pid = data.pruefidentifikator.as_u32();
let ebd = wim_ebd(request_pid).ok_or_else(|| {
WorkflowError::rejected(format!(
"PID {request_pid} is not answered by a WiM Entscheidungsbaum"
))
})?;
let code = mako_pruefung::codes::lookup(ebd, &antwort_code).ok_or_else(|| {
WorkflowError::rejected(format!(
"Antwortcode {antwort_code:?} is not published in {ebd}"
))
})?;
if code.ist_zustimmung() != Some(bestaetigt) {
return Err(WorkflowError::rejected(format!(
"{ebd} publishes {} in the {} cluster",
code.code,
code.cluster.label()
)));
}
let antwort_pid = antwort_pid_for(request_pid, bestaetigt).ok_or_else(|| {
WorkflowError::rejected(format!("PID {request_pid} has no answer PID"))
})?;
if abweichender_termin.is_none() && matches!(code.code, "Z01" | "Z12" | "Z14") {
return Err(WorkflowError::rejected(format!(
"{ebd} {} ({}) requires `abweichender_termin` — the answer states a \
date change and must name the date",
code.code, code.bedeutung
)));
}
let process_date = abweichender_termin
.clone()
.or_else(|| data.process_date.clone())
.unwrap_or_else(|| data.document_date.clone());
let codeliste = code.wire_codeliste().ok_or_else(|| {
WorkflowError::rejected(format!(
"{ebd} {} names no Codeliste for DE 1131",
code.code
))
})?;
let grund = data
.transaktionsgrund
.clone()
.unwrap_or_else(|| TRANSAKTIONSGRUND_WECHSEL.to_owned());
if !transaktionsgruende(request_pid).contains(&grund.as_str()) {
return Err(WorkflowError::rejected(format!(
"Transaktionsgrund {grund:?} is not published for PID {request_pid} \
(expected one of {:?})",
transaktionsgruende(request_pid)
)));
}
let mut payload = serde_json::json!({
"pid": antwort_pid,
"sender": data.grid_operator.as_str(),
"receiver": data.incoming_msb.as_str(),
"melo": data.melo_id.as_str(),
"process_date": process_date,
"transaktionsgrund": grund,
"antwort_code": code.code,
"antwort_codeliste": codeliste,
"antwort_tree": ebd,
});
if let Some(ref vn) = data.vorgangsnummer {
payload["vorgangsnummer"] = serde_json::Value::String(vn.clone());
}
if let Some(ref text) = bemerkung {
payload["bemerkung"] = serde_json::Value::String(text.clone());
}
let outbox =
PendingOutbox::new("UTILMD", data.incoming_msb.as_str(), payload).caused_by(0);
let deadlines = if bestaetigt && matches!(request_pid, 55_042 | 44_042) {
parse_yyyymmdd(&process_date)
.map(|d| {
vec![PendingDeadline::new(
GESAMTVORGANG_AUSBLEIBEN_WINDOW_LABEL,
berlin_cutoff(mako_fristen::add_werktage(
d,
GESAMTVORGANG_AUSBLEIBEN_WT,
HolidayCalendar::BdewMaKo,
)),
)]
})
.unwrap_or_default()
} else {
vec![]
};
Ok(WorkflowOutput::with_outbox_and_deadlines(
vec![DeviceChangeEvent::AntwortGesendet {
pruefidentifikator: Pruefidentifikator::new(antwort_pid)
.map_err(WorkflowError::rejected)?,
bestaetigt,
antwort_code: code.code.to_owned(),
antwort_ebd: ebd.to_owned(),
bemerkung,
abweichender_termin,
}],
vec![outbox],
deadlines,
))
}
DeviceChangeCommand::MeldeGesamtvorgang {
erfolgreich,
zuordnungsbeginn,
} => {
let DeviceChangeState::AuftragBestaetigt(data) = state else {
return Err(WorkflowError::invalid_state(
"AuftragBestaetigt",
state.status_str(),
));
};
if !matches!(data.pruefidentifikator.as_u32(), 55_042 | 44_042) {
return Err(WorkflowError::rejected(format!(
"the Gesamtvorgang belongs to the Beginn Messstellenbetrieb \
(55042 Strom / 44042 Gas); this process is {}",
data.pruefidentifikator
)));
}
let mut payload = serde_json::json!({
"pid": if erfolgreich {
GESAMTVORGANG_ERFOLG_PID
} else {
GESAMTVORGANG_SCHEITERN_PID
},
"sender": data.incoming_msb.as_str(),
"receiver": data.grid_operator.as_str(),
"melo": data.melo_id.as_str(),
});
if erfolgreich {
let Some(ref beginn) = zuordnungsbeginn else {
return Err(WorkflowError::rejected(
"an erfolgreicher Gesamtvorgang must name the Zuordnungsbeginn \
(SG15 DTM+2380) — it is the date the NB assigns from"
.to_owned(),
));
};
if let (Some(datum), Some(bestaetigt)) = (
parse_yyyymmdd(beginn),
data.bestaetigter_zuordnungsbeginn
.as_deref()
.and_then(parse_yyyymmdd),
) && !mako_fristen::vorlauf::VorlaufShape::Korridor(
mako_fristen::vorlauf::REALISIERUNGSKORRIDOR_WT,
)
.check(datum, bestaetigt, HolidayCalendar::BdewMaKo)
.is_ok()
{
let korridor = mako_fristen::vorlauf::realisierungskorridor(
bestaetigt,
HolidayCalendar::BdewMaKo,
);
return Err(WorkflowError::rejected(format!(
"Übernahmezeitpunkt {datum} liegt außerhalb des \
Realisierungskorridors {}..={} um den bestätigten \
Zuordnungsbeginn {bestaetigt} (WiM Teil 1 Kap. 2.3.2 Nr. 5/6)",
korridor.start(),
korridor.end(),
)));
}
payload["zuordnungsbeginn"] = serde_json::Value::String(beginn.clone());
}
let message_ref = data.message_ref.clone();
Ok(WorkflowOutput::with_outbox(
vec![DeviceChangeEvent::GesamtvorgangGemeldet {
erfolgreich,
zuordnungsbeginn,
outbound: true,
message_ref,
}],
vec![
PendingOutbox::new("IFTSTA", data.grid_operator.as_str(), payload)
.caused_by(0),
],
))
}
DeviceChangeCommand::ReceiveGesamtvorgang {
pid,
zuordnungsbeginn,
message_ref,
} => {
if !matches!(state, DeviceChangeState::AntwortGesendet(_)) {
return Err(WorkflowError::invalid_state(
"AntwortGesendet",
state.status_str(),
));
}
let erfolgreich = match pid.as_u32() {
GESAMTVORGANG_ERFOLG_PID => true,
GESAMTVORGANG_SCHEITERN_PID => false,
other => {
return Err(WorkflowError::rejected(format!(
"PID {other} is not a Gesamtvorgang report (expected \
{GESAMTVORGANG_ERFOLG_PID} erfolgreich or \
{GESAMTVORGANG_SCHEITERN_PID} gescheitert)"
)));
}
};
let events = vec![DeviceChangeEvent::GesamtvorgangGemeldet {
erfolgreich,
zuordnungsbeginn,
outbound: false,
message_ref,
}];
if erfolgreich {
Ok(WorkflowOutput::with_outbox_and_deadlines(
events,
vec![],
vec![PendingDeadline::new(
ZUORDNUNG_ANTWORT_WINDOW_LABEL,
deadline_at_werktage(
OffsetDateTime::now_utc(),
1,
HolidayCalendar::BdewMaKo,
),
)],
))
} else {
Ok(events.into())
}
}
DeviceChangeCommand::DispatchZuordnung { zugeordnet } => {
let DeviceChangeState::GesamtvorgangGemeldet(data) = state else {
return Err(WorkflowError::invalid_state(
"GesamtvorgangGemeldet",
state.status_str(),
));
};
let sparte = wim_sparte(data.pruefidentifikator.as_u32()).ok_or_else(|| {
WorkflowError::rejected(format!(
"PID {} is not a WiM Messstellenbetrieb PID",
data.pruefidentifikator
))
})?;
let pid = if zugeordnet {
ZUORDNUNG_ERFOLG_PID
} else {
ZUORDNUNG_SCHEITERN_PID
};
let mut payload = serde_json::json!({
"pid": pid,
"sender": data.grid_operator.as_str(),
"receiver": data.incoming_msb.as_str(),
"melo": data.melo_id.as_str(),
});
if zugeordnet {
let Some(ref beginn) = data.bestaetigter_zuordnungsbeginn else {
return Err(WorkflowError::rejected(
"the Zuordnung needs the Zuordnungsbeginn the MSBN reported — \
without it there is no date to assign from"
.to_owned(),
));
};
if let (Some(datum), Some(vorlaeufig)) = (
parse_yyyymmdd(beginn),
data.process_date.as_deref().and_then(parse_yyyymmdd),
) && !mako_fristen::vorlauf::VorlaufShape::Korridor(
mako_fristen::vorlauf::REALISIERUNGSKORRIDOR_WT,
)
.check(datum, vorlaeufig, HolidayCalendar::BdewMaKo)
.is_ok()
{
let korridor = mako_fristen::vorlauf::realisierungskorridor(
vorlaeufig,
HolidayCalendar::BdewMaKo,
);
return Err(WorkflowError::rejected(format!(
"der gemeldete Übernahmezeitpunkt {datum} liegt außerhalb des \
Realisierungskorridors {}..={} — die Zuordnung ist abzulehnen \
(IFTSTA {ZUORDNUNG_SCHEITERN_PID})",
korridor.start(),
korridor.end(),
)));
}
payload["zuordnungsbeginn"] = serde_json::Value::String(beginn.clone());
}
let mut outbox = vec![
PendingOutbox::new("IFTSTA", data.incoming_msb.as_str(), payload).caused_by(0),
];
if zugeordnet {
outbox.push(
PendingOutbox::new(
"ProcessCompleted",
"",
serde_json::json!({
"pid": pid,
"melo_id": data.melo_id.as_str(),
"msb_mp_id": data.incoming_msb.as_str(),
"zuordnungsbeginn": data.bestaetigter_zuordnungsbeginn,
"zuordnung_stunde": zuordnungs_stunde(sparte),
"sparte": sparte,
"outcome": "zugeordnet",
}),
)
.caused_by(0),
);
}
Ok(WorkflowOutput::with_outbox(
vec![DeviceChangeEvent::ZuordnungEntschieden {
pruefidentifikator: Pruefidentifikator::new(pid)
.map_err(WorkflowError::rejected)?,
zugeordnet,
zuordnungsbeginn: data.bestaetigter_zuordnungsbeginn.clone(),
outbound: true,
}],
outbox,
))
}
DeviceChangeCommand::ReceiveZuordnungsantwort {
pid,
zuordnungsbeginn,
} => {
if !matches!(state, DeviceChangeState::GesamtvorgangGemeldet(_)) {
return Err(WorkflowError::invalid_state(
"GesamtvorgangGemeldet",
state.status_str(),
));
}
let zugeordnet = match pid.as_u32() {
ZUORDNUNG_ERFOLG_PID => true,
ZUORDNUNG_SCHEITERN_PID | GESAMTVORGANG_AUSGEBLIEBEN_PID => false,
other => {
return Err(WorkflowError::rejected(format!(
"PID {other} is not a Zuordnungsantwort (expected \
{ZUORDNUNG_ERFOLG_PID}, {ZUORDNUNG_SCHEITERN_PID} or \
{GESAMTVORGANG_AUSGEBLIEBEN_PID})"
)));
}
};
Ok(vec![DeviceChangeEvent::ZuordnungEntschieden {
pruefidentifikator: pid,
zugeordnet,
zuordnungsbeginn,
outbound: false,
}]
.into())
}
DeviceChangeCommand::Complete { device_id } => {
if !matches!(
state,
DeviceChangeState::AntwortGesendet(_) | DeviceChangeState::AuftragBestaetigt(_)
) {
return Err(WorkflowError::invalid_state(
"AntwortGesendet or AuftragBestaetigt",
state.status_str(),
));
}
Ok(vec![DeviceChangeEvent::Completed { device_id }].into())
}
DeviceChangeCommand::TimeoutExpired { deadline_id, label } => {
if matches!(
state,
DeviceChangeState::Completed(_) | DeviceChangeState::Rejected { .. }
) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![DeviceChangeEvent::DeadlineExpired { deadline_id, label }].into())
}
DeviceChangeCommand::ReceiveInformation {
pid,
sender,
receiver,
message_ref,
..
} => {
Ok(vec![DeviceChangeEvent::InformationEmpfangen {
pid,
sender,
receiver,
message_ref,
}]
.into())
}
}
}
}
#[derive(Debug)]
pub enum DeviceChangeRecord {
New {
event_count: usize,
},
Active {
status: &'static str,
melo_id: MeLo,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
device_id: DeviceId,
pruefidentifikator: Pruefidentifikator,
event_count: usize,
},
}
impl DeviceChangeRecord {
#[must_use]
pub fn status(&self) -> &'static str {
match self {
Self::New { .. } => "New",
Self::Active { status, .. } => status,
}
}
#[must_use]
pub fn event_count(&self) -> usize {
match self {
Self::New { event_count } | Self::Active { event_count, .. } => *event_count,
}
}
#[must_use]
pub fn active_data(&self) -> Option<DeviceChangeRecordData<'_>> {
match self {
Self::New { .. } => None,
Self::Active {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
..
} => Some(DeviceChangeRecordData {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
}),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct DeviceChangeRecordData<'a> {
pub melo_id: &'a MeLo,
pub incoming_msb: &'a MarktpartnerCode,
pub grid_operator: &'a MarktpartnerCode,
pub device_id: &'a DeviceId,
pub pruefidentifikator: &'a Pruefidentifikator,
}
impl Default for DeviceChangeRecord {
fn default() -> Self {
Self::New { event_count: 0 }
}
}
#[derive(Debug, Default)]
pub struct DeviceChangeProjection {
pub records: HashMap<String, DeviceChangeRecord>,
pub last_seq: u64,
}
impl Projection for DeviceChangeProjection {
fn name(&self) -> &'static str {
"DeviceChangeProjection"
}
fn handle_event(&mut self, envelope: &EventEnvelope) {
self.last_seq = self.last_seq.max(envelope.sequence_number);
let record = self
.records
.entry(envelope.stream_id.as_str().to_owned())
.or_default();
let Ok(event) = envelope.decode::<DeviceChangeEvent>() else {
return;
};
match record {
DeviceChangeRecord::New { event_count } => *event_count += 1,
DeviceChangeRecord::Active { event_count, .. } => *event_count += 1,
}
match event {
DeviceChangeEvent::AuftragGesendet {
melo_id,
sender,
receiver,
pruefidentifikator,
..
} => {
let count = record.event_count();
*record = DeviceChangeRecord::Active {
status: "AuftragGesendet",
melo_id,
incoming_msb: sender,
grid_operator: receiver,
device_id: DeviceId::new(""),
pruefidentifikator,
event_count: count,
};
}
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
..
} => {
let count = record.event_count();
*record = DeviceChangeRecord::Active {
status: "Initiated",
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
event_count: count,
};
}
DeviceChangeEvent::ValidationPassed { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "ValidationPassed";
}
}
DeviceChangeEvent::AntwortEmpfangen { is_confirmed, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if is_confirmed {
"AuftragBestaetigt"
} else {
"Rejected"
};
}
}
DeviceChangeEvent::AperakDispatched { positive, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if positive { "AperakSent" } else { "Rejected" };
}
}
DeviceChangeEvent::AntwortGesendet { bestaetigt, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if bestaetigt {
"AntwortGesendet"
} else {
"Rejected"
};
}
}
DeviceChangeEvent::GesamtvorgangGemeldet { erfolgreich, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if erfolgreich {
"GesamtvorgangGemeldet"
} else {
"Rejected"
};
}
}
DeviceChangeEvent::ZuordnungEntschieden { zugeordnet, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if zugeordnet { "Completed" } else { "Rejected" };
}
}
DeviceChangeEvent::Completed { device_id } => {
if let DeviceChangeRecord::Active {
status,
device_id: d,
..
} = record
{
*status = "Completed";
*d = device_id;
}
}
DeviceChangeEvent::Rejected { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "Rejected";
}
}
DeviceChangeEvent::DeadlineExpired { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "Rejected";
}
}
DeviceChangeEvent::InformationEmpfangen { .. } => {
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_receive_cmd(pid: u32, validation_passed: bool) -> DeviceChangeCommand {
DeviceChangeCommand::ReceiveUtilmd {
transaktionsgrund: Some("E03".to_owned()),
pid: Pruefidentifikator::new(pid).expect("test pid must be in range"),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
melo_id: MeLo::new("DE0000000001234567890000000000001"),
device_id: DeviceId::new("ZHR-12345678"),
document_date: "20250115".to_owned(),
message_ref: MessageRef::new("MSG-WIM-001"),
vorgangsnummer: Some("VG-WIM-001".to_owned()),
process_date: Some("20250201".to_owned()),
validation_passed,
validation_errors: if validation_passed {
vec![]
} else {
vec!["AHB rule violation".to_owned()]
},
received_at: time::OffsetDateTime::now_utc(),
}
}
#[test]
fn happy_path_new_to_completed() {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55042, true))
.expect("should accept valid PID 55042");
assert_eq!(events.len(), 2);
assert!(
matches!(&events[0], DeviceChangeEvent::Initiated { pruefidentifikator, .. } if pruefidentifikator.as_u32() == 55042)
);
assert!(matches!(
&events[1],
DeviceChangeEvent::ValidationPassed { .. }
));
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::ValidationPassed(_)),
"expected ValidationPassed, got {}",
state.status_str()
);
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.expect("dispatch APERAK");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::AperakSent(_)),
"expected AperakSent"
);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: true,
antwort_code: "E15".to_owned(),
bemerkung: None,
abweichender_termin: None,
},
)
.expect("dispatch Antwort");
let wire = &out.outbox[0];
assert_eq!(&*wire.message_type, "UTILMD");
assert_eq!(wire.payload["pid"], 55_043);
assert_eq!(wire.payload["antwort_code"], "E15");
assert_eq!(wire.payload["antwort_tree"], "E_0201");
assert_eq!(wire.payload["vorgangsnummer"], "VG-WIM-001");
assert_eq!(wire.payload["sender"], "9900357000004");
assert_eq!(wire.payload["receiver"], "4012345000023");
let events = out.events.clone();
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::AntwortGesendet(_)),
"expected AntwortGesendet, got {}",
state.status_str()
);
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::Complete {
device_id: DeviceId::new("ZHR-99999999"),
},
)
.expect("complete");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::Completed(d) if d.device_id == DeviceId::new("ZHR-99999999")),
"expected Completed with new device_id",
);
}
#[test]
fn an_aperak_alone_does_not_complete_the_process() {
let state = DeviceChangeState::default();
let events =
WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55_042, true)).expect("valid");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.expect("aperak");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::Complete {
device_id: DeviceId::new("ZHR-1"),
},
)
.expect_err("an unanswered order cannot complete");
assert!(err.to_string().contains("AntwortGesendet"), "{err}");
}
#[test]
fn a_foreign_antwortcode_is_refused_before_the_wire() {
let state = answered_state(55_042);
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: false,
antwort_code: "A02".to_owned(),
bemerkung: None,
abweichender_termin: None,
},
)
.expect_err("A02 is not in E_0201");
assert!(err.to_string().contains("E_0201"), "{err}");
}
#[test]
fn a_terminaenderung_must_name_its_date() {
let state = answered_state(55_039);
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: false,
antwort_code: "Z12".to_owned(),
bemerkung: None,
abweichender_termin: None,
},
)
.expect_err("Z12 without a date");
assert!(err.to_string().contains("abweichender_termin"), "{err}");
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: false,
antwort_code: "Z12".to_owned(),
bemerkung: Some("Vertragsbindung bis 30.06.".to_owned()),
abweichender_termin: Some("20260630".to_owned()),
},
)
.expect("Z12 with a date");
assert_eq!(out.outbox[0].payload["pid"], 55_041);
assert_eq!(out.outbox[0].payload["process_date"], "20260630");
}
fn auftrag_bestaetigt(bestaetigter_termin: Option<&str>) -> DeviceChangeState {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::InitiateDeviceChange {
pid: Pruefidentifikator::new(55_042).expect("valid"),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
melo_id: MeLo::new("DE0000000001234567890000000000001"),
process_date: "20260601".to_owned(),
message_ref: MessageRef::new("MSG-OUT-1"),
},
)
.expect("initiate");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::ReceiveAntwort {
pid: Pruefidentifikator::new(55_043).expect("valid"),
sender: MarktpartnerCode::new("9900357000004"),
message_ref: MessageRef::new("MSG-ANT-1"),
reason: None,
bestaetigter_termin: bestaetigter_termin.map(str::to_owned),
},
)
.expect("antwort");
out.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply)
}
#[test]
fn a_confirmed_anmeldung_opens_the_gesamtvorgang_window() {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::InitiateDeviceChange {
pid: Pruefidentifikator::new(55_042).expect("valid"),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
melo_id: MeLo::new("DE0000000001234567890000000000001"),
process_date: "20260601".to_owned(),
message_ref: MessageRef::new("MSG-OUT-1"),
},
)
.expect("initiate");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::ReceiveAntwort {
pid: Pruefidentifikator::new(55_043).expect("valid"),
sender: MarktpartnerCode::new("9900357000004"),
message_ref: MessageRef::new("MSG-ANT-1"),
reason: None,
bestaetigter_termin: None,
},
)
.expect("antwort");
assert!(
out.deadlines
.iter()
.any(|d| d.label == GESAMTVORGANG_MELDUNG_WINDOW_LABEL),
"a confirmed 55042 must open the 10-Werktage Gesamtvorgang window"
);
}
#[test]
fn a_terminaenderung_moves_the_confirmed_zuordnungsbeginn() {
let DeviceChangeState::AuftragBestaetigt(data) = auftrag_bestaetigt(Some("20260701"))
else {
panic!("expected AuftragBestaetigt");
};
assert_eq!(
data.bestaetigter_zuordnungsbeginn.as_deref(),
Some("20260701")
);
assert_eq!(data.process_date.as_deref(), Some("20260601"));
}
#[test]
fn the_gesamtvorgang_date_must_lie_in_the_realisierungskorridor() {
let state = auftrag_bestaetigt(None);
let inside = mako_fristen::add_werktage(
time::Date::from_calendar_date(2026, time::Month::June, 1).expect("valid"),
9,
HolidayCalendar::BdewMaKo,
);
let outside = mako_fristen::add_werktage(inside, 1, HolidayCalendar::BdewMaKo);
let fmt = |d: time::Date| format!("{:04}{:02}{:02}", d.year(), d.month() as u8, d.day());
assert!(
WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::MeldeGesamtvorgang {
erfolgreich: true,
zuordnungsbeginn: Some(fmt(inside)),
},
)
.is_ok()
);
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::MeldeGesamtvorgang {
erfolgreich: true,
zuordnungsbeginn: Some(fmt(outside)),
},
)
.expect_err("one Werktag past the corridor");
assert!(err.to_string().contains("Realisierungskorridor"), "{err}");
}
#[test]
fn an_erfolgreicher_gesamtvorgang_must_name_its_date() {
let err = WimDeviceChangeWorkflow::handle(
&auftrag_bestaetigt(None),
DeviceChangeCommand::MeldeGesamtvorgang {
erfolgreich: true,
zuordnungsbeginn: None,
},
)
.expect_err("no date");
assert!(err.to_string().contains("DTM+2380"), "{err}");
}
#[test]
fn a_gescheiterter_gesamtvorgang_leaves_the_msba_assigned() {
let state = auftrag_bestaetigt(None);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::MeldeGesamtvorgang {
erfolgreich: false,
zuordnungsbeginn: None,
},
)
.expect("Scheitern");
assert_eq!(out.outbox[0].payload["pid"], GESAMTVORGANG_SCHEITERN_PID);
let state = out
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::Rejected { reason } if reason.contains("MSBA")),
"got {}",
state.status_str()
);
}
#[test]
fn the_nb_answers_the_gesamtvorgang_and_assigns_from_the_reported_date() {
let state = answered_state(55_042);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: true,
antwort_code: "E15".to_owned(),
bemerkung: None,
abweichender_termin: None,
},
)
.expect("Bestätigung");
assert!(
out.deadlines
.iter()
.any(|d| d.label == GESAMTVORGANG_AUSBLEIBEN_WINDOW_LABEL),
"confirming an Anmeldung must arm the 11-Werktage Ausbleiben window"
);
let state = out
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
let gemeldet = mako_fristen::add_werktage(
time::Date::from_calendar_date(2025, time::Month::February, 1).expect("valid"),
5,
HolidayCalendar::BdewMaKo,
);
let gemeldet_str = format!(
"{:04}{:02}{:02}",
gemeldet.year(),
gemeldet.month() as u8,
gemeldet.day()
);
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::ReceiveGesamtvorgang {
pid: Pruefidentifikator::new(GESAMTVORGANG_ERFOLG_PID).expect("valid"),
zuordnungsbeginn: Some(gemeldet_str.clone()),
message_ref: MessageRef::new("MSG-IFT-1"),
},
)
.expect("report");
assert!(
out.deadlines
.iter()
.any(|d| d.label == ZUORDNUNG_ANTWORT_WINDOW_LABEL)
);
let state = out
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
assert!(matches!(
&state,
DeviceChangeState::GesamtvorgangGemeldet(_)
));
let out = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchZuordnung { zugeordnet: true },
)
.expect("Zuordnung");
assert_eq!(&*out.outbox[0].message_type, "IFTSTA");
assert_eq!(out.outbox[0].payload["pid"], ZUORDNUNG_ERFOLG_PID);
assert_eq!(out.outbox[0].payload["zuordnungsbeginn"], gemeldet_str);
let derived = &out.outbox[1];
assert_eq!(&*derived.message_type, "ProcessCompleted");
assert_eq!(derived.payload["pid"], ZUORDNUNG_ERFOLG_PID);
assert_eq!(
derived.payload["melo_id"],
"DE0000000001234567890000000000001"
);
assert_eq!(derived.payload["msb_mp_id"], "4012345000023");
assert_eq!(derived.payload["zuordnungsbeginn"], gemeldet_str);
let state = out
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
assert!(matches!(&state, DeviceChangeState::Completed(d)
if d.bestaetigter_zuordnungsbeginn.as_deref() == Some(gemeldet_str.as_str())));
}
#[test]
fn the_nb_refuses_to_assign_from_a_date_outside_the_korridor() {
let state = answered_state(55_042);
let state = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAntwort {
bestaetigt: true,
antwort_code: "E15".to_owned(),
bemerkung: None,
abweichender_termin: None,
},
)
.expect("Bestätigung")
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
let state = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::ReceiveGesamtvorgang {
pid: Pruefidentifikator::new(GESAMTVORGANG_ERFOLG_PID).expect("valid"),
zuordnungsbeginn: Some("20260210".to_owned()),
message_ref: MessageRef::new("MSG-IFT-2"),
},
)
.expect("the report is recorded even when its date is wrong")
.events
.iter()
.fold(state, WimDeviceChangeWorkflow::apply);
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchZuordnung { zugeordnet: true },
)
.expect_err("out of corridor");
assert!(err.to_string().contains("Realisierungskorridor"), "{err}");
assert!(
WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchZuordnung { zugeordnet: false },
)
.is_ok()
);
}
#[test]
fn the_gesamtvorgang_pids_are_not_in_reading_order() {
assert_eq!(GESAMTVORGANG_SCHEITERN_PID, 21_009);
assert_eq!(GESAMTVORGANG_ERFOLG_PID, 21_010);
assert_eq!(ZUORDNUNG_SCHEITERN_PID, 21_011);
assert_eq!(ZUORDNUNG_ERFOLG_PID, 21_012);
assert_eq!(GESAMTVORGANG_AUSGEBLIEBEN_PID, 21_013);
for pid in GESAMTVORGANG_PIDS {
assert!(
IFTSTA_PIDS.contains(pid),
"{pid} must route to this workflow"
);
}
}
fn answered_state(pid: u32) -> DeviceChangeState {
let state = DeviceChangeState::default();
let events =
WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(pid, true)).expect("valid");
events.iter().fold(state, WimDeviceChangeWorkflow::apply)
}
#[test]
fn wrong_pid_is_rejected() {
let state = DeviceChangeState::default();
let err = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55001, true))
.expect_err("should reject wrong PID");
let msg = err.to_string();
assert!(
msg.contains("55001"),
"error should mention the supplied PID: {msg}"
);
}
#[test]
fn validation_failure_rejects_process() {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55042, false))
.expect("should still produce events");
assert!(matches!(&events[1], DeviceChangeEvent::Rejected { .. }));
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::Rejected { .. }),
"expected Rejected"
);
}
#[test]
fn dispatch_aperak_in_wrong_state_is_rejected() {
let state = DeviceChangeState::default();
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.expect_err("should reject dispatch in wrong state");
assert!(err.to_string().contains("ValidationPassed"), "{err}");
}
}