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-kuendigung";
pub const KUENDIGUNG_PIDS: &[u32] = &[55016];
pub const KUENDIGUNG_ANTWORT_WINDOW_LABEL: &str = "gpke-kuendigung-antwortfrist";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum KuendigungEvent {
KuendigungErhalten {
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>,
},
}
impl EventPayload for KuendigungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::KuendigungErhalten { .. } => "KuendigungKuendigungErhalten",
Self::ValidationPassed { .. } => "KuendigungValidationPassed",
Self::AntwortGesendet { .. } => "KuendigungAntwortGesendet",
Self::Beendet => "KuendigungBeendet",
Self::AperakFehlerDispatched { .. } => "KuendigungAperakFehlerDispatched",
Self::Rejected { .. } => "KuendigungRejected",
Self::DeadlineExpired { .. } => "KuendigungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct KuendigungData {
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 KuendigungState {
#[default]
New,
Eingegangen(KuendigungData),
ValidationPassed(KuendigungData),
AntwortGesendet {
data: KuendigungData,
response_pid: Pruefidentifikator,
},
Beendet(KuendigungData),
Rejected {
reason: String,
},
}
impl KuendigungState {
#[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",
}
}
#[must_use]
pub fn data(&self) -> Option<&KuendigungData> {
match self {
Self::Eingegangen(d) | Self::ValidationPassed(d) | Self::Beendet(d) => Some(d),
Self::AntwortGesendet { data, .. } => Some(data),
Self::New | Self::Rejected { .. } => None,
}
}
}
#[derive(Clone)]
pub enum KuendigungCommand {
ReceiveKuendigung {
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,
},
BeendenBestaetigen,
DispatchAperakFehler {
reason: String,
outbound_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for KuendigungCommand {}
pub struct GpkeKuendigungWorkflow;
impl Workflow for GpkeKuendigungWorkflow {
type State = KuendigungState;
type Event = KuendigungEvent;
type Command = KuendigungCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(
KUENDIGUNG_ANTWORT_WINDOW_LABEL,
KuendigungState::Eingegangen(_) | KuendigungState::ValidationPassed(_),
) => Some(KuendigungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
KuendigungEvent::KuendigungErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
pruefidentifikator,
vorgangsnummer,
..
} => KuendigungState::Eingegangen(KuendigungData {
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(),
}),
KuendigungEvent::ValidationPassed { .. } => match state {
KuendigungState::Eingegangen(data) => KuendigungState::ValidationPassed(data),
other => other,
},
KuendigungEvent::AntwortGesendet {
accepted,
response_pid,
..
} => {
if *accepted {
match state {
KuendigungState::ValidationPassed(data) => {
KuendigungState::AntwortGesendet {
response_pid: *response_pid,
data,
}
}
other => other,
}
} else {
KuendigungState::Rejected {
reason: "Anfrage abgelehnt".to_owned(),
}
}
}
KuendigungEvent::Beendet => match state {
KuendigungState::AntwortGesendet { data, .. } => KuendigungState::Beendet(data),
other => other,
},
KuendigungEvent::AperakFehlerDispatched { reason, .. } => KuendigungState::Rejected {
reason: format!("APERAK 29001: {reason}"),
},
KuendigungEvent::Rejected { reason } => KuendigungState::Rejected {
reason: reason.clone(),
},
KuendigungEvent::DeadlineExpired { label, .. } => match state {
KuendigungState::Beendet(_) | KuendigungState::Rejected { .. } => state,
_ => KuendigungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
KuendigungCommand::ReceiveKuendigung {
pid,
sender,
receiver,
location_id,
document_date,
process_date,
message_ref,
vorgang,
validation_passed,
validation_errors,
} => {
if !matches!(state, KuendigungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !KUENDIGUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected Anfrage zur Beendigung der Zuordnung PID (55016), 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![KuendigungEvent::KuendigungErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
vorgangsnummer: vorgang.vorgangsnummer.clone(),
}];
if validation_passed {
events.push(KuendigungEvent::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(KuendigungEvent::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))
}
}
KuendigungCommand::SendAntwort { antwort } => {
let data = match state {
KuendigungState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.label(),
));
}
};
let accepted = antwort.zustimmung;
let response_code: u32 = if accepted { 55017 } else { 55018 };
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![KuendigungEvent::AntwortGesendet {
response_pid,
accepted,
antwort,
}],
outbox,
))
}
KuendigungCommand::BeendenBestaetigen => {
if !matches!(state, KuendigungState::AntwortGesendet { .. }) {
return Err(WorkflowError::invalid_state(
"AntwortGesendet",
state.label(),
));
}
Ok(vec![KuendigungEvent::Beendet].into())
}
KuendigungCommand::DispatchAperakFehler {
reason,
outbound_ref,
} => {
match state {
KuendigungState::Eingegangen(_) | KuendigungState::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![KuendigungEvent::AperakFehlerDispatched {
aperak_pid,
reason,
outbound_ref,
}]
.into())
}
KuendigungCommand::TimeoutExpired { deadline_id, label } => match state {
KuendigungState::Beendet(_) | KuendigungState::Rejected { .. } => Ok(vec![].into()),
_ => Ok(vec![KuendigungEvent::DeadlineExpired { deadline_id, label }].into()),
},
}
}
}
#[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) -> KuendigungCommand {
KuendigungCommand::ReceiveKuendigung {
pid: pid(55016),
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: KuendigungState, events: &[KuendigungEvent]) -> KuendigungState {
events.iter().fold(init, GpkeKuendigungWorkflow::apply)
}
#[test]
fn happy_path_bestaetigung() {
let out = GpkeKuendigungWorkflow::handle(&KuendigungState::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(KuendigungState::New, &out.events);
assert!(matches!(state, KuendigungState::ValidationPassed(_)));
let out = GpkeKuendigungWorkflow::handle(
&state,
KuendigungCommand::SendAntwort {
antwort: crate::lf_antwort::LfAntwort::zustimmung("A36", "E_0624"),
},
)
.unwrap();
if let KuendigungEvent::AntwortGesendet { response_pid, .. } = &out.events[0] {
assert_eq!(response_pid.as_u32(), 55017);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
let out =
GpkeKuendigungWorkflow::handle(&state, KuendigungCommand::BeendenBestaetigen).unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(state, KuendigungState::Beendet(_)));
}
#[test]
fn ablehnung_yields_55018() {
let out = GpkeKuendigungWorkflow::handle(&KuendigungState::New, anfrage_cmd(true)).unwrap();
let state = apply_all(KuendigungState::New, &out.events);
let out = GpkeKuendigungWorkflow::handle(
&state,
KuendigungCommand::SendAntwort {
antwort: crate::lf_antwort::LfAntwort::ablehnung("A35", "E_0624")
.with_bemerkung("Widerspruch"),
},
)
.unwrap();
if let KuendigungEvent::AntwortGesendet { response_pid, .. } = &out.events[0] {
assert_eq!(response_pid.as_u32(), 55018);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
assert!(matches!(state, KuendigungState::Rejected { .. }));
}
#[test]
fn validation_failure_emits_aperak_313() {
let out =
GpkeKuendigungWorkflow::handle(&KuendigungState::New, anfrage_cmd(false)).unwrap();
assert_eq!(out.outbox[0].payload["error_code"], "Z29");
let state = apply_all(KuendigungState::New, &out.events);
assert!(matches!(state, KuendigungState::Rejected { .. }));
}
#[test]
fn wrong_pid_rejected() {
let mut cmd = anfrage_cmd(true);
if let KuendigungCommand::ReceiveKuendigung { pid: p, .. } = &mut cmd {
*p = pid(55001);
}
assert!(GpkeKuendigungWorkflow::handle(&KuendigungState::New, cmd).is_err());
}
}