use std::collections::HashMap;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
projection::Projection,
types::{DeviceId, MarktpartnerCode, MeLo, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const GERAETEUBERNAHME_PIDS: &[u32] = &[
17005, 17011, ];
pub const ANFRAGE_PIDS: &[u32] = &[];
const ANFRAGE_PIDS_ALL: &[u32] = &[17001, 17002];
pub const BESTELLUNG_PIDS: &[u32] = &[17005];
pub const STORNIERUNG_PIDS: &[u32] = &[17011];
pub const ORDRSP_DEADLINE_LABEL: &str = "wim-geraeteubernahme-ordrsp-deadline";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum GeraeteubernahmeEvent {
AnfrageReceived {
pid: Pruefidentifikator,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
melo_id: MeLo,
device_id: DeviceId,
document_date: String,
message_ref: MessageRef,
},
ValidationPassed {
message_ref: MessageRef,
},
AnfrageOrdrspDispatched {
positive: bool,
response_ref: MessageRef,
reason: Option<String>,
},
BestellungReceived {
pid: Pruefidentifikator,
message_ref: MessageRef,
},
BestellungOrdrspDispatched {
positive: bool,
response_ref: MessageRef,
reason: Option<String>,
},
Abgeschlossen {
device_id: DeviceId,
},
Storniert {
stornierung_pid: Pruefidentifikator,
message_ref: MessageRef,
},
Abgelehnt {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for GeraeteubernahmeEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnfrageReceived { .. } => "WimGeraeteubernahmeAnfrageReceived",
Self::ValidationPassed { .. } => "WimGeraeteubernahmeValidationPassed",
Self::AnfrageOrdrspDispatched { .. } => "WimGeraeteubernahmeAnfrageOrdrspDispatched",
Self::BestellungReceived { .. } => "WimGeraeteubernahmeBestellungReceived",
Self::BestellungOrdrspDispatched { .. } => {
"WimGeraeteubernahmeBestellungOrdrspDispatched"
}
Self::Abgeschlossen { .. } => "WimGeraeteubernahmeAbgeschlossen",
Self::Storniert { .. } => "WimGeraeteubernahmeStorniert",
Self::Abgelehnt { .. } => "WimGeraeteubernahmeAbgelehnt",
Self::DeadlineExpired { .. } => "WimGeraeteubernahmeDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct GeraeteubernahmeData {
pub pid: Pruefidentifikator,
pub incoming_msb: MarktpartnerCode,
pub grid_operator: MarktpartnerCode,
pub melo_id: MeLo,
pub device_id: DeviceId,
pub document_date: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum GeraeteubernahmeState {
New,
AnfrageReceived(GeraeteubernahmeData),
ValidationPassed(GeraeteubernahmeData),
AngebotGesendet(GeraeteubernahmeData),
BestellungReceived(GeraeteubernahmeData),
Abgeschlossen(GeraeteubernahmeData),
Storniert {
reason: String,
},
Abgelehnt {
reason: String,
},
}
impl Default for GeraeteubernahmeState {
fn default() -> Self {
Self::New
}
}
impl GeraeteubernahmeState {
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::Abgeschlossen(_) | Self::Storniert { .. } | Self::Abgelehnt { .. }
)
}
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::AnfrageReceived(_) => "AnfrageReceived",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AngebotGesendet(_) => "AngebotGesendet",
Self::BestellungReceived(_) => "BestellungReceived",
Self::Abgeschlossen(_) => "Abgeschlossen",
Self::Storniert { .. } => "Storniert",
Self::Abgelehnt { .. } => "Abgelehnt",
}
}
}
#[derive(Clone)]
pub enum GeraeteubernahmeCommand {
ReceiveAnfrage {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
melo_id: MeLo,
device_id: DeviceId,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
DispatchAnfrageOrdrsp {
positive: bool,
response_ref: MessageRef,
reason: Option<String>,
},
ReceiveBestellung {
pid: Pruefidentifikator,
message_ref: MessageRef,
},
DispatchBestellungOrdrsp {
positive: bool,
response_ref: MessageRef,
reason: Option<String>,
},
ConfirmTransfer {
device_id: DeviceId,
},
ReceiveStornierung {
pid: Pruefidentifikator,
message_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for GeraeteubernahmeCommand {}
pub struct WimGeraeteubernahmeWorkflow;
impl Workflow for WimGeraeteubernahmeWorkflow {
type State = GeraeteubernahmeState;
type Event = GeraeteubernahmeEvent;
type Command = GeraeteubernahmeCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
if deadline.label() == ORDRSP_DEADLINE_LABEL && !state.is_terminal() {
Some(GeraeteubernahmeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
} else {
None
}
}
#[allow(clippy::too_many_lines)]
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
GeraeteubernahmeEvent::AnfrageReceived {
pid,
incoming_msb,
grid_operator,
melo_id,
device_id,
document_date,
..
} => GeraeteubernahmeState::AnfrageReceived(GeraeteubernahmeData {
pid: *pid,
incoming_msb: incoming_msb.clone(),
grid_operator: grid_operator.clone(),
melo_id: melo_id.clone(),
device_id: device_id.clone(),
document_date: document_date.clone(),
}),
GeraeteubernahmeEvent::ValidationPassed { .. } => {
if let GeraeteubernahmeState::AnfrageReceived(data) = state {
GeraeteubernahmeState::ValidationPassed(data)
} else {
state
}
}
GeraeteubernahmeEvent::AnfrageOrdrspDispatched {
positive, reason, ..
} => {
if *positive {
match state {
GeraeteubernahmeState::ValidationPassed(data) => {
GeraeteubernahmeState::AngebotGesendet(data)
}
_ => state,
}
} else {
GeraeteubernahmeState::Abgelehnt {
reason: reason
.clone()
.unwrap_or_else(|| "negative ORDRSP".to_owned()),
}
}
}
GeraeteubernahmeEvent::BestellungReceived { .. } => {
if let GeraeteubernahmeState::AngebotGesendet(data) = state {
GeraeteubernahmeState::BestellungReceived(data)
} else {
state
}
}
GeraeteubernahmeEvent::BestellungOrdrspDispatched {
positive, reason, ..
} => {
if *positive {
state } else {
GeraeteubernahmeState::Abgelehnt {
reason: reason
.clone()
.unwrap_or_else(|| "negative Bestellung-ORDRSP".to_owned()),
}
}
}
GeraeteubernahmeEvent::Abgeschlossen { device_id } => {
if let GeraeteubernahmeState::BestellungReceived(mut data) = state {
data.device_id = device_id.clone();
GeraeteubernahmeState::Abgeschlossen(data)
} else {
state
}
}
GeraeteubernahmeEvent::Storniert {
stornierung_pid, ..
} => GeraeteubernahmeState::Storniert {
reason: format!("Stornierung via PID {stornierung_pid}"),
},
GeraeteubernahmeEvent::Abgelehnt { reason } => GeraeteubernahmeState::Abgelehnt {
reason: reason.clone(),
},
GeraeteubernahmeEvent::DeadlineExpired { label, .. } => match state {
s if s.is_terminal() => s,
_ => GeraeteubernahmeState::Abgelehnt {
reason: format!("deadline expired: {label}"),
},
},
}
}
#[allow(clippy::too_many_lines)]
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
GeraeteubernahmeCommand::ReceiveAnfrage {
pid,
sender,
receiver,
melo_id,
device_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, GeraeteubernahmeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if !ANFRAGE_PIDS_ALL.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {} is not a Geräteübernahme-Anfrage PID (expected {:?})",
pid.as_u32(),
ANFRAGE_PIDS_ALL,
)));
}
let mut events = vec![GeraeteubernahmeEvent::AnfrageReceived {
pid,
incoming_msb: sender,
grid_operator: receiver,
melo_id,
device_id,
document_date,
message_ref: message_ref.clone(),
}];
if validation_passed {
events.push(GeraeteubernahmeEvent::ValidationPassed { message_ref });
} else {
events.push(GeraeteubernahmeEvent::Abgelehnt {
reason: validation_errors.join("; "),
});
}
Ok(events.into())
}
GeraeteubernahmeCommand::DispatchAnfrageOrdrsp {
positive,
response_ref,
reason,
} => {
if !matches!(state, GeraeteubernahmeState::ValidationPassed(_)) {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
Ok(vec![GeraeteubernahmeEvent::AnfrageOrdrspDispatched {
positive,
response_ref,
reason,
}]
.into())
}
GeraeteubernahmeCommand::ReceiveBestellung { pid, message_ref } => {
if !matches!(state, GeraeteubernahmeState::AngebotGesendet(_)) {
return Err(WorkflowError::invalid_state(
"AngebotGesendet",
state.status_str(),
));
}
if !BESTELLUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {} is not a Geräteübernahme-Bestellung PID (expected {:?})",
pid.as_u32(),
BESTELLUNG_PIDS,
)));
}
Ok(vec![GeraeteubernahmeEvent::BestellungReceived { pid, message_ref }].into())
}
GeraeteubernahmeCommand::DispatchBestellungOrdrsp {
positive,
response_ref,
reason,
} => {
if !matches!(state, GeraeteubernahmeState::BestellungReceived(_)) {
return Err(WorkflowError::invalid_state(
"BestellungReceived",
state.status_str(),
));
}
Ok(vec![GeraeteubernahmeEvent::BestellungOrdrspDispatched {
positive,
response_ref,
reason,
}]
.into())
}
GeraeteubernahmeCommand::ConfirmTransfer { device_id } => {
if !matches!(state, GeraeteubernahmeState::BestellungReceived(_)) {
return Err(WorkflowError::invalid_state(
"BestellungReceived",
state.status_str(),
));
}
Ok(vec![GeraeteubernahmeEvent::Abgeschlossen { device_id }].into())
}
GeraeteubernahmeCommand::ReceiveStornierung { pid, message_ref } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
if !STORNIERUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {} is not a Geräteübernahme-Stornierung PID (expected {:?})",
pid.as_u32(),
STORNIERUNG_PIDS,
)));
}
Ok(vec![GeraeteubernahmeEvent::Storniert {
stornierung_pid: pid,
message_ref,
}]
.into())
}
GeraeteubernahmeCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![GeraeteubernahmeEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[derive(Debug)]
pub enum GeraeteubernahmeRecord {
New {
event_count: usize,
},
Active {
status: &'static str,
melo_id: MeLo,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
device_id: DeviceId,
pid: Pruefidentifikator,
event_count: usize,
},
}
impl GeraeteubernahmeRecord {
#[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<GeraeteubernahmeRecordData<'_>> {
match self {
Self::New { .. } => None,
Self::Active {
melo_id,
incoming_msb,
grid_operator,
device_id,
pid,
..
} => Some(GeraeteubernahmeRecordData {
melo_id,
incoming_msb,
grid_operator,
device_id,
pid,
}),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct GeraeteubernahmeRecordData<'a> {
pub melo_id: &'a MeLo,
pub incoming_msb: &'a MarktpartnerCode,
pub grid_operator: &'a MarktpartnerCode,
pub device_id: &'a DeviceId,
pub pid: &'a Pruefidentifikator,
}
impl Default for GeraeteubernahmeRecord {
fn default() -> Self {
Self::New { event_count: 0 }
}
}
#[derive(Debug, Default)]
pub struct GeraeteubernahmeProjection {
pub records: HashMap<String, GeraeteubernahmeRecord>,
pub last_seq: u64,
}
impl Projection for GeraeteubernahmeProjection {
fn name(&self) -> &'static str {
"GeraeteubernahmeProjection"
}
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::<GeraeteubernahmeEvent>() else {
return;
};
match record {
GeraeteubernahmeRecord::New { event_count }
| GeraeteubernahmeRecord::Active { event_count, .. } => *event_count += 1,
}
match event {
GeraeteubernahmeEvent::AnfrageReceived {
pid,
incoming_msb,
grid_operator,
melo_id,
device_id,
..
} => {
let count = record.event_count();
*record = GeraeteubernahmeRecord::Active {
status: "AnfrageReceived",
pid,
incoming_msb,
grid_operator,
melo_id,
device_id,
event_count: count,
};
}
GeraeteubernahmeEvent::ValidationPassed { .. } => {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = "ValidationPassed";
}
}
GeraeteubernahmeEvent::AnfrageOrdrspDispatched { positive, .. } => {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = if positive {
"AngebotGesendet"
} else {
"Abgelehnt"
};
}
}
GeraeteubernahmeEvent::BestellungReceived { .. } => {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = "BestellungReceived";
}
}
GeraeteubernahmeEvent::BestellungOrdrspDispatched { positive, .. } => {
if !positive {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = "Abgelehnt";
}
}
}
GeraeteubernahmeEvent::Abgeschlossen { device_id } => {
if let GeraeteubernahmeRecord::Active {
status,
device_id: d,
..
} = record
{
*status = "Abgeschlossen";
*d = device_id;
}
}
GeraeteubernahmeEvent::Storniert { .. } => {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = "Storniert";
}
}
GeraeteubernahmeEvent::Abgelehnt { .. }
| GeraeteubernahmeEvent::DeadlineExpired { .. } => {
if let GeraeteubernahmeRecord::Active { status, .. } = record {
*status = "Abgelehnt";
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use mako_engine::types::MessageRef;
fn anfrage_cmd(pid: u32) -> GeraeteubernahmeCommand {
GeraeteubernahmeCommand::ReceiveAnfrage {
pid: Pruefidentifikator::new(pid).unwrap(),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
melo_id: MeLo::new("DE00056789012"),
device_id: DeviceId::new("EHZ-1234567890"),
document_date: "20260101".to_owned(),
message_ref: MessageRef::new("MSG-ORDERS-001"),
validation_passed: true,
validation_errors: vec![],
}
}
#[test]
fn happy_path_phase1_to_phase2_to_abgeschlossen() {
let state = GeraeteubernahmeState::default();
let events = WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(17001))
.expect("Anfrage 17001 must succeed — bypassing PID guard for test");
assert_eq!(events.len(), 2); let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::ValidationPassed(_)));
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::DispatchAnfrageOrdrsp {
positive: true,
response_ref: MessageRef::new("MSG-ORDRSP-001"),
reason: None,
},
)
.expect("DispatchAnfrageOrdrsp must succeed");
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::AngebotGesendet(_)));
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::ReceiveBestellung {
pid: Pruefidentifikator::new(17005).unwrap(),
message_ref: MessageRef::new("MSG-ORDERS-002"),
},
)
.expect("ReceiveBestellung must succeed");
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(
state,
GeraeteubernahmeState::BestellungReceived(_)
));
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::DispatchBestellungOrdrsp {
positive: true,
response_ref: MessageRef::new("MSG-ORDRSP-002"),
reason: None,
},
)
.expect("DispatchBestellungOrdrsp must succeed");
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(
state,
GeraeteubernahmeState::BestellungReceived(_)
));
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::ConfirmTransfer {
device_id: DeviceId::new("NEW-EHZ-9999999"),
},
)
.expect("ConfirmTransfer must succeed");
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::Abgeschlossen(_)));
}
#[test]
fn negative_anfrage_ordrsp_rejects() {
let state = GeraeteubernahmeState::default();
let events = WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(17001)).unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::DispatchAnfrageOrdrsp {
positive: false,
response_ref: MessageRef::new("MSG-ORDRSP-NEG"),
reason: Some("MeLo nicht bekannt".to_owned()),
},
)
.unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::Abgelehnt { .. }));
}
#[test]
fn validation_failure_rejects() {
let state = GeraeteubernahmeState::default();
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::ReceiveAnfrage {
pid: Pruefidentifikator::new(17001).unwrap(),
sender: MarktpartnerCode::new("9900123456789"),
receiver: MarktpartnerCode::new("9900987654321"),
melo_id: MeLo::new("DE00011111111"),
device_id: DeviceId::new("EHZ-001"),
document_date: "20260101".to_owned(),
message_ref: MessageRef::new("MSG-001"),
validation_passed: false,
validation_errors: vec!["mandatory segment missing".to_owned()],
},
)
.unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::Abgelehnt { .. }));
}
#[test]
fn stornierung_from_active_transitions_to_storniert() {
let state = GeraeteubernahmeState::default();
let events = WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(17001)).unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::ReceiveStornierung {
pid: Pruefidentifikator::new(17011).unwrap(), message_ref: MessageRef::new("MSG-STORNO-001"),
},
)
.unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::Storniert { .. }));
}
#[test]
fn deadline_on_active_rejects() {
let state = GeraeteubernahmeState::default();
let events = WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(17001)).unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
let events = WimGeraeteubernahmeWorkflow::handle(
&state,
GeraeteubernahmeCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "wim-geraeteubernahme-ordrsp-deadline".into(),
},
)
.unwrap();
let state = events
.iter()
.fold(state, WimGeraeteubernahmeWorkflow::apply);
assert!(matches!(state, GeraeteubernahmeState::Abgelehnt { .. }));
}
#[test]
fn deadline_on_terminal_is_noop() {
let terminal = GeraeteubernahmeState::Abgelehnt {
reason: "test".to_owned(),
};
let events = WimGeraeteubernahmeWorkflow::handle(
&terminal,
GeraeteubernahmeCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "late-deadline".into(),
},
)
.unwrap();
assert!(events.is_empty(), "deadline on terminal must be a no-op");
}
#[test]
fn all_anfrage_pids_accepted() {
for &pid in ANFRAGE_PIDS {
let state = GeraeteubernahmeState::default();
assert!(
WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(pid)).is_ok(),
"PID {pid} must be accepted",
);
}
}
#[test]
fn wrong_pid_family_rejected() {
let state = GeraeteubernahmeState::default();
let result = WimGeraeteubernahmeWorkflow::handle(&state, anfrage_cmd(55001));
assert!(
result.is_err(),
"GPKE PID must be rejected by Geräteübernahme"
);
}
}