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 STORNIERUNG_PID: u32 = 39_002;
pub const BESTAETIGUNG_PID: u32 = 19_013;
pub const ABLEHNUNG_PID: u32 = 19_014;
pub const STORNIERUNG_DEADLINE_LABEL: &str = "wim-stornierung-deadline";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum StornierungEvent {
StornierungReceived {
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
document_date: String,
message_ref: MessageRef,
cancelled_ref: Option<MessageRef>,
},
ValidationPassed {
message_ref: MessageRef,
},
Bestaetigt {
response_ref: MessageRef,
},
Abgelehnt {
reason: String,
response_ref: MessageRef,
},
ValidationFailed {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for StornierungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::StornierungReceived { .. } => "WimStornierungReceived",
Self::ValidationPassed { .. } => "WimStornierungValidationPassed",
Self::Bestaetigt { .. } => "WimStornierungBestaetigt",
Self::Abgelehnt { .. } => "WimStornierungAbgelehnt",
Self::ValidationFailed { .. } => "WimStornierungValidationFailed",
Self::DeadlineExpired { .. } => "WimStornierungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StornierungData {
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub melo_id: MeLo,
pub document_date: String,
pub cancelled_ref: Option<MessageRef>,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum StornierungState {
#[default]
New,
StornierungReceived(StornierungData),
ValidationPassed(StornierungData),
Bestaetigt(StornierungData),
Abgelehnt {
reason: String,
},
}
impl StornierungState {
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(self, Self::Bestaetigt(_) | Self::Abgelehnt { .. })
}
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::StornierungReceived(_) => "StornierungReceived",
Self::ValidationPassed(_) => "ValidationPassed",
Self::Bestaetigt(_) => "Bestaetigt",
Self::Abgelehnt { .. } => "Abgelehnt",
}
}
}
#[derive(Clone)]
pub enum StornierungCommand {
ReceiveOrdchg {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
document_date: String,
message_ref: MessageRef,
cancelled_ref: Option<MessageRef>,
validation_passed: bool,
validation_errors: Vec<String>,
},
Accept {
response_ref: MessageRef,
},
Reject {
reason: String,
response_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for StornierungCommand {}
pub struct WimStornierungWorkflow;
impl Workflow for WimStornierungWorkflow {
type State = StornierungState;
type Event = StornierungEvent;
type Command = StornierungCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
STORNIERUNG_DEADLINE_LABEL,
StornierungState::StornierungReceived(_) | StornierungState::ValidationPassed(_),
) => Some(StornierungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
StornierungEvent::StornierungReceived {
sender,
receiver,
melo_id,
document_date,
cancelled_ref,
..
} => StornierungState::StornierungReceived(StornierungData {
sender: sender.clone(),
receiver: receiver.clone(),
melo_id: melo_id.clone(),
document_date: document_date.clone(),
cancelled_ref: cancelled_ref.clone(),
}),
StornierungEvent::ValidationPassed { .. } => {
if let StornierungState::StornierungReceived(data) = state {
StornierungState::ValidationPassed(data)
} else {
state
}
}
StornierungEvent::Bestaetigt { .. } => {
if let StornierungState::ValidationPassed(data) = state {
StornierungState::Bestaetigt(data)
} else {
state
}
}
StornierungEvent::Abgelehnt { reason, .. } => StornierungState::Abgelehnt {
reason: reason.clone(),
},
StornierungEvent::ValidationFailed { reason } => StornierungState::Abgelehnt {
reason: reason.clone(),
},
StornierungEvent::DeadlineExpired { label, .. } => match state {
s if s.is_terminal() => s,
_ => StornierungState::Abgelehnt {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
StornierungCommand::ReceiveOrdchg {
pid,
sender,
receiver,
melo_id,
document_date,
message_ref,
cancelled_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, StornierungState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if pid.as_u32() != STORNIERUNG_PID {
return Err(WorkflowError::rejected(format!(
"PID {} is not the WiM Stornierung PID (expected {STORNIERUNG_PID})",
pid.as_u32()
)));
}
let mut events = vec![StornierungEvent::StornierungReceived {
sender,
receiver,
melo_id,
document_date,
message_ref: message_ref.clone(),
cancelled_ref,
}];
if validation_passed {
events.push(StornierungEvent::ValidationPassed { message_ref });
} else {
events.push(StornierungEvent::ValidationFailed {
reason: validation_errors.join("; "),
});
}
Ok(events.into())
}
StornierungCommand::Accept { response_ref } => {
if !matches!(state, StornierungState::ValidationPassed(_)) {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
Ok(vec![StornierungEvent::Bestaetigt { response_ref }].into())
}
StornierungCommand::Reject {
reason,
response_ref,
} => {
if !matches!(state, StornierungState::ValidationPassed(_)) {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
Ok(vec![StornierungEvent::Abgelehnt {
reason,
response_ref,
}]
.into())
}
StornierungCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![StornierungEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[derive(Debug)]
pub enum StornierungRecord {
New {
event_count: usize,
},
Active {
status: &'static str,
melo_id: MeLo,
sender: MarktpartnerCode,
event_count: usize,
},
}
impl StornierungRecord {
#[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<StornierungRecordData<'_>> {
match self {
Self::New { .. } => None,
Self::Active {
melo_id, sender, ..
} => Some(StornierungRecordData { melo_id, sender }),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct StornierungRecordData<'a> {
pub melo_id: &'a MeLo,
pub sender: &'a MarktpartnerCode,
}
impl Default for StornierungRecord {
fn default() -> Self {
Self::New { event_count: 0 }
}
}
#[derive(Debug, Default)]
pub struct StornierungProjection {
pub records: HashMap<String, StornierungRecord>,
pub last_seq: u64,
}
impl Projection for StornierungProjection {
fn name(&self) -> &'static str {
"StornierungProjection"
}
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::<StornierungEvent>() else {
return;
};
match record {
StornierungRecord::New { event_count }
| StornierungRecord::Active { event_count, .. } => *event_count += 1,
}
match event {
StornierungEvent::StornierungReceived {
sender, melo_id, ..
} => {
let count = record.event_count();
*record = StornierungRecord::Active {
status: "StornierungReceived",
melo_id,
sender,
event_count: count,
};
}
StornierungEvent::ValidationPassed { .. } => {
if let StornierungRecord::Active { status, .. } = record {
*status = "ValidationPassed";
}
}
StornierungEvent::Bestaetigt { .. } => {
if let StornierungRecord::Active { status, .. } = record {
*status = "Bestaetigt";
}
}
StornierungEvent::Abgelehnt { .. }
| StornierungEvent::ValidationFailed { .. }
| StornierungEvent::DeadlineExpired { .. } => {
if let StornierungRecord::Active { status, .. } = record {
*status = "Abgelehnt";
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn stornierung_cmd(valid: bool) -> StornierungCommand {
StornierungCommand::ReceiveOrdchg {
pid: Pruefidentifikator::new(39_002).unwrap(),
sender: MarktpartnerCode::new("9900123456789"),
receiver: MarktpartnerCode::new("4012345000023"),
melo_id: MeLo::new("DE00056789012"),
document_date: "20260101".to_owned(),
message_ref: MessageRef::new("MSG-ORDCHG-001"),
cancelled_ref: Some(MessageRef::new("MSG-ORDERS-001")),
validation_passed: valid,
validation_errors: if valid {
vec![]
} else {
vec!["missing segment".to_owned()]
},
}
}
#[test]
fn happy_path_stornierung_accepted() {
let state = StornierungState::default();
let events = WimStornierungWorkflow::handle(&state, stornierung_cmd(true)).unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
assert!(matches!(state, StornierungState::ValidationPassed(_)));
let events = WimStornierungWorkflow::handle(
&state,
StornierungCommand::Accept {
response_ref: MessageRef::new("MSG-ORDRSP-39001"),
},
)
.unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
assert!(matches!(state, StornierungState::Bestaetigt(_)));
}
#[test]
fn stornierung_rejected_by_nb() {
let state = StornierungState::default();
let events = WimStornierungWorkflow::handle(&state, stornierung_cmd(true)).unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
let events = WimStornierungWorkflow::handle(
&state,
StornierungCommand::Reject {
reason: "Auftrag bereits ausgeführt".to_owned(),
response_ref: MessageRef::new("MSG-ORDRSP-39002"),
},
)
.unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
assert!(matches!(state, StornierungState::Abgelehnt { .. }));
}
#[test]
fn validation_failure_rejects() {
let state = StornierungState::default();
let events = WimStornierungWorkflow::handle(&state, stornierung_cmd(false)).unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
assert!(matches!(state, StornierungState::Abgelehnt { .. }));
}
#[test]
fn wrong_pid_is_rejected() {
let state = StornierungState::default();
let result = WimStornierungWorkflow::handle(
&state,
StornierungCommand::ReceiveOrdchg {
pid: Pruefidentifikator::new(39_000).unwrap(), sender: MarktpartnerCode::new("9900123456789"),
receiver: MarktpartnerCode::new("4012345000023"),
melo_id: MeLo::new("DE00056789012"),
document_date: "20260101".to_owned(),
message_ref: MessageRef::new("MSG-001"),
cancelled_ref: None,
validation_passed: true,
validation_errors: vec![],
},
);
assert!(
result.is_err(),
"PID 39000 (Gas Sperrung) must be rejected by WiM Stornierung"
);
}
#[test]
fn deadline_on_active_rejects() {
let state = StornierungState::default();
let events = WimStornierungWorkflow::handle(&state, stornierung_cmd(true)).unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
let events = WimStornierungWorkflow::handle(
&state,
StornierungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "wim-stornierung-deadline".into(),
},
)
.unwrap();
let state = events.iter().fold(state, WimStornierungWorkflow::apply);
assert!(matches!(state, StornierungState::Abgelehnt { .. }));
}
#[test]
fn deadline_on_terminal_is_noop() {
let terminal = StornierungState::Bestaetigt(StornierungData {
sender: MarktpartnerCode::new("X"),
receiver: MarktpartnerCode::new("Y"),
melo_id: MeLo::new("Z"),
document_date: String::new(),
cancelled_ref: None,
});
let events = WimStornierungWorkflow::handle(
&terminal,
StornierungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "late".into(),
},
)
.unwrap();
assert!(events.is_empty());
}
}