use std::collections::HashMap;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
projection::Projection,
types::{MarktpartnerCode, MeLo, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-ersteinbau";
pub const VORABINFORMATION_PID: u32 = mako_fristen::antwort::ERSTEINBAU_VORABINFORMATION_PID;
pub const ZUSTIMMUNG_PID: u32 = 21_030;
pub const ABLEHNUNG_PID: u32 = 21_031;
pub const SCHEITERN_PID: u32 = 21_027;
pub const KEIN_ERSTEINBAU_LF_PID: u32 = 21_025;
pub const ERSTEINBAU_PIDS: &[u32] = &[VORABINFORMATION_PID, ZUSTIMMUNG_PID, ABLEHNUNG_PID];
pub const ANTWORT_WINDOW_LABEL: &str = "wim-ersteinbau-antwort-frist";
pub const ERSTEINBAU_EBD: &str = mako_pruefung::codes::EBD_ERSTEINBAU_IMS;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum ErsteinbauEvent {
VorabinformationEmpfangen {
melo_id: MeLo,
gmsb: MarktpartnerCode,
wmsb: MarktpartnerCode,
umstellungszeitpunkt: String,
message_ref: MessageRef,
},
VorabinformationGesendet {
melo_id: MeLo,
gmsb: MarktpartnerCode,
wmsb: MarktpartnerCode,
umstellungszeitpunkt: String,
message_ref: MessageRef,
},
AntwortGesendet {
pruefidentifikator: Pruefidentifikator,
antwort_code: String,
},
ScheiternGemeldet {
reason: String,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for ErsteinbauEvent {
fn event_type(&self) -> &'static str {
match self {
Self::VorabinformationEmpfangen { .. } => "WimErsteinbauVorabinformationEmpfangen",
Self::VorabinformationGesendet { .. } => "WimErsteinbauVorabinformationGesendet",
Self::AntwortGesendet { .. } => "WimErsteinbauAntwortGesendet",
Self::ScheiternGemeldet { .. } => "WimErsteinbauScheiternGemeldet",
Self::Rejected { .. } => "WimErsteinbauRejected",
Self::DeadlineExpired { .. } => "WimErsteinbauDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ErsteinbauData {
pub melo_id: MeLo,
pub gmsb: MarktpartnerCode,
pub wmsb: MarktpartnerCode,
pub umstellungszeitpunkt: String,
pub message_ref: MessageRef,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
#[serde(tag = "status", content = "data")]
pub enum ErsteinbauState {
#[default]
New,
Angekuendigt(ErsteinbauData),
Zugestimmt(ErsteinbauData),
Abgelehnt {
antwort_code: Option<String>,
reason: String,
},
Gescheitert {
reason: String,
},
Rejected {
reason: String,
},
}
impl ErsteinbauState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Angekuendigt(_) => "Angekuendigt",
Self::Zugestimmt(_) => "Zugestimmt",
Self::Abgelehnt { .. } => "Abgelehnt",
Self::Gescheitert { .. } => "Gescheitert",
Self::Rejected { .. } => "Rejected",
}
}
#[must_use]
pub const fn rollout_freigegeben(&self) -> bool {
matches!(self, Self::Zugestimmt(_))
}
}
impl mako_engine::workflow::OccupiesBusinessKey for ErsteinbauState {
fn occupies_business_key(&self) -> bool {
matches!(self, Self::Angekuendigt(_))
}
}
#[derive(Clone)]
pub enum ErsteinbauCommand {
ReceiveVorabinformation {
pid: Pruefidentifikator,
gmsb: MarktpartnerCode,
wmsb: MarktpartnerCode,
melo_id: MeLo,
umstellungszeitpunkt: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
SendVorabinformation {
gmsb: MarktpartnerCode,
wmsb: MarktpartnerCode,
melo_id: MeLo,
umstellungszeitpunkt: String,
message_ref: MessageRef,
},
DispatchAntwort {
antwort_code: String,
},
MeldeScheitern {
reason: String,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for ErsteinbauCommand {}
pub struct WimErsteinbauWorkflow;
impl Workflow for WimErsteinbauWorkflow {
type State = ErsteinbauState;
type Event = ErsteinbauEvent;
type Command = ErsteinbauCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(ANTWORT_WINDOW_LABEL, ErsteinbauState::Angekuendigt(_)) => {
Some(ErsteinbauCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
ErsteinbauEvent::VorabinformationEmpfangen {
melo_id,
gmsb,
wmsb,
umstellungszeitpunkt,
message_ref,
}
| ErsteinbauEvent::VorabinformationGesendet {
melo_id,
wmsb,
umstellungszeitpunkt,
message_ref,
gmsb,
} => ErsteinbauState::Angekuendigt(ErsteinbauData {
melo_id: melo_id.clone(),
gmsb: gmsb.clone(),
wmsb: wmsb.clone(),
umstellungszeitpunkt: umstellungszeitpunkt.clone(),
message_ref: message_ref.clone(),
}),
ErsteinbauEvent::AntwortGesendet {
pruefidentifikator,
antwort_code,
} => match state {
ErsteinbauState::Angekuendigt(d) => {
if pruefidentifikator.as_u32() == ZUSTIMMUNG_PID {
ErsteinbauState::Zugestimmt(d)
} else {
ErsteinbauState::Abgelehnt {
antwort_code: Some(antwort_code.clone()),
reason: format!(
"E_0233 {antwort_code} — der Ersteinbau eines iMS durch den gMSB \
darf nicht erfolgen"
),
}
}
}
other => other,
},
ErsteinbauEvent::ScheiternGemeldet { reason } => ErsteinbauState::Gescheitert {
reason: reason.clone(),
},
ErsteinbauEvent::Rejected { reason } => ErsteinbauState::Rejected {
reason: reason.clone(),
},
ErsteinbauEvent::DeadlineExpired { label, .. } => match state {
ErsteinbauState::Angekuendigt(_) | ErsteinbauState::New => {
ErsteinbauState::Abgelehnt {
antwort_code: None,
reason: format!(
"Antwortfrist abgelaufen ({label}) — ohne Zustimmung des wMSB darf \
der gMSB kein iMS einbauen"
),
}
}
other => other,
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
ErsteinbauCommand::ReceiveVorabinformation {
pid,
gmsb,
wmsb,
melo_id,
umstellungszeitpunkt,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, ErsteinbauState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if pid.as_u32() != VORABINFORMATION_PID {
return Err(WorkflowError::rejected(format!(
"PID {} is not the Vorabinformation zum Ersteinbau ({VORABINFORMATION_PID})",
pid.as_u32()
)));
}
if !validation_passed {
return Ok(vec![ErsteinbauEvent::Rejected {
reason: validation_errors.join("; "),
}]
.into());
}
Ok(vec![ErsteinbauEvent::VorabinformationEmpfangen {
melo_id,
gmsb,
wmsb,
umstellungszeitpunkt,
message_ref,
}]
.into())
}
ErsteinbauCommand::SendVorabinformation {
gmsb,
wmsb,
melo_id,
umstellungszeitpunkt,
message_ref,
} => {
if !matches!(state, ErsteinbauState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
let payload = serde_json::json!({
"pid": VORABINFORMATION_PID,
"sender": gmsb.as_str(),
"receiver": wmsb.as_str(),
"melo": melo_id.as_str(),
"umstellungszeitpunkt": umstellungszeitpunkt,
});
Ok(WorkflowOutput::with_outbox(
vec![ErsteinbauEvent::VorabinformationGesendet {
melo_id,
gmsb,
wmsb: wmsb.clone(),
umstellungszeitpunkt,
message_ref,
}],
vec![PendingOutbox::new("IFTSTA", wmsb.as_str(), payload).caused_by(0)],
))
}
ErsteinbauCommand::DispatchAntwort { antwort_code } => {
let ErsteinbauState::Angekuendigt(data) = state else {
return Err(WorkflowError::invalid_state("Angekuendigt", state.label()));
};
let code = mako_pruefung::codes::lookup(ERSTEINBAU_EBD, &antwort_code).ok_or_else(
|| {
WorkflowError::rejected(format!(
"Antwortcode {antwort_code:?} is not published in {ERSTEINBAU_EBD}"
))
},
)?;
let zustimmung = code.ist_zustimmung().ok_or_else(|| {
WorkflowError::rejected(format!("{} sits off the agreement axis", code.code))
})?;
let antwort_pid = if zustimmung {
ZUSTIMMUNG_PID
} else {
ABLEHNUNG_PID
};
let payload = serde_json::json!({
"pid": antwort_pid,
"sender": data.wmsb.as_str(),
"receiver": data.gmsb.as_str(),
"melo": data.melo_id.as_str(),
"antwort_code": code.code,
"antwort_codeliste": ERSTEINBAU_EBD,
"orig_message_ref": data.message_ref.as_str(),
});
Ok(WorkflowOutput::with_outbox(
vec![ErsteinbauEvent::AntwortGesendet {
pruefidentifikator: Pruefidentifikator::new(antwort_pid)
.map_err(WorkflowError::rejected)?,
antwort_code: code.code.to_owned(),
}],
vec![PendingOutbox::new("IFTSTA", data.gmsb.as_str(), payload).caused_by(0)],
))
}
ErsteinbauCommand::MeldeScheitern { reason } => {
let ErsteinbauState::Zugestimmt(data) = state else {
return Err(WorkflowError::invalid_state("Zugestimmt", state.label()));
};
let payload = serde_json::json!({
"pid": SCHEITERN_PID,
"sender": data.wmsb.as_str(),
"receiver": data.gmsb.as_str(),
"melo": data.melo_id.as_str(),
"grund": reason,
});
Ok(WorkflowOutput::with_outbox(
vec![ErsteinbauEvent::ScheiternGemeldet { reason }],
vec![PendingOutbox::new("IFTSTA", data.gmsb.as_str(), payload).caused_by(0)],
))
}
ErsteinbauCommand::TimeoutExpired { deadline_id, label } => {
if !matches!(state, ErsteinbauState::Angekuendigt(_)) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![ErsteinbauEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[derive(Debug, Default)]
pub struct ErsteinbauRecord {
pub status: &'static str,
pub melo_id: Option<String>,
pub antwort_code: Option<String>,
pub rollout_freigegeben: bool,
}
#[derive(Debug, Default)]
pub struct ErsteinbauProjection {
pub records: HashMap<String, ErsteinbauRecord>,
}
impl Projection for ErsteinbauProjection {
fn name(&self) -> &'static str {
"ErsteinbauProjection"
}
fn handle_event(&mut self, envelope: &EventEnvelope) {
let Ok(event) = envelope.decode::<ErsteinbauEvent>() else {
return;
};
let record = self
.records
.entry(envelope.stream_id.as_str().to_owned())
.or_default();
match event {
ErsteinbauEvent::VorabinformationEmpfangen { melo_id, .. }
| ErsteinbauEvent::VorabinformationGesendet { melo_id, .. } => {
record.status = "Angekuendigt";
record.melo_id = Some(melo_id.as_str().to_owned());
}
ErsteinbauEvent::AntwortGesendet {
pruefidentifikator,
antwort_code,
} => {
let zustimmung = pruefidentifikator.as_u32() == ZUSTIMMUNG_PID;
record.status = if zustimmung {
"Zugestimmt"
} else {
"Abgelehnt"
};
record.rollout_freigegeben = zustimmung;
record.antwort_code = Some(antwort_code);
}
ErsteinbauEvent::ScheiternGemeldet { .. } => {
record.status = "Gescheitert";
record.rollout_freigegeben = false;
}
ErsteinbauEvent::Rejected { .. } => record.status = "Rejected",
ErsteinbauEvent::DeadlineExpired { .. } => {
if record.status == "Angekuendigt" {
record.status = "Abgelehnt";
record.rollout_freigegeben = false;
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn mp(s: &str) -> MarktpartnerCode {
MarktpartnerCode::new(s)
}
fn melo() -> MeLo {
MeLo::new("DE0000000001234567890000000000001")
}
fn vorabinformation() -> ErsteinbauCommand {
ErsteinbauCommand::ReceiveVorabinformation {
pid: Pruefidentifikator::new(VORABINFORMATION_PID).expect("valid PID"),
gmsb: mp("9900000000001"),
wmsb: mp("9900000000003"),
melo_id: melo(),
umstellungszeitpunkt: "20261201".to_owned(),
message_ref: MessageRef::new("MSG-1"),
validation_passed: true,
validation_errors: vec![],
}
}
fn angekuendigt() -> ErsteinbauState {
let out = WimErsteinbauWorkflow::handle(&ErsteinbauState::New, vorabinformation())
.expect("accepted");
out.events
.iter()
.fold(ErsteinbauState::New, WimErsteinbauWorkflow::apply)
}
#[test]
fn a_vorabinformation_opens_the_answer_window() {
assert!(matches!(angekuendigt(), ErsteinbauState::Angekuendigt(_)));
}
#[test]
fn a03_clears_the_rollout_and_rides_21030() {
let out = WimErsteinbauWorkflow::handle(
&angekuendigt(),
ErsteinbauCommand::DispatchAntwort {
antwort_code: "A03".to_owned(),
},
)
.expect("A03 is published in E_0233");
let ErsteinbauEvent::AntwortGesendet {
pruefidentifikator, ..
} = &out.events[0]
else {
panic!("expected an AntwortGesendet");
};
assert_eq!(pruefidentifikator.as_u32(), ZUSTIMMUNG_PID);
let state = out
.events
.iter()
.fold(angekuendigt(), WimErsteinbauWorkflow::apply);
assert!(state.rollout_freigegeben());
}
#[test]
fn a04_is_a_refusal_and_blocks_the_rollout() {
let out = WimErsteinbauWorkflow::handle(
&angekuendigt(),
ErsteinbauCommand::DispatchAntwort {
antwort_code: "A04".to_owned(),
},
)
.expect("A04 is published in E_0233");
let ErsteinbauEvent::AntwortGesendet {
pruefidentifikator, ..
} = &out.events[0]
else {
panic!("expected an AntwortGesendet");
};
assert_eq!(pruefidentifikator.as_u32(), ABLEHNUNG_PID);
let state = out
.events
.iter()
.fold(angekuendigt(), WimErsteinbauWorkflow::apply);
assert!(!state.rollout_freigegeben());
assert!(matches!(state, ErsteinbauState::Abgelehnt { .. }));
}
#[test]
fn every_ablehnungscode_rides_21031() {
for c in ["A01", "A02", "A04"] {
let out = WimErsteinbauWorkflow::handle(
&angekuendigt(),
ErsteinbauCommand::DispatchAntwort {
antwort_code: c.to_owned(),
},
)
.expect("published in E_0233");
let ErsteinbauEvent::AntwortGesendet {
pruefidentifikator, ..
} = &out.events[0]
else {
panic!("expected an AntwortGesendet");
};
assert_eq!(pruefidentifikator.as_u32(), ABLEHNUNG_PID, "code {c}");
}
}
#[test]
fn an_expired_antwortfrist_blocks_the_rollout() {
let state = WimErsteinbauWorkflow::apply(
angekuendigt(),
&ErsteinbauEvent::DeadlineExpired {
deadline_id: DeadlineId::new(),
label: ANTWORT_WINDOW_LABEL.into(),
},
);
assert!(!state.rollout_freigegeben());
let ErsteinbauState::Abgelehnt { antwort_code, .. } = &state else {
panic!("expected Abgelehnt, got {}", state.label());
};
assert_eq!(*antwort_code, None, "no code was ever stated");
}
#[test]
fn a_code_from_another_tree_is_refused() {
assert!(
WimErsteinbauWorkflow::handle(
&angekuendigt(),
ErsteinbauCommand::DispatchAntwort {
antwort_code: "E15".to_owned(),
},
)
.is_err()
);
}
#[test]
fn a_scheitermeldung_needs_a_consented_rollout() {
assert!(
WimErsteinbauWorkflow::handle(
&angekuendigt(),
ErsteinbauCommand::MeldeScheitern {
reason: "kein Zugang".to_owned(),
},
)
.is_err(),
"a rollout that was never consented to cannot fail"
);
}
#[test]
fn the_wrong_pid_is_refused() {
let mut cmd = vorabinformation();
if let ErsteinbauCommand::ReceiveVorabinformation { ref mut pid, .. } = cmd {
*pid = Pruefidentifikator::new(21_007).expect("valid PID");
}
assert!(WimErsteinbauWorkflow::handle(&ErsteinbauState::New, cmd).is_err());
}
#[test]
fn the_answer_window_is_the_published_one() {
assert_eq!(mako_fristen::antwort::ERSTEINBAU_ANTWORT_WERKTAGE, 3);
let o = mako_fristen::antwort::WIM
.iter()
.find(|o| o.trigger_pid == VORABINFORMATION_PID)
.expect("21029 is in the WiM Antwortfrist table");
assert_eq!(o.antwort_pids, (ZUSTIMMUNG_PID, ABLEHNUNG_PID));
assert_eq!(o.ebd, Some(ERSTEINBAU_EBD));
let via_lookup = mako_fristen::antwort::antwort_obligation(VORABINFORMATION_PID)
.expect("21029 is discoverable by PID");
assert_eq!(via_lookup.family.as_str(), "wim");
}
}