use mako_engine::types::Pruefidentifikator;
use mako_engine::{
deadline::Deadline,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, PendingDeadline, Workflow, WorkflowOutput},
};
use mako_fristen::{
APERAK_STROM_WINDOW_LABEL, HolidayCalendar, aperak_strom_due_at, deadline_at_werktage,
};
pub const WORKFLOW_NAME: &str = "gpke-stammdatenaenderung";
pub const RUECKMELDUNG_WINDOW_LABEL: &str = "gpke-stammdaten-rueckmeldung-window";
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum StammdatenObjekt {
Marktlokation,
Messlokation,
Netzlokation,
SteuerbareRessource,
TechnischeRessource,
Tranche,
PaketId,
}
impl StammdatenObjekt {
#[must_use]
pub fn loc_qualifier(self) -> &'static str {
match self {
Self::Marktlokation => "Z16",
Self::Messlokation => "Z17",
Self::Netzlokation => "Z18",
Self::SteuerbareRessource => "Z19",
Self::TechnischeRessource => "Z20",
Self::Tranche => "Z21",
Self::PaketId => "Z16", }
}
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Self::Marktlokation => "MARKTLOKATION",
Self::Messlokation => "MESSLOKATION",
Self::Netzlokation => "NETZLOKATION",
Self::SteuerbareRessource => "STEUERBARE_RESSOURCE",
Self::TechnischeRessource => "TECHNISCHE_RESSOURCE",
Self::Tranche => "TRANCHE",
Self::PaketId => "PAKET_ID",
}
}
}
pub const STAMMDATEN_PAIRS: &[(u32, u32, StammdatenObjekt)] = &[
(55615, 55621, StammdatenObjekt::Netzlokation),
(55616, 55622, StammdatenObjekt::Marktlokation),
(55617, 55623, StammdatenObjekt::TechnischeRessource),
(55618, 55624, StammdatenObjekt::SteuerbareRessource),
(55619, 55625, StammdatenObjekt::Tranche),
(55620, 55626, StammdatenObjekt::Messlokation),
(55691, 55692, StammdatenObjekt::PaketId),
(55627, 55633, StammdatenObjekt::Netzlokation),
(55628, 55634, StammdatenObjekt::Marktlokation),
(55629, 55635, StammdatenObjekt::TechnischeRessource),
(55630, 55636, StammdatenObjekt::SteuerbareRessource),
(55632, 55638, StammdatenObjekt::Messlokation),
(55688, 55689, StammdatenObjekt::Marktlokation),
(55109, 55137, StammdatenObjekt::Marktlokation),
(55230, 55232, StammdatenObjekt::Netzlokation),
(55693, 55694, StammdatenObjekt::TechnischeRessource),
(55110, 55136, StammdatenObjekt::Marktlokation),
(55557, 55559, StammdatenObjekt::Marktlokation),
(55639, 55644, StammdatenObjekt::Netzlokation),
(55640, 55645, StammdatenObjekt::Marktlokation),
(55641, 55646, StammdatenObjekt::SteuerbareRessource),
(55642, 55647, StammdatenObjekt::Tranche),
(55643, 55648, StammdatenObjekt::Messlokation),
(55649, 55654, StammdatenObjekt::Netzlokation),
(55650, 55655, StammdatenObjekt::Marktlokation),
(55651, 55656, StammdatenObjekt::SteuerbareRessource),
(55652, 55657, StammdatenObjekt::Tranche),
(55653, 55658, StammdatenObjekt::Messlokation),
(55659, 55664, StammdatenObjekt::Netzlokation),
(55660, 55665, StammdatenObjekt::Marktlokation),
(55661, 55666, StammdatenObjekt::SteuerbareRessource),
(55662, 55667, StammdatenObjekt::Tranche),
(55663, 55669, StammdatenObjekt::Messlokation),
(55684, 55685, StammdatenObjekt::Marktlokation),
(55686, 55687, StammdatenObjekt::Tranche),
(55670, 55671, StammdatenObjekt::Marktlokation),
];
#[must_use]
pub fn objekt_of(aenderung_pid: u32) -> Option<StammdatenObjekt> {
STAMMDATEN_PAIRS
.iter()
.find(|(a, _, _)| *a == aenderung_pid)
.map(|(_, _, o)| *o)
}
#[must_use]
pub fn rueckmeldung_pid_for(aenderung_pid: u32) -> Option<u32> {
STAMMDATEN_PAIRS
.iter()
.find(|(a, _, _)| *a == aenderung_pid)
.map(|(_, r, _)| *r)
}
#[must_use]
pub fn is_rueckmeldung_pid(pid: u32) -> bool {
STAMMDATEN_PAIRS.iter().any(|(_, r, _)| *r == pid)
}
#[must_use]
pub fn is_aenderung_pid(pid: u32) -> bool {
STAMMDATEN_PAIRS.iter().any(|(a, _, _)| *a == pid)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum Qualitaet {
Uebernommen,
UebernommenMitKorrektur,
}
impl Qualitaet {
#[must_use]
pub fn code(self) -> &'static str {
match self {
Self::Uebernommen => "A01",
Self::UebernommenMitKorrektur => "A02",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum StammdatenEvent {
AenderungErhalten {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
pruefidentifikator: Pruefidentifikator,
objekt: StammdatenObjekt,
aenderungsdatum: String,
message_ref: MessageRef,
},
ValidationPassed {
message_ref: MessageRef,
},
RueckmeldungGesendet {
response_pid: Pruefidentifikator,
qualitaet: Qualitaet,
},
AenderungGesendet {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
pruefidentifikator: Pruefidentifikator,
objekt: StammdatenObjekt,
aenderungsdatum: String,
},
RueckmeldungErhalten {
response_pid: Pruefidentifikator,
qualitaet: Qualitaet,
},
StillschweigendAngenommen {
deadline_id: DeadlineId,
},
Rejected {
reason: String,
},
}
impl EventPayload for StammdatenEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AenderungErhalten { .. } => "StammdatenAenderungErhalten",
Self::ValidationPassed { .. } => "StammdatenValidationPassed",
Self::RueckmeldungGesendet { .. } => "StammdatenRueckmeldungGesendet",
Self::AenderungGesendet { .. } => "StammdatenAenderungGesendet",
Self::RueckmeldungErhalten { .. } => "StammdatenRueckmeldungErhalten",
Self::StillschweigendAngenommen { .. } => "StammdatenStillschweigendAngenommen",
Self::Rejected { .. } => "StammdatenRejected",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StammdatenData {
pub location_id: MaLo,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub pruefidentifikator: Pruefidentifikator,
pub objekt: StammdatenObjekt,
pub aenderungsdatum: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum StammdatenState {
#[default]
New,
Eingegangen(StammdatenData),
ValidationPassed(StammdatenData),
Beantwortet {
data: StammdatenData,
qualitaet: Qualitaet,
},
Gesendet(StammdatenData),
Abgeschlossen {
data: StammdatenData,
qualitaet: Option<Qualitaet>,
},
Rejected {
reason: String,
},
}
impl StammdatenState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Eingegangen(_) => "Eingegangen",
Self::ValidationPassed(_) => "ValidationPassed",
Self::Beantwortet { .. } => "Beantwortet",
Self::Gesendet(_) => "Gesendet",
Self::Abgeschlossen { .. } => "Abgeschlossen",
Self::Rejected { .. } => "Rejected",
}
}
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::Beantwortet { .. } | Self::Abgeschlossen { .. } | Self::Rejected { .. }
)
}
}
#[derive(Clone)]
pub enum StammdatenCommand {
ReceiveAenderung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
aenderungsdatum: String,
patch: serde_json::Value,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
received_at: time::OffsetDateTime,
},
SendRueckmeldung {
qualitaet: Qualitaet,
},
SendAenderung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
aenderungsdatum: String,
},
ReceiveRueckmeldung {
response_pid: Pruefidentifikator,
qualitaet: Qualitaet,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for StammdatenCommand {}
fn apply_outbox(data: &StammdatenData, patch: &serde_json::Value) -> PendingOutbox {
PendingOutbox::new(
"ProcessCompleted",
"",
serde_json::json!({
"pid": data.pruefidentifikator.as_u32(),
"malo_id": data.location_id.as_str(),
"objekt": data.objekt.as_str(),
"aenderungsdatum": data.aenderungsdatum,
"stammdaten_patch": patch,
}),
)
}
pub struct GpkeStammdatenaenderungWorkflow;
impl Workflow for GpkeStammdatenaenderungWorkflow {
type State = StammdatenState;
type Event = StammdatenEvent;
type Command = StammdatenCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(
RUECKMELDUNG_WINDOW_LABEL | APERAK_STROM_WINDOW_LABEL,
StammdatenState::Eingegangen(_)
| StammdatenState::ValidationPassed(_)
| StammdatenState::Gesendet(_),
) => 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::AenderungErhalten {
location_id,
sender,
receiver,
pruefidentifikator,
objekt,
aenderungsdatum,
..
} => StammdatenState::Eingegangen(StammdatenData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
pruefidentifikator: *pruefidentifikator,
objekt: *objekt,
aenderungsdatum: aenderungsdatum.clone(),
}),
StammdatenEvent::ValidationPassed { .. } => match state {
StammdatenState::Eingegangen(data) => StammdatenState::ValidationPassed(data),
other => other,
},
StammdatenEvent::RueckmeldungGesendet { qualitaet, .. } => match state {
StammdatenState::ValidationPassed(data) => StammdatenState::Beantwortet {
data,
qualitaet: *qualitaet,
},
other => other,
},
StammdatenEvent::AenderungGesendet {
location_id,
sender,
receiver,
pruefidentifikator,
objekt,
aenderungsdatum,
} => StammdatenState::Gesendet(StammdatenData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
pruefidentifikator: *pruefidentifikator,
objekt: *objekt,
aenderungsdatum: aenderungsdatum.clone(),
}),
StammdatenEvent::RueckmeldungErhalten { qualitaet, .. } => match state {
StammdatenState::Gesendet(data) => StammdatenState::Abgeschlossen {
data,
qualitaet: Some(*qualitaet),
},
other => other,
},
StammdatenEvent::StillschweigendAngenommen { .. } => match state {
StammdatenState::Gesendet(data) => StammdatenState::Abgeschlossen {
data,
qualitaet: None,
},
StammdatenState::Eingegangen(data) | StammdatenState::ValidationPassed(data) => {
StammdatenState::Beantwortet {
data,
qualitaet: Qualitaet::Uebernommen,
}
}
other => other,
},
StammdatenEvent::Rejected { reason } => StammdatenState::Rejected {
reason: reason.clone(),
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
StammdatenCommand::ReceiveAenderung {
pid,
sender,
receiver,
location_id,
aenderungsdatum,
patch,
message_ref,
validation_passed,
validation_errors,
received_at,
} => {
if !matches!(state, StammdatenState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
let Some(objekt) = objekt_of(pid.as_u32()) else {
return Err(WorkflowError::rejected(format!(
"PID {pid} is not a Stammdatenänderung Änderung PID",
)));
};
let sender_mp_id = sender.clone();
let receiver_gln = receiver.clone();
let mut events = vec![StammdatenEvent::AenderungErhalten {
location_id: location_id.clone(),
sender,
receiver,
pruefidentifikator: pid,
objekt,
aenderungsdatum: aenderungsdatum.clone(),
message_ref: message_ref.clone(),
}];
if !validation_passed {
let reason = if validation_errors.is_empty() {
"structural validation failed".to_owned()
} else {
validation_errors.join("; ")
};
events.push(StammdatenEvent::Rejected {
reason: reason.clone(),
});
let outbox = vec![
PendingOutbox::aperak_fehler(
receiver_gln.as_str(),
sender_mp_id.as_str(),
message_ref.as_str(),
mako_engine::erc::codes::Z29,
reason,
)
.caused_by(0),
];
return Ok(WorkflowOutput::with_outbox(events, outbox));
}
events.push(StammdatenEvent::ValidationPassed {
message_ref: message_ref.clone(),
});
let mut outbox = vec![
PendingOutbox::aperak_anerkennung(
receiver_gln.as_str(),
sender_mp_id.as_str(),
message_ref.as_str(),
)
.caused_by(1),
];
if patch.as_object().is_some_and(|o| !o.is_empty()) {
let data = StammdatenData {
location_id,
sender: sender_mp_id,
receiver: receiver_gln,
pruefidentifikator: pid,
objekt,
aenderungsdatum,
};
outbox.push(apply_outbox(&data, &patch).caused_by(1));
}
let deadlines = vec![
PendingDeadline::new(
APERAK_STROM_WINDOW_LABEL,
aperak_strom_due_at(received_at),
),
PendingDeadline::new(
RUECKMELDUNG_WINDOW_LABEL,
deadline_at_werktage(
received_at,
mako_fristen::antwort::STAMMDATEN_RUECKMELDUNG_WERKTAGE,
HolidayCalendar::BdewMaKo,
),
),
];
Ok(WorkflowOutput::with_outbox_and_deadlines(
events, outbox, deadlines,
))
}
StammdatenCommand::SendRueckmeldung { qualitaet } => {
let data = match state {
StammdatenState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.label(),
));
}
};
let response_pid = rueckmeldung_pid_for(data.pruefidentifikator.as_u32())
.and_then(|p| Pruefidentifikator::new(p).ok())
.ok_or_else(|| {
WorkflowError::rejected(format!(
"no Rückmeldung PID for Änderung {}",
data.pruefidentifikator
))
})?;
let outbox = vec![PendingOutbox::new(
"UTILMD",
data.sender.as_str(),
serde_json::json!({
"direction": "outbound",
"pid": response_pid.as_u32(),
"sender": data.receiver.as_str(),
"receiver": data.sender.as_str(),
"malo": data.location_id.as_str(),
"process_date": data.aenderungsdatum,
"qualitaet": qualitaet.code(),
}),
)];
Ok(WorkflowOutput::with_outbox(
vec![StammdatenEvent::RueckmeldungGesendet {
response_pid,
qualitaet,
}],
outbox,
))
}
StammdatenCommand::SendAenderung {
pid,
sender,
receiver,
location_id,
aenderungsdatum,
} => {
if !matches!(state, StammdatenState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
let Some(objekt) = objekt_of(pid.as_u32()) else {
return Err(WorkflowError::rejected(format!(
"PID {pid} is not a Stammdatenänderung Änderung PID",
)));
};
let utilmd = PendingOutbox::new(
"UTILMD",
receiver.as_str(),
serde_json::json!({
"direction": "outbound",
"pid": pid.as_u32(),
"sender": sender.as_str(),
"receiver": receiver.as_str(),
"malo": location_id.as_str(),
"process_date": aenderungsdatum,
}),
);
Ok(WorkflowOutput::with_outbox(
vec![StammdatenEvent::AenderungGesendet {
location_id,
sender,
receiver,
pruefidentifikator: pid,
objekt,
aenderungsdatum,
}],
vec![utilmd],
))
}
StammdatenCommand::ReceiveRueckmeldung {
response_pid,
qualitaet,
} => {
if !matches!(state, StammdatenState::Gesendet(_)) {
return Err(WorkflowError::invalid_state("Gesendet", state.label()));
}
if !is_rueckmeldung_pid(response_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {response_pid} is not a Stammdatenänderung Rückmeldung PID",
)));
}
Ok(vec![StammdatenEvent::RueckmeldungErhalten {
response_pid,
qualitaet,
}]
.into())
}
StammdatenCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
Ok(vec![].into())
} else if label.as_ref() == RUECKMELDUNG_WINDOW_LABEL {
Ok(vec![StammdatenEvent::StillschweigendAngenommen { deadline_id }].into())
} else {
Ok(vec![].into())
}
}
}
}
}
#[cfg(test)]
mod tests {
use mako_engine::{ids::DeadlineId, workflow::Workflow};
use super::*;
fn pid(code: u32) -> Pruefidentifikator {
Pruefidentifikator::new(code).unwrap()
}
fn mcod(s: &str) -> MarktpartnerCode {
MarktpartnerCode::new(s)
}
fn malo(s: &str) -> MaLo {
MaLo::new(s)
}
fn mref(s: &str) -> MessageRef {
MessageRef::new(s)
}
fn now() -> time::OffsetDateTime {
time::macros::datetime!(2026-07-01 10:00:00 UTC)
}
fn receive_malo_cmd(patch: serde_json::Value, ok: bool) -> StammdatenCommand {
StammdatenCommand::ReceiveAenderung {
pid: pid(55616), sender: mcod("9900357000004"),
receiver: mcod("9900357000011"),
location_id: malo("51238696781"),
aenderungsdatum: "20260801".to_owned(),
patch,
message_ref: mref("SD-001"),
validation_passed: ok,
validation_errors: if ok {
vec![]
} else {
vec!["missing mandatory segment".to_owned()]
},
received_at: now(),
}
}
fn apply_all(init: StammdatenState, events: &[StammdatenEvent]) -> StammdatenState {
events
.iter()
.fold(init, GpkeStammdatenaenderungWorkflow::apply)
}
#[test]
fn pid_classification_and_pairing() {
assert_eq!(objekt_of(55616), Some(StammdatenObjekt::Marktlokation));
assert_eq!(objekt_of(55620), Some(StammdatenObjekt::Messlokation));
assert_eq!(rueckmeldung_pid_for(55616), Some(55622));
assert_eq!(rueckmeldung_pid_for(55109), Some(55137));
assert!(is_aenderung_pid(55616));
assert!(is_rueckmeldung_pid(55622));
assert!(!is_rueckmeldung_pid(55616));
assert_eq!(rueckmeldung_pid_for(55230), Some(55232));
assert_eq!(objekt_of(55230), Some(StammdatenObjekt::Netzlokation));
assert_eq!(rueckmeldung_pid_for(55557), Some(55559));
assert_eq!(objekt_of(55557), Some(StammdatenObjekt::Marktlokation));
assert!(objekt_of(21047).is_none());
}
#[test]
fn pairs_table_has_no_duplicate_pids() {
let mut seen = std::collections::HashSet::new();
for (a, r, _) in STAMMDATEN_PAIRS {
assert!(seen.insert(*a), "duplicate Änderung PID {a}");
assert!(seen.insert(*r), "duplicate Rückmeldung PID {r}");
}
}
#[test]
fn malo_change_applied_and_answered_a01() {
let patch = serde_json::json!({ "bilanzierungsmethode": "RLM", "netzebene": "NSP" });
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
receive_malo_cmd(patch, true),
)
.unwrap();
assert_eq!(out.outbox.len(), 2);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
assert_eq!(out.outbox[1].message_type.as_ref(), "ProcessCompleted");
assert_eq!(out.outbox[1].payload["objekt"], "MARKTLOKATION");
assert_eq!(
out.outbox[1].payload["stammdaten_patch"]["bilanzierungsmethode"],
"RLM"
);
assert_eq!(out.deadlines.len(), 2); let state = apply_all(StammdatenState::New, &out.events);
assert!(matches!(state, StammdatenState::ValidationPassed(_)));
let out = GpkeStammdatenaenderungWorkflow::handle(
&state,
StammdatenCommand::SendRueckmeldung {
qualitaet: Qualitaet::Uebernommen,
},
)
.unwrap();
assert_eq!(out.outbox[0].payload["pid"], 55622);
assert_eq!(out.outbox[0].payload["qualitaet"], "A01");
let state = apply_all(state, &out.events);
assert!(matches!(
state,
StammdatenState::Beantwortet {
qualitaet: Qualitaet::Uebernommen,
..
}
));
}
#[test]
fn empty_patch_acknowledged_without_apply_outbox() {
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
receive_malo_cmd(serde_json::json!({}), true),
)
.unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
}
#[test]
fn nelo_change_emits_object_tagged_apply() {
let mut cmd = receive_malo_cmd(serde_json::json!({ "netzebene": "NSP" }), true);
if let StammdatenCommand::ReceiveAenderung { pid: p, .. } = &mut cmd {
*p = pid(55615); }
let out = GpkeStammdatenaenderungWorkflow::handle(&StammdatenState::New, cmd).unwrap();
assert_eq!(out.outbox.len(), 2);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
assert_eq!(out.outbox[1].message_type.as_ref(), "ProcessCompleted");
assert_eq!(out.outbox[1].payload["objekt"], "NETZLOKATION");
assert_eq!(
out.outbox[1].payload["stammdaten_patch"]["netzebene"],
"NSP"
);
}
#[test]
fn non_malo_object_without_grounded_attributes_is_acknowledged_only() {
let mut cmd = receive_malo_cmd(serde_json::json!({}), true);
if let StammdatenCommand::ReceiveAenderung { pid: p, .. } = &mut cmd {
*p = pid(55618); }
let out = GpkeStammdatenaenderungWorkflow::handle(&StammdatenState::New, cmd).unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
}
#[test]
fn a02_reports_correction() {
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
receive_malo_cmd(serde_json::json!({ "regelzone": "10YDE-EON------1" }), true),
)
.unwrap();
let state = apply_all(StammdatenState::New, &out.events);
let out = GpkeStammdatenaenderungWorkflow::handle(
&state,
StammdatenCommand::SendRueckmeldung {
qualitaet: Qualitaet::UebernommenMitKorrektur,
},
)
.unwrap();
assert_eq!(out.outbox[0].payload["qualitaet"], "A02");
}
#[test]
fn validation_failure_rejects_with_aperak_313() {
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
receive_malo_cmd(serde_json::json!({}), false),
)
.unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
assert_eq!(out.outbox[0].payload["error_code"], "Z29");
let state = apply_all(StammdatenState::New, &out.events);
assert!(matches!(state, StammdatenState::Rejected { .. }));
}
#[test]
fn rueckmeldung_timeout_is_tacit_acceptance() {
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
receive_malo_cmd(serde_json::json!({ "netzebene": "NSP" }), true),
)
.unwrap();
let state = apply_all(StammdatenState::New, &out.events);
let out = GpkeStammdatenaenderungWorkflow::handle(
&state,
StammdatenCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: RUECKMELDUNG_WINDOW_LABEL.into(),
},
)
.unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(
state,
StammdatenState::Beantwortet {
qualitaet: Qualitaet::Uebernommen,
..
}
));
}
#[test]
fn initiator_send_and_receive_rueckmeldung() {
let out = GpkeStammdatenaenderungWorkflow::handle(
&StammdatenState::New,
StammdatenCommand::SendAenderung {
pid: pid(55109), sender: mcod("9900357000011"),
receiver: mcod("9900357000004"),
location_id: malo("51238696781"),
aenderungsdatum: "20260801".to_owned(),
},
)
.unwrap();
assert_eq!(out.outbox[0].message_type.as_ref(), "UTILMD");
assert_eq!(out.outbox[0].payload["pid"], 55109);
let state = apply_all(StammdatenState::New, &out.events);
assert!(matches!(state, StammdatenState::Gesendet(_)));
let out = GpkeStammdatenaenderungWorkflow::handle(
&state,
StammdatenCommand::ReceiveRueckmeldung {
response_pid: pid(55137),
qualitaet: Qualitaet::Uebernommen,
},
)
.unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(state, StammdatenState::Abgeschlossen { .. }));
}
#[test]
fn wrong_pid_rejected() {
let mut cmd = receive_malo_cmd(serde_json::json!({}), true);
if let StammdatenCommand::ReceiveAenderung { pid: p, .. } = &mut cmd {
*p = pid(55001); }
assert!(GpkeStammdatenaenderungWorkflow::handle(&StammdatenState::New, cmd).is_err());
}
}