use mako_engine::types::Pruefidentifikator;
use mako_engine::{
deadline::Deadline,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "gpke-beendigung-zuordnung";
pub const BEENDIGUNG_ZUORDNUNG_PIDS: &[u32] = &[55010, 55011, 55012];
pub const ANFRAGE_PID: u32 = 55_010;
pub const ANTWORT_PIDS: &[u32] = &[55_011, 55_012];
pub const NB_ANFRAGE_WINDOW_LABEL: &str = "gpke-beendigung-zuordnung-lfa-antwort";
pub const BEENDIGUNG_ZUORDNUNG_ANTWORT_WINDOW_LABEL: &str =
"gpke-beendigung-zuordnung-antwortfrist";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum BeendigungZuordnungEvent {
AnfrageErhalten {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
document_date: String,
process_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
vorgangsnummer: Option<String>,
},
ValidationPassed {
message_ref: MessageRef,
},
AntwortGesendet {
response_pid: Pruefidentifikator,
accepted: bool,
antwort: crate::lf_antwort::LfAntwort,
},
Beendet,
AperakFehlerDispatched {
aperak_pid: Pruefidentifikator,
reason: String,
outbound_ref: MessageRef,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
AnfrageGesendet {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
process_date: String,
vorgangsnummer: String,
anmeldung_process_id: String,
},
LfaAntwortErhalten {
response_pid: Option<Pruefidentifikator>,
antwortcode: Option<String>,
zustimmung: bool,
grund: Option<String>,
zuordnungsende: Option<String>,
fristablauf: bool,
},
}
impl EventPayload for BeendigungZuordnungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnfrageErhalten { .. } => "BeendigungZuordnungAnfrageErhalten",
Self::ValidationPassed { .. } => "BeendigungZuordnungValidationPassed",
Self::AntwortGesendet { .. } => "BeendigungZuordnungAntwortGesendet",
Self::Beendet => "BeendigungZuordnungBeendet",
Self::AperakFehlerDispatched { .. } => "BeendigungZuordnungAperakFehlerDispatched",
Self::Rejected { .. } => "BeendigungZuordnungRejected",
Self::DeadlineExpired { .. } => "BeendigungZuordnungDeadlineExpired",
Self::AnfrageGesendet { .. } => "BeendigungZuordnungAnfrageGesendet",
Self::LfaAntwortErhalten { .. } => "BeendigungZuordnungLfaAntwortErhalten",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct BeendigungZuordnungData {
pub location_id: MaLo,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub document_date: String,
pub process_date: String,
pub pruefidentifikator: Pruefidentifikator,
pub vorgangsnummer: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum BeendigungZuordnungState {
#[default]
New,
Eingegangen(BeendigungZuordnungData),
ValidationPassed(BeendigungZuordnungData),
AntwortGesendet {
data: BeendigungZuordnungData,
response_pid: Pruefidentifikator,
},
Beendet(BeendigungZuordnungData),
Rejected {
reason: String,
},
AnfrageGesendet {
data: BeendigungZuordnungData,
anmeldung_process_id: String,
},
LfaAntwort {
data: BeendigungZuordnungData,
anmeldung_process_id: String,
zustimmung: bool,
antwortcode: Option<String>,
},
}
impl BeendigungZuordnungState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Eingegangen(_) => "Eingegangen",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AntwortGesendet { .. } => "AntwortGesendet",
Self::Beendet(_) => "Beendet",
Self::Rejected { .. } => "Rejected",
Self::AnfrageGesendet { .. } => "AnfrageGesendet",
Self::LfaAntwort { .. } => "LfaAntwort",
}
}
#[must_use]
pub fn data(&self) -> Option<&BeendigungZuordnungData> {
match self {
Self::Eingegangen(d) | Self::ValidationPassed(d) | Self::Beendet(d) => Some(d),
Self::AntwortGesendet { data, .. }
| Self::AnfrageGesendet { data, .. }
| Self::LfaAntwort { data, .. } => Some(data),
Self::New | Self::Rejected { .. } => None,
}
}
}
#[derive(Clone)]
pub enum BeendigungZuordnungCommand {
ReceiveAnfrage {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
document_date: String,
process_date: String,
message_ref: MessageRef,
vorgang: Box<crate::lf_antwort::LfVorgangsdaten>,
validation_passed: bool,
validation_errors: Vec<String>,
},
SendAntwort {
antwort: crate::lf_antwort::LfAntwort,
},
Anfragen {
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
tranche: bool,
process_date: String,
vorgangsnummer: String,
anmeldung_process_id: String,
kunde_name: Option<String>,
lfn_mp_id: Option<String>,
},
ReceiveAntwort {
response_pid: Pruefidentifikator,
antwortcode: String,
zustimmung: bool,
grund: Option<String>,
zuordnungsende: Option<String>,
},
AntwortfristAbgelaufen,
BeendenBestaetigen,
DispatchAperakFehler {
reason: String,
outbound_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for BeendigungZuordnungCommand {}
pub struct GpkeBeendigungZuordnungWorkflow;
impl Workflow for GpkeBeendigungZuordnungWorkflow {
type State = BeendigungZuordnungState;
type Event = BeendigungZuordnungEvent;
type Command = BeendigungZuordnungCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(
BEENDIGUNG_ZUORDNUNG_ANTWORT_WINDOW_LABEL,
BeendigungZuordnungState::Eingegangen(_)
| BeendigungZuordnungState::ValidationPassed(_),
) => Some(BeendigungZuordnungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
(NB_ANFRAGE_WINDOW_LABEL, BeendigungZuordnungState::AnfrageGesendet { .. }) => {
Some(BeendigungZuordnungCommand::AntwortfristAbgelaufen)
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
BeendigungZuordnungEvent::AnfrageErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
pruefidentifikator,
vorgangsnummer,
..
} => BeendigungZuordnungState::Eingegangen(BeendigungZuordnungData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
document_date: document_date.clone(),
process_date: process_date.clone(),
pruefidentifikator: *pruefidentifikator,
vorgangsnummer: vorgangsnummer.clone(),
}),
BeendigungZuordnungEvent::ValidationPassed { .. } => match state {
BeendigungZuordnungState::Eingegangen(data) => {
BeendigungZuordnungState::ValidationPassed(data)
}
other => other,
},
BeendigungZuordnungEvent::AntwortGesendet {
accepted,
response_pid,
..
} => {
if *accepted {
match state {
BeendigungZuordnungState::ValidationPassed(data) => {
BeendigungZuordnungState::AntwortGesendet {
response_pid: *response_pid,
data,
}
}
other => other,
}
} else {
BeendigungZuordnungState::Rejected {
reason: "Anfrage abgelehnt".to_owned(),
}
}
}
BeendigungZuordnungEvent::Beendet => match state {
BeendigungZuordnungState::AntwortGesendet { data, .. } => {
BeendigungZuordnungState::Beendet(data)
}
other => other,
},
BeendigungZuordnungEvent::AperakFehlerDispatched { reason, .. } => {
BeendigungZuordnungState::Rejected {
reason: format!("APERAK 29001: {reason}"),
}
}
BeendigungZuordnungEvent::AnfrageGesendet {
location_id,
sender,
receiver,
process_date,
vorgangsnummer,
anmeldung_process_id,
} => BeendigungZuordnungState::AnfrageGesendet {
data: BeendigungZuordnungData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
document_date: process_date.clone(),
process_date: process_date.clone(),
pruefidentifikator: Pruefidentifikator::new(ANFRAGE_PID)
.unwrap_or_else(|_| unreachable!("55010 is a valid Prüfidentifikator")),
vorgangsnummer: Some(vorgangsnummer.clone()),
},
anmeldung_process_id: anmeldung_process_id.clone(),
},
BeendigungZuordnungEvent::LfaAntwortErhalten {
zustimmung,
antwortcode,
..
} => match state {
BeendigungZuordnungState::AnfrageGesendet {
data,
anmeldung_process_id,
} => BeendigungZuordnungState::LfaAntwort {
data,
anmeldung_process_id,
zustimmung: *zustimmung,
antwortcode: antwortcode.clone(),
},
other => other,
},
BeendigungZuordnungEvent::Rejected { reason } => BeendigungZuordnungState::Rejected {
reason: reason.clone(),
},
BeendigungZuordnungEvent::DeadlineExpired { label, .. } => match state {
BeendigungZuordnungState::Beendet(_)
| BeendigungZuordnungState::Rejected { .. } => state,
_ => BeendigungZuordnungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
BeendigungZuordnungCommand::ReceiveAnfrage {
pid,
sender,
receiver,
location_id,
document_date,
process_date,
message_ref,
vorgang,
validation_passed,
validation_errors,
} => {
if !matches!(state, BeendigungZuordnungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !BEENDIGUNG_ZUORDNUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected Anfrage zur Beendigung der Zuordnung PID (55010), got {pid}",
)));
}
let sender_mp_id = sender.clone();
let receiver_gln = receiver.clone();
let notify_malo = location_id.clone();
let notify_termin = process_date.clone();
let mut events = vec![BeendigungZuordnungEvent::AnfrageErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
vorgangsnummer: vorgang.vorgangsnummer.clone(),
}];
if validation_passed {
events.push(BeendigungZuordnungEvent::ValidationPassed {
message_ref: message_ref.clone(),
});
let outbox = vec![
vorgang
.process_initiated(
pid,
¬ify_malo,
&sender_mp_id,
&receiver_gln,
¬ify_termin,
&serde_json::Value::Null,
)
.caused_by(1),
PendingOutbox::aperak_anerkennung(
receiver_gln.as_str(),
sender_mp_id.as_str(),
message_ref.as_str(),
)
.caused_by(1),
];
Ok(WorkflowOutput::with_outbox(events, outbox))
} else {
let reason = validation_errors.join("; ");
events.push(BeendigungZuordnungEvent::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),
];
Ok(WorkflowOutput::with_outbox(events, outbox))
}
}
BeendigungZuordnungCommand::SendAntwort { antwort } => {
let data = match state {
BeendigungZuordnungState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.label(),
));
}
};
let accepted = antwort.zustimmung;
let response_code: u32 = if accepted { 55011 } else { 55012 };
let response_pid = Pruefidentifikator::new(response_code)
.map_err(|e| WorkflowError::rejected(e.clone()))?;
let outbox = vec![
crate::lf_antwort::antwort_outbox(
response_code,
&antwort,
&data.location_id,
&data.sender,
&data.receiver,
&data.process_date,
data.vorgangsnummer.as_deref(),
)
.caused_by(0),
];
Ok(WorkflowOutput::with_outbox(
vec![BeendigungZuordnungEvent::AntwortGesendet {
response_pid,
accepted,
antwort,
}],
outbox,
))
}
BeendigungZuordnungCommand::Anfragen {
sender,
receiver,
location_id,
tranche,
process_date,
vorgangsnummer,
anmeldung_process_id,
kunde_name,
lfn_mp_id,
} => {
if !matches!(state, BeendigungZuordnungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if process_date.trim().is_empty() {
return Err(WorkflowError::rejected(
"the Anfrage zur Beendigung der Zuordnung names the Zuordnungsende it \
asks for (SG4 DTM+93); UTILMD AHB Strom marks it Muss on 55010"
.to_owned(),
));
}
let mut payload = serde_json::json!({
"direction": "outbound",
"pid": ANFRAGE_PID,
"sender": sender.as_str(),
"receiver": receiver.as_str(),
"malo": location_id.as_str(),
"process_date": process_date,
"vorgangsnummer": vorgangsnummer,
"document_code": "E02",
"lokationstyp": if tranche { "Z21" } else { "Z16" },
});
let obj = payload.as_object_mut().expect("json! built an object");
if let Some(name) = kunde_name.filter(|n| !n.is_empty()) {
obj.insert("kunde_name".into(), name.into());
}
if let Some(lfn) = lfn_mp_id.filter(|m| !m.is_empty()) {
obj.insert(
"beteiligte_marktpartner".into(),
serde_json::Value::Array(vec![serde_json::Value::String(lfn)]),
);
}
let outbox = vec![PendingOutbox::new("UTILMD", receiver.as_str(), payload)];
Ok(WorkflowOutput::with_outbox(
vec![BeendigungZuordnungEvent::AnfrageGesendet {
location_id,
sender,
receiver,
process_date,
vorgangsnummer,
anmeldung_process_id,
}],
outbox,
))
}
BeendigungZuordnungCommand::ReceiveAntwort {
response_pid,
antwortcode,
zustimmung,
grund,
zuordnungsende,
} => {
let BeendigungZuordnungState::AnfrageGesendet {
data,
anmeldung_process_id,
} = state
else {
return Err(WorkflowError::invalid_state(
"AnfrageGesendet",
state.label(),
));
};
if !ANTWORT_PIDS.contains(&response_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected an Antwort auf die Anfrage zur Beendigung der Zuordnung \
({ANTWORT_PIDS:?}), got {response_pid}",
)));
}
Ok(WorkflowOutput::with_outbox(
vec![BeendigungZuordnungEvent::LfaAntwortErhalten {
response_pid: Some(response_pid),
antwortcode: Some(antwortcode.clone()),
zustimmung,
grund: grund.clone(),
zuordnungsende: zuordnungsende.clone(),
fristablauf: false,
}],
vec![lfa_antwort_notification(
data,
anmeldung_process_id,
Some(&antwortcode),
zustimmung,
grund.as_deref(),
zuordnungsende.as_deref(),
false,
)],
))
}
BeendigungZuordnungCommand::AntwortfristAbgelaufen => {
let BeendigungZuordnungState::AnfrageGesendet {
data,
anmeldung_process_id,
} = state
else {
return Ok(vec![].into());
};
Ok(WorkflowOutput::with_outbox(
vec![BeendigungZuordnungEvent::LfaAntwortErhalten {
response_pid: None,
antwortcode: None,
zustimmung: true,
grund: None,
zuordnungsende: None,
fristablauf: true,
}],
vec![lfa_antwort_notification(
data,
anmeldung_process_id,
None,
true,
None,
None,
true,
)],
))
}
BeendigungZuordnungCommand::BeendenBestaetigen => {
if !matches!(state, BeendigungZuordnungState::AntwortGesendet { .. }) {
return Err(WorkflowError::invalid_state(
"AntwortGesendet",
state.label(),
));
}
Ok(vec![BeendigungZuordnungEvent::Beendet].into())
}
BeendigungZuordnungCommand::DispatchAperakFehler {
reason,
outbound_ref,
} => {
match state {
BeendigungZuordnungState::Eingegangen(_)
| BeendigungZuordnungState::ValidationPassed(_) => {}
_ => {
return Err(WorkflowError::invalid_state(
"Eingegangen or ValidationPassed",
state.label(),
));
}
}
let aperak_pid = Pruefidentifikator::new(29_001)
.map_err(|e| WorkflowError::rejected(e.clone()))?;
Ok(vec![BeendigungZuordnungEvent::AperakFehlerDispatched {
aperak_pid,
reason,
outbound_ref,
}]
.into())
}
BeendigungZuordnungCommand::TimeoutExpired { deadline_id, label } => match state {
BeendigungZuordnungState::Beendet(_)
| BeendigungZuordnungState::Rejected { .. } => Ok(vec![].into()),
_ => Ok(
vec![BeendigungZuordnungEvent::DeadlineExpired { deadline_id, label }].into(),
),
},
}
}
}
#[allow(clippy::too_many_arguments)]
fn lfa_antwort_notification(
data: &BeendigungZuordnungData,
anmeldung_process_id: &str,
antwortcode: Option<&str>,
zustimmung: bool,
grund: Option<&str>,
zuordnungsende: Option<&str>,
fristablauf: bool,
) -> PendingOutbox {
let mut payload = serde_json::json!({
"type": "LfaAntwortAufAbmeldeanfrage",
"pid": ANFRAGE_PID,
"malo_id": data.location_id.as_str(),
"grid_operator": data.sender.as_str(),
"lfa_mp_id": data.receiver.as_str(),
"anmeldung_process_id": anmeldung_process_id,
"zustimmung": zustimmung,
"fristablauf": fristablauf,
});
let obj = payload.as_object_mut().expect("json! built an object");
if let Some(c) = antwortcode {
obj.insert("antwortcode".into(), c.into());
}
if let Some(g) = grund {
obj.insert("grund".into(), g.into());
}
if let Some(ende) = zuordnungsende {
obj.insert("zuordnungsende".into(), ende.into());
}
PendingOutbox::new("LfaAntwortAufAbmeldeanfrage", data.sender.as_str(), payload)
}
#[cfg(test)]
mod tests {
use mako_engine::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 anfrage_cmd(ok: bool) -> BeendigungZuordnungCommand {
BeendigungZuordnungCommand::ReceiveAnfrage {
pid: pid(55010),
sender: mcod("9900357000004"),
receiver: mcod("4012345000023"),
location_id: malo("51238696781"),
document_date: "20251001".to_owned(),
process_date: "20260101".to_owned(),
message_ref: mref("BEEND-001"),
vorgang: Box::new(crate::LfVorgangsdaten::default()),
validation_passed: ok,
validation_errors: if ok {
vec![]
} else {
vec!["missing mandatory segment".to_owned()]
},
}
}
fn apply_all(
init: BeendigungZuordnungState,
events: &[BeendigungZuordnungEvent],
) -> BeendigungZuordnungState {
events
.iter()
.fold(init, GpkeBeendigungZuordnungWorkflow::apply)
}
#[test]
fn happy_path_bestaetigung() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&BeendigungZuordnungState::New,
anfrage_cmd(true),
)
.unwrap();
assert_eq!(out.events.len(), 2);
assert_eq!(out.outbox.len(), 2);
assert_eq!(out.outbox[0].message_type.as_ref(), "ProcessInitiated");
assert_eq!(out.outbox[1].payload["pid"], 29002);
assert!(
out.outbox[1].payload["orig_message_ref"].is_string(),
"SG2 RFF+ACE is Muss on 29002 — an APERAK without the reference \
cannot be rendered",
);
let state = apply_all(BeendigungZuordnungState::New, &out.events);
assert!(matches!(
state,
BeendigungZuordnungState::ValidationPassed(_)
));
let out = GpkeBeendigungZuordnungWorkflow::handle(
&state,
BeendigungZuordnungCommand::SendAntwort {
antwort: crate::lf_antwort::LfAntwort::zustimmung("A36", "E_0624"),
},
)
.unwrap();
if let BeendigungZuordnungEvent::AntwortGesendet { response_pid, .. } = &out.events[0] {
assert_eq!(response_pid.as_u32(), 55011);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
let out = GpkeBeendigungZuordnungWorkflow::handle(
&state,
BeendigungZuordnungCommand::BeendenBestaetigen,
)
.unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(state, BeendigungZuordnungState::Beendet(_)));
}
#[test]
fn ablehnung_yields_55012() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&BeendigungZuordnungState::New,
anfrage_cmd(true),
)
.unwrap();
let state = apply_all(BeendigungZuordnungState::New, &out.events);
let out = GpkeBeendigungZuordnungWorkflow::handle(
&state,
BeendigungZuordnungCommand::SendAntwort {
antwort: crate::lf_antwort::LfAntwort::ablehnung("A35", "E_0624")
.with_bemerkung("Widerspruch"),
},
)
.unwrap();
if let BeendigungZuordnungEvent::AntwortGesendet { response_pid, .. } = &out.events[0] {
assert_eq!(response_pid.as_u32(), 55012);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
assert!(matches!(state, BeendigungZuordnungState::Rejected { .. }));
}
#[test]
fn validation_failure_emits_aperak_313() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&BeendigungZuordnungState::New,
anfrage_cmd(false),
)
.unwrap();
assert_eq!(out.outbox[0].payload["error_code"], "Z29");
let state = apply_all(BeendigungZuordnungState::New, &out.events);
assert!(matches!(state, BeendigungZuordnungState::Rejected { .. }));
}
#[test]
fn wrong_pid_rejected() {
let mut cmd = anfrage_cmd(true);
if let BeendigungZuordnungCommand::ReceiveAnfrage { pid: p, .. } = &mut cmd {
*p = pid(55001);
}
assert!(
GpkeBeendigungZuordnungWorkflow::handle(&BeendigungZuordnungState::New, cmd).is_err()
);
}
}
#[cfg(test)]
mod nb_initiator_tests {
use super::*;
fn mp(v: &str) -> MarktpartnerCode {
MarktpartnerCode::new(v.to_owned())
}
fn anfragen() -> BeendigungZuordnungCommand {
BeendigungZuordnungCommand::Anfragen {
sender: mp("9900357000004"),
receiver: mp("9900111000002"),
location_id: MaLo::new("51238696781".to_owned()),
tranche: false,
process_date: "20261101".to_owned(),
vorgangsnummer: "ANF-1".to_owned(),
anmeldung_process_id: "11111111-1111-1111-1111-111111111111".to_owned(),
kunde_name: Some("Mustermann".to_owned()),
lfn_mp_id: Some("9900555000005".to_owned()),
}
}
fn after_anfrage() -> BeendigungZuordnungState {
let out =
GpkeBeendigungZuordnungWorkflow::handle(&BeendigungZuordnungState::New, anfragen())
.expect("Anfragen accepted");
out.events
.iter()
.fold(BeendigungZuordnungState::New, |st, e| {
GpkeBeendigungZuordnungWorkflow::apply(st, e)
})
}
#[test]
fn the_nb_renders_the_anfrage_with_its_sg12_parties() {
let out =
GpkeBeendigungZuordnungWorkflow::handle(&BeendigungZuordnungState::New, anfragen())
.expect("Anfragen accepted");
let p = &out.outbox[0].payload;
assert_eq!(&*out.outbox[0].message_type, "UTILMD");
assert_eq!(p["pid"], 55_010);
assert_eq!(p["document_code"], "E02");
assert_eq!(p["process_date"], "20261101");
assert_eq!(p["kunde_name"], "Mustermann");
assert_eq!(p["beteiligte_marktpartner"][0], "9900555000005");
}
#[test]
fn a_tranche_anfrage_names_the_tranche_qualifier() {
let BeendigungZuordnungCommand::Anfragen {
sender,
receiver,
location_id,
process_date,
vorgangsnummer,
anmeldung_process_id,
kunde_name,
lfn_mp_id,
..
} = anfragen()
else {
unreachable!()
};
let out = GpkeBeendigungZuordnungWorkflow::handle(
&BeendigungZuordnungState::New,
BeendigungZuordnungCommand::Anfragen {
sender,
receiver,
location_id,
tranche: true,
process_date,
vorgangsnummer,
anmeldung_process_id,
kunde_name,
lfn_mp_id,
},
)
.expect("accepted");
assert_eq!(out.outbox[0].payload["lokationstyp"], "Z21");
}
#[test]
fn a_lapsed_window_is_a_zustimmung_not_a_timeout() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&after_anfrage(),
BeendigungZuordnungCommand::AntwortfristAbgelaufen,
)
.expect("lapse accepted");
let BeendigungZuordnungEvent::LfaAntwortErhalten {
zustimmung,
fristablauf,
antwortcode,
..
} = &out.events[0]
else {
panic!("expected LfaAntwortErhalten, got {:?}", out.events[0]);
};
assert!(zustimmung, "silence releases the Marktlokation");
assert!(fristablauf);
assert!(antwortcode.is_none(), "the LFA named no code");
assert_eq!(&*out.outbox[0].message_type, "LfaAntwortAufAbmeldeanfrage");
assert_eq!(out.outbox[0].payload["zustimmung"], true);
}
#[test]
fn the_nb_window_and_the_lfa_window_do_not_share_a_label() {
assert_ne!(
NB_ANFRAGE_WINDOW_LABEL,
BEENDIGUNG_ZUORDNUNG_ANTWORT_WINDOW_LABEL
);
}
#[test]
fn a_widerspruch_reaches_processd_with_its_grund() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&after_anfrage(),
BeendigungZuordnungCommand::ReceiveAntwort {
response_pid: Pruefidentifikator::new(55_012).expect("valid"),
antwortcode: "A35".to_owned(),
zustimmung: false,
grund: Some("Vertragsbindung bis 31.12.2026".to_owned()),
zuordnungsende: None,
},
)
.expect("answer accepted");
let p = &out.outbox[0].payload;
assert_eq!(p["antwortcode"], "A35");
assert_eq!(p["zustimmung"], false);
assert_eq!(p["grund"], "Vertragsbindung bis 31.12.2026");
assert_eq!(p["fristablauf"], false);
assert_eq!(
p["anmeldung_process_id"],
"11111111-1111-1111-1111-111111111111"
);
}
#[test]
fn fall_b_carries_the_earlier_zuordnungsende() {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&after_anfrage(),
BeendigungZuordnungCommand::ReceiveAntwort {
response_pid: Pruefidentifikator::new(55_011).expect("valid"),
antwortcode: "A34".to_owned(),
zustimmung: true,
grund: None,
zuordnungsende: Some("20261015".to_owned()),
},
)
.expect("answer accepted");
assert_eq!(out.outbox[0].payload["zuordnungsende"], "20261015");
}
#[test]
fn a_late_answer_is_refused_rather_than_reopening_the_process() {
let closed = {
let out = GpkeBeendigungZuordnungWorkflow::handle(
&after_anfrage(),
BeendigungZuordnungCommand::AntwortfristAbgelaufen,
)
.expect("lapse");
out.events.iter().fold(after_anfrage(), |st, e| {
GpkeBeendigungZuordnungWorkflow::apply(st, e)
})
};
let err = GpkeBeendigungZuordnungWorkflow::handle(
&closed,
BeendigungZuordnungCommand::ReceiveAntwort {
response_pid: Pruefidentifikator::new(55_011).expect("valid"),
antwortcode: "A36".to_owned(),
zustimmung: true,
grund: None,
zuordnungsende: None,
},
)
.expect_err("a late answer must not reopen the process");
assert!(format!("{err}").contains("AnfrageGesendet"), "{err}");
}
#[test]
fn an_anfrage_without_a_zuordnungsende_is_refused() {
let BeendigungZuordnungCommand::Anfragen {
sender,
receiver,
location_id,
tranche,
vorgangsnummer,
anmeldung_process_id,
kunde_name,
lfn_mp_id,
..
} = anfragen()
else {
unreachable!()
};
let err = GpkeBeendigungZuordnungWorkflow::handle(
&BeendigungZuordnungState::New,
BeendigungZuordnungCommand::Anfragen {
sender,
receiver,
location_id,
tranche,
process_date: String::new(),
vorgangsnummer,
anmeldung_process_id,
kunde_name,
lfn_mp_id,
},
)
.expect_err("55010 must name a Zuordnungsende");
assert!(format!("{err}").contains("Zuordnungsende"), "{err}");
}
}