use mako_engine::{
deadline::Deadline,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const ABRECHNUNGSDATEN_PIDS: &[u32] = &[55_156, 55_220, 55_673];
pub const BEARBEITUNGSSTAND_PID: u32 = 21_047;
pub const BESTELLUNG_EBD: &str = "E_0595";
pub const WORKFLOW_NAME: &str = "gpke-abrechnungsdaten";
pub const BEARBEITUNGSSTAND_WINDOW_LABEL: &str = "gpke-abrechnungsdaten-bearbeitungsstand";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum AbrechnungsdatenEvent {
RueckmeldungErhalten {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
BearbeitungsstandGesendet {
antwort_code: String,
sendet_stammdatenaenderung: bool,
bemerkung: Option<String>,
},
Abgebrochen {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for AbrechnungsdatenEvent {
fn event_type(&self) -> &'static str {
match self {
Self::RueckmeldungErhalten { .. } => "AbrechnungsdatenRueckmeldungErhalten",
Self::BearbeitungsstandGesendet { .. } => "AbrechnungsdatenBearbeitungsstandGesendet",
Self::Abgebrochen { .. } => "AbrechnungsdatenAbgebrochen",
Self::DeadlineExpired { .. } => "AbrechnungsdatenDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AbrechnungsdatenData {
pub pruefidentifikator: Pruefidentifikator,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub location_id: MaLo,
pub document_date: String,
pub message_ref: MessageRef,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum AbrechnungsdatenState {
#[default]
New,
Eingegangen(AbrechnungsdatenData),
Beantwortet {
data: AbrechnungsdatenData,
antwort_code: String,
},
Abgebrochen {
reason: String,
},
}
impl AbrechnungsdatenState {
#[must_use]
pub const fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Eingegangen(_) => "Eingegangen",
Self::Beantwortet { .. } => "Beantwortet",
Self::Abgebrochen { .. } => "Abgebrochen",
}
}
#[must_use]
pub const fn is_terminal(&self) -> bool {
matches!(self, Self::Beantwortet { .. } | Self::Abgebrochen { .. })
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "command", content = "data")]
pub enum AbrechnungsdatenCommand {
ReceiveRueckmeldung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
SendBearbeitungsstand {
antwort_code: String,
sendet_stammdatenaenderung: bool,
bemerkung: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for AbrechnungsdatenCommand {}
#[derive(Debug, Clone, Copy, Default)]
pub struct GpkeAbrechnungsdatenWorkflow;
impl Workflow for GpkeAbrechnungsdatenWorkflow {
type State = AbrechnungsdatenState;
type Command = AbrechnungsdatenCommand;
type Event = AbrechnungsdatenEvent;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
(deadline.label() == BEARBEITUNGSSTAND_WINDOW_LABEL && !state.is_terminal()).then(|| {
AbrechnungsdatenCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}
})
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
AbrechnungsdatenCommand::ReceiveRueckmeldung {
pid,
sender,
receiver,
location_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, AbrechnungsdatenState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !ABRECHNUNGSDATEN_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::other(format!(
"PID {} is not a GPKE Teil 2 § 3.1 Abrechnungsdaten-Rückmeldung",
pid.as_u32()
)));
}
Ok(vec![AbrechnungsdatenEvent::RueckmeldungErhalten {
pruefidentifikator: pid,
sender,
receiver,
location_id,
document_date,
message_ref,
validation_passed,
validation_errors,
}]
.into())
}
AbrechnungsdatenCommand::SendBearbeitungsstand {
antwort_code,
sendet_stammdatenaenderung,
bemerkung,
} => {
let AbrechnungsdatenState::Eingegangen(data) = state else {
return Err(WorkflowError::invalid_state("Eingegangen", state.label()));
};
let events = vec![AbrechnungsdatenEvent::BearbeitungsstandGesendet {
antwort_code: antwort_code.clone(),
sendet_stammdatenaenderung,
bemerkung: bemerkung.clone(),
}];
let outbox = vec![PendingOutbox::new(
"IFTSTA",
data.sender.as_str(),
serde_json::json!({
"pid": BEARBEITUNGSSTAND_PID,
"anfrage_pid": data.pruefidentifikator.as_u32(),
"sender": data.receiver.as_str(),
"receiver": data.sender.as_str(),
"malo": data.location_id.as_str(),
"antwort_code": antwort_code,
"antwort_codeliste": BESTELLUNG_EBD,
"sendet_stammdatenaenderung": sendet_stammdatenaenderung,
"bemerkung": bemerkung,
}),
)];
Ok(WorkflowOutput::with_outbox(events, outbox))
}
AbrechnungsdatenCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(Vec::new().into());
}
Ok(vec![AbrechnungsdatenEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
AbrechnungsdatenEvent::RueckmeldungErhalten {
pruefidentifikator,
sender,
receiver,
location_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if *validation_passed {
AbrechnungsdatenState::Eingegangen(AbrechnungsdatenData {
pruefidentifikator: *pruefidentifikator,
sender: sender.clone(),
receiver: receiver.clone(),
location_id: location_id.clone(),
document_date: document_date.clone(),
message_ref: message_ref.clone(),
})
} else {
AbrechnungsdatenState::Abgebrochen {
reason: format!("AHB validation failed: {}", validation_errors.join("; ")),
}
}
}
AbrechnungsdatenEvent::BearbeitungsstandGesendet { antwort_code, .. } => match state {
AbrechnungsdatenState::Eingegangen(data) => AbrechnungsdatenState::Beantwortet {
data,
antwort_code: antwort_code.clone(),
},
other => other,
},
AbrechnungsdatenEvent::Abgebrochen { reason } => AbrechnungsdatenState::Abgebrochen {
reason: reason.clone(),
},
AbrechnungsdatenEvent::DeadlineExpired { label, .. } => {
if state.is_terminal() {
state
} else {
AbrechnungsdatenState::Abgebrochen {
reason: format!(
"{label} expired without a Bearbeitungsstandsmeldung — GPKE Teil 2 \
§ 3.1 gives the NB until the 2. WT nach dem ÜT"
),
}
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn pid(code: u32) -> Pruefidentifikator {
Pruefidentifikator::new(code).expect("valid PID")
}
fn receive(code: u32) -> AbrechnungsdatenCommand {
AbrechnungsdatenCommand::ReceiveRueckmeldung {
pid: pid(code),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9900357000004"),
location_id: MaLo::new("51238696012"),
document_date: "20260304".to_owned(),
message_ref: MessageRef::new("ABR-001"),
validation_passed: true,
validation_errors: vec![],
}
}
fn eingegangen(code: u32) -> AbrechnungsdatenState {
let out =
GpkeAbrechnungsdatenWorkflow::handle(&AbrechnungsdatenState::default(), receive(code))
.expect("accepted");
out.events.iter().fold(
AbrechnungsdatenState::default(),
GpkeAbrechnungsdatenWorkflow::apply,
)
}
#[test]
fn every_routed_pid_is_accepted_and_nothing_else_is() {
for code in ABRECHNUNGSDATEN_PIDS {
assert!(matches!(
eingegangen(*code),
AbrechnungsdatenState::Eingegangen(_)
));
}
assert!(
GpkeAbrechnungsdatenWorkflow::handle(
&AbrechnungsdatenState::default(),
receive(55_001)
)
.is_err()
);
}
#[test]
fn the_bearbeitungsstand_reaches_the_wire_with_its_code() {
let out = GpkeAbrechnungsdatenWorkflow::handle(
&eingegangen(55_156),
AbrechnungsdatenCommand::SendBearbeitungsstand {
antwort_code: "A05".to_owned(),
sendet_stammdatenaenderung: true,
bemerkung: None,
},
)
.expect("answered");
let iftsta = out
.outbox
.iter()
.find(|o| &*o.message_type == "IFTSTA")
.expect("the answer must produce an outbound IFTSTA");
assert_eq!(iftsta.payload["pid"], BEARBEITUNGSSTAND_PID);
assert_eq!(iftsta.payload["antwort_code"], "A05");
assert_eq!(iftsta.payload["antwort_codeliste"], "E_0595");
assert_eq!(iftsta.payload["sendet_stammdatenaenderung"], true);
assert_eq!(iftsta.payload["sender"], "9900357000004");
assert_eq!(iftsta.payload["receiver"], "4012345000023");
}
#[test]
fn both_clusters_ride_the_same_pid() {
for (code, sendet) in [("A02", true), ("A03", false)] {
let out = GpkeAbrechnungsdatenWorkflow::handle(
&eingegangen(55_220),
AbrechnungsdatenCommand::SendBearbeitungsstand {
antwort_code: code.to_owned(),
sendet_stammdatenaenderung: sendet,
bemerkung: None,
},
)
.expect("answered");
let iftsta = out
.outbox
.iter()
.find(|o| &*o.message_type == "IFTSTA")
.expect("outbound IFTSTA");
assert_eq!(iftsta.payload["pid"], BEARBEITUNGSSTAND_PID, "{code}");
}
}
#[test]
fn the_codes_belong_to_the_clearing_branch() {
assert_eq!(BESTELLUNG_EBD, mako_pruefung::codes::EBD_BESTELLUNG);
for code in mako_pruefung::codes::E_0595_CLEARING_CODES {
assert!(
mako_pruefung::codes::lookup(BESTELLUNG_EBD, code).is_some(),
"{code}"
);
}
}
#[test]
fn a_deadline_on_a_settled_process_is_a_no_op() {
let answered = {
let out = GpkeAbrechnungsdatenWorkflow::handle(
&eingegangen(55_673),
AbrechnungsdatenCommand::SendBearbeitungsstand {
antwort_code: "A01".to_owned(),
sendet_stammdatenaenderung: false,
bemerkung: None,
},
)
.expect("answered");
out.events
.iter()
.fold(eingegangen(55_673), GpkeAbrechnungsdatenWorkflow::apply)
};
let out = GpkeAbrechnungsdatenWorkflow::handle(
&answered,
AbrechnungsdatenCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: BEARBEITUNGSSTAND_WINDOW_LABEL.into(),
},
)
.expect("no-op");
assert!(out.events.is_empty());
}
#[test]
fn an_expired_window_closes_the_process() {
let out = GpkeAbrechnungsdatenWorkflow::handle(
&eingegangen(55_156),
AbrechnungsdatenCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: BEARBEITUNGSSTAND_WINDOW_LABEL.into(),
},
)
.expect("expired");
let state = out
.events
.iter()
.fold(eingegangen(55_156), GpkeAbrechnungsdatenWorkflow::apply);
assert!(matches!(state, AbrechnungsdatenState::Abgebrochen { .. }));
}
}