use std::collections::HashMap;
use mako_engine::types::Pruefidentifikator;
use mako_engine::{
envelope::EventEnvelope,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
projection::Projection,
types::{DeviceId, MarktpartnerCode, MeLo, MessageRef},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-device-change";
pub const APERAK_WINDOW_LABEL: &str = "wim-aperak-5-werktage";
pub const IFTSTA_PIDS: &[u32] = &[
21_007, 21_009, 21_010, 21_011, 21_012, 21_013, 21_015, 21_018, 21_029, 21_030, 21_031, 21_032,
];
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum DeviceChangeEvent {
Initiated {
melo_id: MeLo,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
device_id: DeviceId,
document_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
},
ValidationPassed {
message_ref: MessageRef,
},
AperakDispatched {
positive: bool,
reason: Option<String>,
},
Completed {
device_id: DeviceId,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
IftstaStatusReceived {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
},
}
impl EventPayload for DeviceChangeEvent {
fn event_type(&self) -> &'static str {
match self {
Self::Initiated { .. } => "WimDeviceChangeInitiated",
Self::ValidationPassed { .. } => "WimDeviceChangeValidationPassed",
Self::AperakDispatched { .. } => "WimDeviceChangeAperakDispatched",
Self::Completed { .. } => "WimDeviceChangeCompleted",
Self::Rejected { .. } => "WimDeviceChangeRejected",
Self::DeadlineExpired { .. } => "WimDeviceChangeDeadlineExpired",
Self::IftstaStatusReceived { .. } => "WimDeviceChangeIftstaStatusReceived",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DeviceChangeData {
pub melo_id: MeLo,
pub incoming_msb: MarktpartnerCode,
pub grid_operator: MarktpartnerCode,
pub device_id: DeviceId,
pub document_date: String,
pub pruefidentifikator: Pruefidentifikator,
#[serde(default)]
pub message_ref: Option<MessageRef>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum DeviceChangeState {
New,
Initiated(DeviceChangeData),
ValidationPassed(DeviceChangeData),
AperakSent(DeviceChangeData),
Completed(DeviceChangeData),
Rejected {
reason: String,
},
}
impl Default for DeviceChangeState {
fn default() -> Self {
Self::New
}
}
impl DeviceChangeState {
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::Initiated(_) => "Initiated",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AperakSent(_) => "AperakSent",
Self::Completed(_) => "Completed",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum DeviceChangeCommand {
ReceiveUtilmd {
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>,
},
ReceiveRestOrder {
tx_id: String,
sender_gln: MarktpartnerCode,
melo_id: MeLo,
device_category: String,
process_date: String,
},
DispatchAperak {
positive: bool,
reason: Option<String>,
},
Complete {
device_id: DeviceId,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
ReceiveIftsta {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
}
impl CommandPayload for DeviceChangeCommand {}
pub struct WimDeviceChangeWorkflow;
impl Workflow for WimDeviceChangeWorkflow {
type State = DeviceChangeState;
type Event = DeviceChangeEvent;
type Command = DeviceChangeCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
APERAK_WINDOW_LABEL,
DeviceChangeState::Initiated(_) | DeviceChangeState::ValidationPassed(_),
) => Some(DeviceChangeCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb,
grid_operator,
device_id,
document_date,
message_ref,
pruefidentifikator,
} => DeviceChangeState::Initiated(DeviceChangeData {
melo_id: melo_id.clone(),
incoming_msb: incoming_msb.clone(),
grid_operator: grid_operator.clone(),
device_id: device_id.clone(),
document_date: document_date.clone(),
pruefidentifikator: *pruefidentifikator,
message_ref: Some(message_ref.clone()),
}),
DeviceChangeEvent::ValidationPassed { .. } => {
if let DeviceChangeState::Initiated(data) = state {
DeviceChangeState::ValidationPassed(data)
} else {
state
}
}
DeviceChangeEvent::AperakDispatched { positive, .. } => match state {
DeviceChangeState::ValidationPassed(data) => {
if *positive {
DeviceChangeState::AperakSent(data)
} else {
DeviceChangeState::Rejected {
reason: "negative APERAK".to_owned(),
}
}
}
_ => state,
},
DeviceChangeEvent::Completed { device_id } => {
if let DeviceChangeState::AperakSent(mut data) = state {
data.device_id = device_id.clone();
DeviceChangeState::Completed(data)
} else {
state
}
}
DeviceChangeEvent::Rejected { reason } => DeviceChangeState::Rejected {
reason: reason.clone(),
},
DeviceChangeEvent::DeadlineExpired { label, .. } => match state {
DeviceChangeState::Completed(_) | DeviceChangeState::Rejected { .. } => state,
_ => DeviceChangeState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
DeviceChangeEvent::IftstaStatusReceived { .. } => state,
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
DeviceChangeCommand::ReceiveRestOrder {
tx_id,
sender_gln,
melo_id,
device_category,
process_date,
} => {
if !matches!(state, DeviceChangeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
let pid = Pruefidentifikator::new(55_042).map_err(|e| {
WorkflowError::rejected(format!(
"constant PID 55042 (WiM Anmeldung MSB) invalid: {e}"
))
})?;
let device_id = DeviceId::new(&*tx_id);
let message_ref = MessageRef::new(&*tx_id);
Ok(vec![
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb: sender_gln,
grid_operator: MarktpartnerCode::new(""),
device_id,
document_date: format!("{process_date}|category={device_category}"),
message_ref: message_ref.clone(),
pruefidentifikator: pid,
},
DeviceChangeEvent::ValidationPassed { message_ref },
]
.into())
}
DeviceChangeCommand::ReceiveUtilmd {
pid,
sender,
receiver,
melo_id,
device_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, DeviceChangeState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
let valid_pids = [55_039_u32, 55_042, 55_051, 55_168];
if !valid_pids.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {} is not a WiM Messstellenbetrieb PID (expected 55039, 55042, 55051, or 55168)",
pid.as_u32()
)));
}
let sender_gln = sender.clone();
let receiver_gln = receiver.clone();
let mut events = vec![DeviceChangeEvent::Initiated {
melo_id,
incoming_msb: sender,
grid_operator: receiver,
device_id,
document_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
}];
if validation_passed {
events.push(DeviceChangeEvent::ValidationPassed { message_ref });
Ok(WorkflowOutput::events(events))
} else {
let reason = validation_errors.join("; ");
events.push(DeviceChangeEvent::Rejected {
reason: reason.clone(),
});
let outbox = vec![
PendingOutbox::new(
"APERAK",
sender_gln.as_str(),
serde_json::json!({
"sender": receiver_gln.as_str(),
"receiver": sender_gln.as_str(),
"pid": 29001_u32,
"error_code": "Z29",
"reason": reason,
}),
)
.caused_by(0),
];
Ok(WorkflowOutput::with_outbox(events, outbox))
}
}
DeviceChangeCommand::DispatchAperak { positive, reason } => {
let data = match state {
DeviceChangeState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
};
let mut aperak_payload = serde_json::json!({
"sender": data.grid_operator.as_str(),
"pid": data.pruefidentifikator.as_u32(),
"melo": data.melo_id.as_str(),
"positive": positive,
});
if let Some(ref mr) = data.message_ref {
aperak_payload["orig_message_ref"] =
serde_json::Value::String(mr.as_str().to_owned());
}
if let Some(ref r) = reason {
aperak_payload["reason"] = serde_json::Value::String(r.clone());
}
let outbox_entry =
PendingOutbox::new("APERAK", data.incoming_msb.as_str(), aperak_payload)
.caused_by(0);
Ok(WorkflowOutput::with_outbox(
vec![DeviceChangeEvent::AperakDispatched { positive, reason }],
vec![outbox_entry],
))
}
DeviceChangeCommand::Complete { device_id } => {
if !matches!(state, DeviceChangeState::AperakSent(_)) {
return Err(WorkflowError::invalid_state(
"AperakSent",
state.status_str(),
));
}
Ok(vec![DeviceChangeEvent::Completed { device_id }].into())
}
DeviceChangeCommand::TimeoutExpired { deadline_id, label } => {
if matches!(
state,
DeviceChangeState::Completed(_) | DeviceChangeState::Rejected { .. }
) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![DeviceChangeEvent::DeadlineExpired { deadline_id, label }].into())
}
DeviceChangeCommand::ReceiveIftsta {
pid,
sender,
receiver,
message_ref,
..
} => {
Ok(vec![DeviceChangeEvent::IftstaStatusReceived {
pid,
sender,
receiver,
message_ref,
}]
.into())
}
}
}
}
#[derive(Debug)]
pub enum DeviceChangeRecord {
New {
event_count: usize,
},
Active {
status: &'static str,
melo_id: MeLo,
incoming_msb: MarktpartnerCode,
grid_operator: MarktpartnerCode,
device_id: DeviceId,
pruefidentifikator: Pruefidentifikator,
event_count: usize,
},
}
impl DeviceChangeRecord {
#[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<DeviceChangeRecordData<'_>> {
match self {
Self::New { .. } => None,
Self::Active {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
..
} => Some(DeviceChangeRecordData {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
}),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct DeviceChangeRecordData<'a> {
pub melo_id: &'a MeLo,
pub incoming_msb: &'a MarktpartnerCode,
pub grid_operator: &'a MarktpartnerCode,
pub device_id: &'a DeviceId,
pub pruefidentifikator: &'a Pruefidentifikator,
}
impl Default for DeviceChangeRecord {
fn default() -> Self {
Self::New { event_count: 0 }
}
}
#[derive(Debug, Default)]
pub struct DeviceChangeProjection {
pub records: HashMap<String, DeviceChangeRecord>,
pub last_seq: u64,
}
impl Projection for DeviceChangeProjection {
fn name(&self) -> &'static str {
"DeviceChangeProjection"
}
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::<DeviceChangeEvent>() else {
return;
};
match record {
DeviceChangeRecord::New { event_count } => *event_count += 1,
DeviceChangeRecord::Active { event_count, .. } => *event_count += 1,
}
match event {
DeviceChangeEvent::Initiated {
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
..
} => {
let count = record.event_count();
*record = DeviceChangeRecord::Active {
status: "Initiated",
melo_id,
incoming_msb,
grid_operator,
device_id,
pruefidentifikator,
event_count: count,
};
}
DeviceChangeEvent::ValidationPassed { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "ValidationPassed";
}
}
DeviceChangeEvent::AperakDispatched { positive, .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = if positive { "AperakSent" } else { "Rejected" };
}
}
DeviceChangeEvent::Completed { device_id } => {
if let DeviceChangeRecord::Active {
status,
device_id: d,
..
} = record
{
*status = "Completed";
*d = device_id;
}
}
DeviceChangeEvent::Rejected { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "Rejected";
}
}
DeviceChangeEvent::DeadlineExpired { .. } => {
if let DeviceChangeRecord::Active { status, .. } = record {
*status = "Rejected";
}
}
DeviceChangeEvent::IftstaStatusReceived { .. } => {
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_receive_cmd(pid: u32, validation_passed: bool) -> DeviceChangeCommand {
DeviceChangeCommand::ReceiveUtilmd {
pid: Pruefidentifikator::new(pid).expect("test pid must be in range"),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
melo_id: MeLo::new("DE0000000001234567890000000000001"),
device_id: DeviceId::new("ZHR-12345678"),
document_date: "20250115".to_owned(),
message_ref: MessageRef::new("MSG-WIM-001"),
validation_passed,
validation_errors: if validation_passed {
vec![]
} else {
vec!["AHB rule violation".to_owned()]
},
}
}
#[test]
fn happy_path_new_to_completed() {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55042, true))
.expect("should accept valid PID 55042");
assert_eq!(events.len(), 2);
assert!(
matches!(&events[0], DeviceChangeEvent::Initiated { pruefidentifikator, .. } if pruefidentifikator.as_u32() == 55042)
);
assert!(matches!(
&events[1],
DeviceChangeEvent::ValidationPassed { .. }
));
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::ValidationPassed(_)),
"expected ValidationPassed, got {}",
state.status_str()
);
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.expect("dispatch APERAK");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::AperakSent(_)),
"expected AperakSent"
);
let events = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::Complete {
device_id: DeviceId::new("ZHR-99999999"),
},
)
.expect("complete");
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::Completed(d) if d.device_id == DeviceId::new("ZHR-99999999")),
"expected Completed with new device_id",
);
}
#[test]
fn wrong_pid_is_rejected() {
let state = DeviceChangeState::default();
let err = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55001, true))
.expect_err("should reject wrong PID");
let msg = err.to_string();
assert!(
msg.contains("55001"),
"error should mention the supplied PID: {msg}"
);
}
#[test]
fn validation_failure_rejects_process() {
let state = DeviceChangeState::default();
let events = WimDeviceChangeWorkflow::handle(&state, make_receive_cmd(55042, false))
.expect("should still produce events");
assert!(matches!(&events[1], DeviceChangeEvent::Rejected { .. }));
let state = events.iter().fold(state, WimDeviceChangeWorkflow::apply);
assert!(
matches!(&state, DeviceChangeState::Rejected { .. }),
"expected Rejected"
);
}
#[test]
fn dispatch_aperak_in_wrong_state_is_rejected() {
let state = DeviceChangeState::default();
let err = WimDeviceChangeWorkflow::handle(
&state,
DeviceChangeCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.expect_err("should reject dispatch in wrong state");
assert!(err.to_string().contains("ValidationPassed"), "{err}");
}
}