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-lf-abmeldung";
pub const LF_ABMELDUNG_PIDS: &[u32] = &[55007];
pub const LF_ABMELDUNG_APERAK_WINDOW_LABEL: &str = "gpke-lf-abmeldung-aperak-window";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum LfAbmeldungEvent {
AnkuendigungErhalten {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
document_date: String,
process_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
},
ValidationPassed {
message_ref: MessageRef,
},
AntwortGesendet {
response_pid: Pruefidentifikator,
accepted: bool,
reason: Option<String>,
},
Beendet,
AperakFehlerDispatched {
aperak_pid: Pruefidentifikator,
reason: String,
outbound_ref: MessageRef,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for LfAbmeldungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnkuendigungErhalten { .. } => "LfAbmeldungAnkuendigungErhalten",
Self::ValidationPassed { .. } => "LfAbmeldungValidationPassed",
Self::AntwortGesendet { .. } => "LfAbmeldungAntwortGesendet",
Self::Beendet => "LfAbmeldungBeendet",
Self::AperakFehlerDispatched { .. } => "LfAbmeldungAperakFehlerDispatched",
Self::Rejected { .. } => "LfAbmeldungRejected",
Self::DeadlineExpired { .. } => "LfAbmeldungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct LfAbmeldungData {
pub location_id: MaLo,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub document_date: String,
pub process_date: String,
pub pruefidentifikator: Pruefidentifikator,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum LfAbmeldungState {
New,
Eingegangen(LfAbmeldungData),
ValidationPassed(LfAbmeldungData),
AntwortGesendet {
data: LfAbmeldungData,
response_pid: Pruefidentifikator,
},
Beendet(LfAbmeldungData),
Rejected {
reason: String,
},
}
impl Default for LfAbmeldungState {
fn default() -> Self {
Self::New
}
}
impl LfAbmeldungState {
#[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<&LfAbmeldungData> {
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 LfAbmeldungCommand {
ReceiveAnkuendigung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
document_date: String,
process_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
SendAntwort {
accepted: bool,
reason: Option<String>,
},
BeendenBestaetigen,
DispatchAperakFehler {
reason: String,
outbound_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for LfAbmeldungCommand {}
pub struct GpkeLfAbmeldungWorkflow;
impl Workflow for GpkeLfAbmeldungWorkflow {
type State = LfAbmeldungState;
type Event = LfAbmeldungEvent;
type Command = LfAbmeldungCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(LF_ABMELDUNG_APERAK_WINDOW_LABEL, LfAbmeldungState::Eingegangen(_))
| (LF_ABMELDUNG_APERAK_WINDOW_LABEL, LfAbmeldungState::ValidationPassed(_)) => {
Some(LfAbmeldungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
LfAbmeldungEvent::AnkuendigungErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
pruefidentifikator,
..
} => LfAbmeldungState::Eingegangen(LfAbmeldungData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
document_date: document_date.clone(),
process_date: process_date.clone(),
pruefidentifikator: *pruefidentifikator,
}),
LfAbmeldungEvent::ValidationPassed { .. } => match state {
LfAbmeldungState::Eingegangen(data) => LfAbmeldungState::ValidationPassed(data),
other => other,
},
LfAbmeldungEvent::AntwortGesendet {
accepted,
response_pid,
..
} => {
if *accepted {
match state {
LfAbmeldungState::ValidationPassed(data) => {
LfAbmeldungState::AntwortGesendet {
response_pid: *response_pid,
data,
}
}
other => other,
}
} else {
LfAbmeldungState::Rejected {
reason: "Ankündigung abgelehnt".to_owned(),
}
}
}
LfAbmeldungEvent::Beendet => match state {
LfAbmeldungState::AntwortGesendet { data, .. } => LfAbmeldungState::Beendet(data),
other => other,
},
LfAbmeldungEvent::AperakFehlerDispatched { reason, .. } => LfAbmeldungState::Rejected {
reason: format!("APERAK 29001: {reason}"),
},
LfAbmeldungEvent::Rejected { reason } => LfAbmeldungState::Rejected {
reason: reason.clone(),
},
LfAbmeldungEvent::DeadlineExpired { label, .. } => match state {
LfAbmeldungState::Beendet(_) | LfAbmeldungState::Rejected { .. } => state,
_ => LfAbmeldungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
LfAbmeldungCommand::ReceiveAnkuendigung {
pid,
sender,
receiver,
location_id,
document_date,
process_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, LfAbmeldungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !LF_ABMELDUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected NB-initiated Lieferende PID (55007), got {pid}",
)));
}
let sender_mp_id = sender.clone();
let receiver_gln = receiver.clone();
let mut events = vec![LfAbmeldungEvent::AnkuendigungErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
}];
if validation_passed {
events.push(LfAbmeldungEvent::ValidationPassed { message_ref });
let outbox = vec![
PendingOutbox::new(
"APERAK",
sender_mp_id.as_str(),
serde_json::json!({
"sender": receiver_gln.as_str(),
"receiver": sender_mp_id.as_str(),
"pid": 29001_u32,
"document_code": "312",
}),
)
.caused_by(1),
];
Ok(WorkflowOutput::with_outbox(events, outbox))
} else {
let reason = validation_errors.join("; ");
events.push(LfAbmeldungEvent::Rejected {
reason: reason.clone(),
});
let outbox = vec![
PendingOutbox::new(
"APERAK",
sender_mp_id.as_str(),
serde_json::json!({
"sender": receiver_gln.as_str(),
"receiver": sender_mp_id.as_str(),
"pid": 29001_u32,
"error_code": mako_engine::erc::codes::Z29,
"reason": reason,
}),
)
.caused_by(0),
];
Ok(WorkflowOutput::with_outbox(events, outbox))
}
}
LfAbmeldungCommand::SendAntwort { accepted, reason } => {
match state {
LfAbmeldungState::ValidationPassed(_) => {}
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.label(),
));
}
}
let response_code: u32 = if accepted { 55008 } else { 55009 };
let response_pid = Pruefidentifikator::new(response_code)
.map_err(|e| WorkflowError::rejected(e.to_string()))?;
Ok(vec![LfAbmeldungEvent::AntwortGesendet {
response_pid,
accepted,
reason,
}]
.into())
}
LfAbmeldungCommand::BeendenBestaetigen => {
if !matches!(state, LfAbmeldungState::AntwortGesendet { .. }) {
return Err(WorkflowError::invalid_state(
"AntwortGesendet",
state.label(),
));
}
Ok(vec![LfAbmeldungEvent::Beendet].into())
}
LfAbmeldungCommand::DispatchAperakFehler {
reason,
outbound_ref,
} => {
match state {
LfAbmeldungState::Eingegangen(_) | LfAbmeldungState::ValidationPassed(_) => {}
_ => {
return Err(WorkflowError::invalid_state(
"Eingegangen or ValidationPassed",
state.label(),
));
}
}
let aperak_pid = Pruefidentifikator::new(29_001)
.map_err(|e| WorkflowError::rejected(e.to_string()))?;
Ok(vec![LfAbmeldungEvent::AperakFehlerDispatched {
aperak_pid,
reason,
outbound_ref,
}]
.into())
}
LfAbmeldungCommand::TimeoutExpired { deadline_id, label } => match state {
LfAbmeldungState::Beendet(_) | LfAbmeldungState::Rejected { .. } => {
Ok(vec![].into())
}
_ => Ok(vec![LfAbmeldungEvent::DeadlineExpired { deadline_id, label }].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 ankuendigung_cmd(ok: bool) -> LfAbmeldungCommand {
LfAbmeldungCommand::ReceiveAnkuendigung {
pid: pid(55007),
sender: mcod("9900357000004"),
receiver: mcod("4012345000023"),
location_id: malo("51238696781"),
document_date: "20251001".to_owned(),
process_date: "20260101".to_owned(),
message_ref: mref("ABMELD-001"),
validation_passed: ok,
validation_errors: if ok {
vec![]
} else {
vec!["missing mandatory segment".to_owned()]
},
}
}
fn apply_all(init: LfAbmeldungState, events: &[LfAbmeldungEvent]) -> LfAbmeldungState {
events.iter().fold(init, GpkeLfAbmeldungWorkflow::apply)
}
#[test]
fn lf_abmeldung_happy_path_bestaetigung() {
let out = GpkeLfAbmeldungWorkflow::handle(&LfAbmeldungState::New, ankuendigung_cmd(true))
.unwrap();
assert_eq!(out.events.len(), 2); let state = apply_all(LfAbmeldungState::New, &out.events);
assert!(matches!(state, LfAbmeldungState::ValidationPassed(_)));
let out = GpkeLfAbmeldungWorkflow::handle(
&state,
LfAbmeldungCommand::SendAntwort {
accepted: true,
reason: None,
},
)
.unwrap();
if let LfAbmeldungEvent::AntwortGesendet {
response_pid,
accepted,
..
} = &out.events[0]
{
assert!(accepted);
assert_eq!(response_pid.as_u32(), 55008);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
assert!(matches!(state, LfAbmeldungState::AntwortGesendet { .. }));
let out = GpkeLfAbmeldungWorkflow::handle(&state, LfAbmeldungCommand::BeendenBestaetigen)
.unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(state, LfAbmeldungState::Beendet(_)));
}
#[test]
fn lf_abmeldung_ablehnung() {
let out = GpkeLfAbmeldungWorkflow::handle(&LfAbmeldungState::New, ankuendigung_cmd(true))
.unwrap();
let state = apply_all(LfAbmeldungState::New, &out.events);
let out = GpkeLfAbmeldungWorkflow::handle(
&state,
LfAbmeldungCommand::SendAntwort {
accepted: false,
reason: Some("Widerspruch".to_owned()),
},
)
.unwrap();
if let LfAbmeldungEvent::AntwortGesendet {
response_pid,
accepted,
..
} = &out.events[0]
{
assert!(!accepted);
assert_eq!(response_pid.as_u32(), 55009);
} else {
panic!("expected AntwortGesendet");
}
let state = apply_all(state, &out.events);
assert!(matches!(state, LfAbmeldungState::Rejected { .. }));
}
#[test]
fn lf_abmeldung_wrong_pid_rejected() {
let result = GpkeLfAbmeldungWorkflow::handle(
&LfAbmeldungState::New,
LfAbmeldungCommand::ReceiveAnkuendigung {
pid: pid(55001),
sender: mcod("9900357000004"),
receiver: mcod("4012345000023"),
location_id: malo("51238696781"),
document_date: "20251001".to_owned(),
process_date: "20260101".to_owned(),
message_ref: mref("X"),
validation_passed: true,
validation_errors: vec![],
},
);
assert!(result.is_err());
}
#[test]
fn timeout_in_beendet_is_noop() {
let data = LfAbmeldungData {
location_id: malo("51238696781"),
sender: mcod("9900357000004"),
receiver: mcod("4012345000023"),
document_date: "20251001".to_owned(),
process_date: "20260101".to_owned(),
pruefidentifikator: pid(55007),
};
let state = LfAbmeldungState::Beendet(data);
let dl_id = DeadlineId::new();
let out = GpkeLfAbmeldungWorkflow::handle(
&state,
LfAbmeldungCommand::TimeoutExpired {
deadline_id: dl_id,
label: LF_ABMELDUNG_APERAK_WINDOW_LABEL.into(),
},
)
.unwrap();
assert!(out.events.is_empty());
}
}