use std::collections::HashMap;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
projection::Projection,
types::{MarktpartnerCode, MeLo, MessageRef, Pruefidentifikator, Sparte},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-weiterverpflichtung";
pub const AUFTRAG_PID: u32 = 17_002;
pub const ANTWORT_PIDS: (u32, u32) = (19_003, 19_004);
pub const WEITERVERPFLICHTUNG_PIDS: &[u32] = &[AUFTRAG_PID, ANTWORT_PIDS.0, ANTWORT_PIDS.1];
pub const ANTWORT_WINDOW_LABEL: &str = "wim-weiterverpflichtung-antwort";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum WeiterverpflichtungEvent {
AuftragEmpfangen {
melo_id: MeLo,
nb: MarktpartnerCode,
msba: MarktpartnerCode,
verschobenes_zuordnungsende: String,
message_ref: MessageRef,
sparte: Sparte,
},
AntwortGesendet {
pruefidentifikator: Pruefidentifikator,
antwort_code: String,
abweichender_termin: Option<String>,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for WeiterverpflichtungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AuftragEmpfangen { .. } => "WimWeiterverpflichtungAuftragEmpfangen",
Self::AntwortGesendet { .. } => "WimWeiterverpflichtungAntwortGesendet",
Self::Rejected { .. } => "WimWeiterverpflichtungRejected",
Self::DeadlineExpired { .. } => "WimWeiterverpflichtungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct WeiterverpflichtungData {
pub melo_id: MeLo,
pub nb: MarktpartnerCode,
pub msba: MarktpartnerCode,
pub verschobenes_zuordnungsende: String,
pub message_ref: MessageRef,
pub sparte: Sparte,
}
#[must_use]
pub const fn weiterverpflichtung_ebd(sparte: Sparte) -> &'static str {
match sparte {
Sparte::Strom => mako_pruefung::codes::EBD_WEITERVERPFLICHTUNG,
Sparte::Gas => mako_pruefung::codes::EBD_WEITERVERPFLICHTUNG_GAS,
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
#[serde(tag = "status", content = "data")]
pub enum WeiterverpflichtungState {
#[default]
New,
AuftragEmpfangen(WeiterverpflichtungData),
Beantwortet(WeiterverpflichtungData),
Rejected {
reason: String,
},
}
impl WeiterverpflichtungState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::AuftragEmpfangen(_) => "AuftragEmpfangen",
Self::Beantwortet(_) => "Beantwortet",
Self::Rejected { .. } => "Rejected",
}
}
}
impl mako_engine::workflow::OccupiesBusinessKey for WeiterverpflichtungState {
fn occupies_business_key(&self) -> bool {
matches!(self, Self::AuftragEmpfangen(_))
}
}
#[derive(Clone)]
pub enum WeiterverpflichtungCommand {
ReceiveAuftrag {
pid: Pruefidentifikator,
nb: MarktpartnerCode,
msba: MarktpartnerCode,
melo_id: MeLo,
verschobenes_zuordnungsende: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
sparte: Sparte,
},
DispatchAntwort {
antwort_code: String,
abweichender_termin: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for WeiterverpflichtungCommand {}
pub struct WimWeiterverpflichtungWorkflow;
impl Workflow for WimWeiterverpflichtungWorkflow {
type State = WeiterverpflichtungState;
type Event = WeiterverpflichtungEvent;
type Command = WeiterverpflichtungCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(ANTWORT_WINDOW_LABEL, WeiterverpflichtungState::AuftragEmpfangen(_)) => {
Some(WeiterverpflichtungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
WeiterverpflichtungEvent::AuftragEmpfangen {
melo_id,
nb,
msba,
verschobenes_zuordnungsende,
message_ref,
sparte,
} => WeiterverpflichtungState::AuftragEmpfangen(WeiterverpflichtungData {
sparte: *sparte,
melo_id: melo_id.clone(),
nb: nb.clone(),
msba: msba.clone(),
verschobenes_zuordnungsende: verschobenes_zuordnungsende.clone(),
message_ref: message_ref.clone(),
}),
WeiterverpflichtungEvent::AntwortGesendet { .. } => match state {
WeiterverpflichtungState::AuftragEmpfangen(d) => {
WeiterverpflichtungState::Beantwortet(d)
}
other => other,
},
WeiterverpflichtungEvent::Rejected { reason } => WeiterverpflichtungState::Rejected {
reason: reason.clone(),
},
WeiterverpflichtungEvent::DeadlineExpired { label, .. } => match state {
WeiterverpflichtungState::Beantwortet(_)
| WeiterverpflichtungState::Rejected { .. } => state,
_ => WeiterverpflichtungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
WeiterverpflichtungCommand::ReceiveAuftrag {
pid,
nb,
msba,
melo_id,
verschobenes_zuordnungsende,
message_ref,
validation_passed,
validation_errors,
sparte,
} => {
if !matches!(state, WeiterverpflichtungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if pid.as_u32() != AUFTRAG_PID {
return Err(WorkflowError::rejected(format!(
"PID {} is not the Weiterverpflichtung ({AUFTRAG_PID})",
pid.as_u32()
)));
}
if !validation_passed {
return Ok(vec![WeiterverpflichtungEvent::Rejected {
reason: validation_errors.join("; "),
}]
.into());
}
Ok(vec![WeiterverpflichtungEvent::AuftragEmpfangen {
melo_id,
nb,
msba,
verschobenes_zuordnungsende,
message_ref,
sparte,
}]
.into())
}
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code,
abweichender_termin,
} => {
let WeiterverpflichtungState::AuftragEmpfangen(data) = state else {
return Err(WorkflowError::invalid_state(
"AuftragEmpfangen",
state.label(),
));
};
let tree = weiterverpflichtung_ebd(data.sparte);
let code = mako_pruefung::codes::lookup(tree, &antwort_code).ok_or_else(|| {
WorkflowError::rejected(format!(
"Antwortcode {antwort_code:?} is not published in {tree}"
))
})?;
if abweichender_termin.is_none() && matches!(code.code, "Z14" | "Z22") {
return Err(WorkflowError::rejected(format!(
"{tree} {} ({}) requires the corrected Abmeldetermin in DTM DE 2380",
code.code, code.bedeutung
)));
}
let bestaetigt = code.ist_zustimmung().ok_or_else(|| {
WorkflowError::rejected(format!("{} sits off the agreement axis", code.code))
})?;
let antwort_pid = if bestaetigt {
ANTWORT_PIDS.0
} else {
ANTWORT_PIDS.1
};
let mut payload = serde_json::json!({
"pid": antwort_pid,
"sender": data.msba.as_str(),
"receiver": data.nb.as_str(),
"melo": data.melo_id.as_str(),
"antwort_code": code.code,
"antwort_codeliste": code.wire_codeliste().ok_or_else(|| {
WorkflowError::rejected(format!("{tree} {} names no Codeliste", code.code))
})?,
"antwort_tree": tree,
"orig_message_ref": data.message_ref.as_str(),
});
if let Some(ref t) = abweichender_termin {
payload["abmeldetermin"] = serde_json::Value::String(t.clone());
}
Ok(WorkflowOutput::with_outbox(
vec![WeiterverpflichtungEvent::AntwortGesendet {
pruefidentifikator: Pruefidentifikator::new(antwort_pid)
.map_err(WorkflowError::rejected)?,
antwort_code: code.code.to_owned(),
abweichender_termin,
}],
vec![PendingOutbox::new("ORDRSP", data.nb.as_str(), payload).caused_by(0)],
))
}
WeiterverpflichtungCommand::TimeoutExpired { deadline_id, label } => {
if matches!(
state,
WeiterverpflichtungState::Beantwortet(_)
| WeiterverpflichtungState::Rejected { .. }
) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![WeiterverpflichtungEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[derive(Debug, Default)]
pub struct WeiterverpflichtungRecord {
pub status: &'static str,
pub melo_id: Option<String>,
pub antwort_code: Option<String>,
}
#[derive(Debug, Default)]
pub struct WeiterverpflichtungProjection {
pub records: HashMap<String, WeiterverpflichtungRecord>,
}
impl Projection for WeiterverpflichtungProjection {
fn name(&self) -> &'static str {
"WeiterverpflichtungProjection"
}
fn handle_event(&mut self, envelope: &EventEnvelope) {
let Ok(event) = envelope.decode::<WeiterverpflichtungEvent>() else {
return;
};
let record = self
.records
.entry(envelope.stream_id.as_str().to_owned())
.or_default();
match event {
WeiterverpflichtungEvent::AuftragEmpfangen { melo_id, .. } => {
record.status = "AuftragEmpfangen";
record.melo_id = Some(melo_id.as_str().to_owned());
}
WeiterverpflichtungEvent::AntwortGesendet { antwort_code, .. } => {
record.status = "Beantwortet";
record.antwort_code = Some(antwort_code);
}
WeiterverpflichtungEvent::Rejected { .. }
| WeiterverpflichtungEvent::DeadlineExpired { .. } => {
record.status = "Rejected";
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn auftrag_in(sparte: Sparte) -> WeiterverpflichtungCommand {
WeiterverpflichtungCommand::ReceiveAuftrag {
pid: Pruefidentifikator::new(AUFTRAG_PID).expect("17002 is a valid PID"),
nb: MarktpartnerCode::new("9900357000004"),
msba: MarktpartnerCode::new("4012345000023"),
melo_id: MeLo::new("DE0000000001234567890000000000001"),
verschobenes_zuordnungsende: "20260501".to_owned(),
message_ref: MessageRef::new("ORD-17002-1"),
validation_passed: true,
validation_errors: vec![],
sparte,
}
}
fn empfangen_in(sparte: Sparte) -> WeiterverpflichtungState {
let s = WeiterverpflichtungState::default();
let ev = WimWeiterverpflichtungWorkflow::handle(&s, auftrag_in(sparte)).expect("valid");
ev.iter().fold(s, WimWeiterverpflichtungWorkflow::apply)
}
fn empfangen() -> WeiterverpflichtungState {
empfangen_in(Sparte::Strom)
}
#[test]
fn the_sparte_picks_the_codeliste_on_a_shared_pid() {
for (sparte, ebd, codeliste) in [
(Sparte::Strom, "E_0203", "S_0061"),
(Sparte::Gas, "E_2004", "G_0072"),
] {
let out = WimWeiterverpflichtungWorkflow::handle(
&empfangen_in(sparte),
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code: "Z13".to_owned(),
abweichender_termin: None,
},
)
.expect("Z13");
assert_eq!(out.outbox[0].payload["pid"], 19_003, "{sparte}");
assert_eq!(out.outbox[0].payload["antwort_tree"], ebd);
assert_eq!(out.outbox[0].payload["antwort_codeliste"], codeliste);
}
}
#[test]
fn a_plain_agreement_answers_on_19003() {
let out = WimWeiterverpflichtungWorkflow::handle(
&empfangen(),
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code: "Z13".to_owned(),
abweichender_termin: None,
},
)
.expect("Z13");
assert_eq!(&*out.outbox[0].message_type, "ORDRSP");
assert_eq!(out.outbox[0].payload["pid"], 19_003);
assert_eq!(out.outbox[0].payload["antwort_tree"], "E_0203");
assert_eq!(out.outbox[0].payload["antwort_codeliste"], "S_0061");
}
#[test]
fn a_refusal_answers_on_19004_and_names_the_capped_date() {
let out = WimWeiterverpflichtungWorkflow::handle(
&empfangen(),
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code: "Z22".to_owned(),
abweichender_termin: Some("20260401".to_owned()),
},
)
.expect("Z22");
assert_eq!(out.outbox[0].payload["pid"], 19_004);
assert_eq!(out.outbox[0].payload["abmeldetermin"], "20260401");
}
#[test]
fn a_terminaenderung_without_its_date_is_refused() {
let err = WimWeiterverpflichtungWorkflow::handle(
&empfangen(),
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code: "Z14".to_owned(),
abweichender_termin: None,
},
)
.expect_err("Z14 without a date");
assert!(err.to_string().contains("DE 2380"), "{err}");
}
#[test]
fn a_foreign_code_is_refused() {
let err = WimWeiterverpflichtungWorkflow::handle(
&empfangen(),
WeiterverpflichtungCommand::DispatchAntwort {
antwort_code: "Z34".to_owned(),
abweichender_termin: None,
},
)
.expect_err("Z34 is E_0200");
assert!(err.to_string().contains("E_0203"), "{err}");
}
#[test]
fn the_geraeteubernahme_bestellung_is_not_a_weiterverpflichtung() {
let mut cmd = auftrag_in(Sparte::Strom);
if let WeiterverpflichtungCommand::ReceiveAuftrag { ref mut pid, .. } = cmd {
*pid = Pruefidentifikator::new(17_001).expect("valid");
}
assert!(
WimWeiterverpflichtungWorkflow::handle(&WeiterverpflichtungState::default(), cmd)
.is_err()
);
}
#[test]
fn the_answer_window_is_one_werktag() {
use mako_fristen::antwort::{FristShape, antwort_obligation};
let o = antwort_obligation(AUFTRAG_PID).expect("published");
assert_eq!(o.frist, FristShape::WerktageAtCutoff(1));
assert_eq!(o.antwort_pids, ANTWORT_PIDS);
assert_eq!(o.ebd, Some(mako_pruefung::codes::EBD_WEITERVERPFLICHTUNG));
}
}