use std::collections::HashMap;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
projection::Projection,
types::{MarktpartnerCode, MeLo, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-stammdaten";
pub const ANFORDERUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(17132);
pub const ANFORDERUNG_PID_GAS: Pruefidentifikator = Pruefidentifikator::const_new(17101);
pub const UEBERMITTLUNG_PIDS: std::ops::RangeInclusive<u32> = 17102..=17133;
pub const STAMMDATEN_DEADLINE_LABEL: &str = "wim-stammdaten-deadline";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum StammdatenEvent {
AnforderungReceived {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
document_date: String,
message_ref: MessageRef,
},
ValidationPassed {
message_ref: MessageRef,
},
StammdatenUebermittelt {
response_pid: Pruefidentifikator,
response_ref: MessageRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
standorteigenschaften: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
zaehlwerke: Vec<serde_json::Value>,
},
Abgelehnt {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for StammdatenEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnforderungReceived { .. } => "WimStammdatenAnforderungReceived",
Self::ValidationPassed { .. } => "WimStammdatenValidationPassed",
Self::StammdatenUebermittelt { .. } => "WimStammdatenUebermittelt",
Self::Abgelehnt { .. } => "WimStammdatenAbgelehnt",
Self::DeadlineExpired { .. } => "WimStammdatenDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StammdatenData {
pub pid: Pruefidentifikator,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub melo_id: MeLo,
pub document_date: String,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum StammdatenState {
#[default]
New,
AnforderungReceived(StammdatenData),
ValidationPassed(StammdatenData),
Uebermittelt {
data: StammdatenData,
response_pid: Pruefidentifikator,
},
Abgelehnt {
reason: String,
},
}
impl StammdatenState {
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(self, Self::Uebermittelt { .. } | Self::Abgelehnt { .. })
}
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::AnforderungReceived(_) => "AnforderungReceived",
Self::ValidationPassed(_) => "ValidationPassed",
Self::Uebermittelt { .. } => "Uebermittelt",
Self::Abgelehnt { .. } => "Abgelehnt",
}
}
}
#[derive(Clone)]
pub enum StammdatenCommand {
ReceiveAnforderung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
TransmitStammdaten {
response_pid: Pruefidentifikator,
response_ref: MessageRef,
standorteigenschaften: Option<serde_json::Value>,
zaehlwerke: Vec<serde_json::Value>,
},
RejectAnforderung {
reason: String,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for StammdatenCommand {}
pub struct WimStammdatenWorkflow;
impl Workflow for WimStammdatenWorkflow {
type State = StammdatenState;
type Event = StammdatenEvent;
type Command = StammdatenCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
STAMMDATEN_DEADLINE_LABEL,
StammdatenState::AnforderungReceived(_) | StammdatenState::ValidationPassed(_),
) => Some(StammdatenCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
StammdatenEvent::AnforderungReceived {
pid,
sender,
receiver,
melo_id,
document_date,
..
} => StammdatenState::AnforderungReceived(StammdatenData {
pid: *pid,
sender: sender.clone(),
receiver: receiver.clone(),
melo_id: melo_id.clone(),
document_date: document_date.clone(),
}),
StammdatenEvent::ValidationPassed { .. } => {
if let StammdatenState::AnforderungReceived(data) = state {
StammdatenState::ValidationPassed(data)
} else {
state
}
}
StammdatenEvent::StammdatenUebermittelt { response_pid, .. } => {
if let StammdatenState::ValidationPassed(data) = state {
StammdatenState::Uebermittelt {
data,
response_pid: *response_pid,
}
} else {
state
}
}
StammdatenEvent::Abgelehnt { reason } => StammdatenState::Abgelehnt {
reason: reason.clone(),
},
StammdatenEvent::DeadlineExpired { label, .. } => match state {
s if s.is_terminal() => s,
_ => StammdatenState::Abgelehnt {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
StammdatenCommand::ReceiveAnforderung {
pid,
sender,
receiver,
melo_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, StammdatenState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if pid != ANFORDERUNG_PID {
return Err(WorkflowError::rejected(format!(
"PID {} is not a Stammdaten-Anforderung PID (expected {ANFORDERUNG_PID})",
pid.as_u32()
)));
}
let mut events = vec![StammdatenEvent::AnforderungReceived {
pid,
sender,
receiver,
melo_id,
document_date,
message_ref: message_ref.clone(),
}];
if validation_passed {
events.push(StammdatenEvent::ValidationPassed { message_ref });
} else {
events.push(StammdatenEvent::Abgelehnt {
reason: validation_errors.join("; "),
});
}
Ok(events.into())
}
StammdatenCommand::TransmitStammdaten {
response_pid,
response_ref,
standorteigenschaften,
zaehlwerke,
} => {
let StammdatenState::ValidationPassed(data) = &state else {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
};
if !UEBERMITTLUNG_PIDS.contains(&response_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {} is not a Stammdaten-Übermittlung PID (expected 17102–17133)",
response_pid.as_u32()
)));
}
let melo_id = data.melo_id.as_str().to_owned();
let event = StammdatenEvent::StammdatenUebermittelt {
response_pid,
response_ref,
standorteigenschaften: standorteigenschaften.clone(),
zaehlwerke: zaehlwerke.clone(),
};
if !zaehlwerke.is_empty() || standorteigenschaften.is_some() {
let mut payload = serde_json::json!({
"melo_id": melo_id,
"pid": response_pid.as_u32(),
});
if !zaehlwerke.is_empty() {
payload["zaehlwerke"] = serde_json::Value::Array(zaehlwerke);
}
if let Some(se) = standorteigenschaften {
payload["standorteigenschaften"] = se;
}
let outbox =
mako_engine::outbox::PendingOutbox::new("ProcessCompleted", "", payload);
Ok(mako_engine::workflow::WorkflowOutput::with_outbox(
vec![event],
vec![outbox],
))
} else {
Ok(vec![event].into())
}
}
StammdatenCommand::RejectAnforderung { reason } => {
if !matches!(
state,
StammdatenState::AnforderungReceived(_) | StammdatenState::ValidationPassed(_)
) {
return Err(WorkflowError::invalid_state(
"AnforderungReceived or ValidationPassed",
state.status_str(),
));
}
Ok(vec![StammdatenEvent::Abgelehnt { reason }].into())
}
StammdatenCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![StammdatenEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[derive(Debug)]
pub enum StammdatenRecord {
New {
event_count: usize,
},
Active {
status: &'static str,
melo_id: MeLo,
sender: MarktpartnerCode,
response_pid: Option<Pruefidentifikator>,
event_count: usize,
},
}
impl StammdatenRecord {
#[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<StammdatenRecordData<'_>> {
match self {
Self::New { .. } => None,
Self::Active {
melo_id,
sender,
response_pid,
..
} => Some(StammdatenRecordData {
melo_id,
sender,
response_pid: response_pid.as_ref(),
}),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct StammdatenRecordData<'a> {
pub melo_id: &'a MeLo,
pub sender: &'a MarktpartnerCode,
pub response_pid: Option<&'a Pruefidentifikator>,
}
impl Default for StammdatenRecord {
fn default() -> Self {
Self::New { event_count: 0 }
}
}
#[derive(Debug, Default)]
pub struct StammdatenProjection {
pub records: HashMap<String, StammdatenRecord>,
pub last_seq: u64,
}
impl Projection for StammdatenProjection {
fn name(&self) -> &'static str {
"StammdatenProjection"
}
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::<StammdatenEvent>() else {
return;
};
match record {
StammdatenRecord::New { event_count }
| StammdatenRecord::Active { event_count, .. } => *event_count += 1,
}
match event {
StammdatenEvent::AnforderungReceived {
sender, melo_id, ..
} => {
let count = record.event_count();
*record = StammdatenRecord::Active {
status: "AnforderungReceived",
melo_id,
sender,
response_pid: None,
event_count: count,
};
}
StammdatenEvent::ValidationPassed { .. } => {
if let StammdatenRecord::Active { status, .. } = record {
*status = "ValidationPassed";
}
}
StammdatenEvent::StammdatenUebermittelt { response_pid, .. } => {
if let StammdatenRecord::Active {
status,
response_pid: rp,
..
} = record
{
*status = "Uebermittelt";
*rp = Some(response_pid);
}
}
StammdatenEvent::Abgelehnt { .. } | StammdatenEvent::DeadlineExpired { .. } => {
if let StammdatenRecord::Active { status, .. } = record {
*status = "Abgelehnt";
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn anforderung_cmd(pid: u32, valid: bool) -> StammdatenCommand {
StammdatenCommand::ReceiveAnforderung {
pid: Pruefidentifikator::new(pid).unwrap(),
sender: MarktpartnerCode::new("9900123456789"),
receiver: MarktpartnerCode::new("4012345000023"),
melo_id: MeLo::new("DE00056789012"),
document_date: "20260101".to_owned(),
message_ref: MessageRef::new("MSG-17132-001"),
validation_passed: valid,
validation_errors: if valid {
vec![]
} else {
vec!["error".to_owned()]
},
}
}
#[test]
fn happy_path_anforderung_to_uebermittlung() {
let state = StammdatenState::default();
let events = WimStammdatenWorkflow::handle(&state, anforderung_cmd(17132, true)).unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
assert!(matches!(state, StammdatenState::ValidationPassed(_)));
let events = WimStammdatenWorkflow::handle(
&state,
StammdatenCommand::TransmitStammdaten {
response_pid: Pruefidentifikator::new(17102).unwrap(),
response_ref: MessageRef::new("MSG-17102-001"),
standorteigenschaften: None,
zaehlwerke: vec![],
},
)
.unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
assert!(matches!(state, StammdatenState::Uebermittelt { .. }));
}
#[test]
fn validation_failure_rejects() {
let state = StammdatenState::default();
let events = WimStammdatenWorkflow::handle(&state, anforderung_cmd(17132, false)).unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
assert!(matches!(state, StammdatenState::Abgelehnt { .. }));
}
#[test]
fn wrong_anforderung_pid_is_rejected() {
let state = StammdatenState::default();
let result = WimStammdatenWorkflow::handle(&state, anforderung_cmd(17102, true));
assert!(
result.is_err(),
"PID 17102 is Übermittlung, not Anforderung"
);
}
#[test]
fn deadline_on_active_rejects() {
let state = StammdatenState::default();
let events = WimStammdatenWorkflow::handle(&state, anforderung_cmd(17132, true)).unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
let events = WimStammdatenWorkflow::handle(
&state,
StammdatenCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "wim-stammdaten-deadline".into(),
},
)
.unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
assert!(matches!(state, StammdatenState::Abgelehnt { .. }));
}
#[test]
fn deadline_on_terminal_is_noop() {
let terminal = StammdatenState::Abgelehnt {
reason: "test".to_owned(),
};
let events = WimStammdatenWorkflow::handle(
&terminal,
StammdatenCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "late".into(),
},
)
.unwrap();
assert!(events.is_empty());
}
#[test]
fn all_uebermittlung_pids_accepted() {
for pid in UEBERMITTLUNG_PIDS {
let state = StammdatenState::default();
let events =
WimStammdatenWorkflow::handle(&state, anforderung_cmd(17132, true)).unwrap();
let state = events.iter().fold(state, WimStammdatenWorkflow::apply);
let result = WimStammdatenWorkflow::handle(
&state,
StammdatenCommand::TransmitStammdaten {
response_pid: Pruefidentifikator::new(pid).unwrap(),
response_ref: MessageRef::new("MSG-RESP"),
standorteigenschaften: None,
zaehlwerke: vec![],
},
);
assert!(result.is_ok(), "PID {pid} must be accepted as Übermittlung");
}
}
}