use mako_engine::{
error::WorkflowError,
outbox::PendingOutbox,
types::{MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AnforderungKind {
NormierteProfile,
Lieferantenclearingliste,
Bilanzkreiszuordnungsliste,
ClearinglisteBas,
ClearinglisteDzr,
Bilanzierungsgebietsclearingliste,
BkSzrAggregationsebene,
ClearinglisteUenbDzr,
LieferantenausfallarbeitsClearingliste,
}
impl AnforderungKind {
#[must_use]
pub fn from_pid(pid: u32) -> Option<Self> {
Some(match pid {
17201 => Self::NormierteProfile,
17202 => Self::Lieferantenclearingliste,
17203 => Self::Bilanzkreiszuordnungsliste,
17204 => Self::ClearinglisteBas,
17205 => Self::ClearinglisteDzr,
17206 => Self::Bilanzierungsgebietsclearingliste,
17207 => Self::BkSzrAggregationsebene,
17208 => Self::ClearinglisteUenbDzr,
17210 => Self::LieferantenausfallarbeitsClearingliste,
_ => return None,
})
}
#[must_use]
pub fn pid(self) -> u32 {
match self {
Self::NormierteProfile => 17201,
Self::Lieferantenclearingliste => 17202,
Self::Bilanzkreiszuordnungsliste => 17203,
Self::ClearinglisteBas => 17204,
Self::ClearinglisteDzr => 17205,
Self::Bilanzierungsgebietsclearingliste => 17206,
Self::BkSzrAggregationsebene => 17207,
Self::ClearinglisteUenbDzr => 17208,
Self::LieferantenausfallarbeitsClearingliste => 17210,
}
}
#[must_use]
pub fn ablehnung(self, vorgang: AbonnementVorgang) -> Option<(u32, &'static str)> {
match (self, vorgang) {
(Self::BkSzrAggregationsebene, AbonnementVorgang::Bestellung) => {
Some((ABLEHNUNG_PID, "E_0003"))
}
(Self::BkSzrAggregationsebene, AbonnementVorgang::Abbestellung) => {
Some((ABLEHNUNG_PID, "E_0022"))
}
_ => None,
}
}
#[must_use]
pub fn label(self) -> &'static str {
match self {
Self::NormierteProfile => "Anforderung normierter Profile und Profilschar",
Self::Lieferantenclearingliste => "Anforderung Lieferantenclearingliste",
Self::Bilanzkreiszuordnungsliste => "Anforderung Bilanzkreiszuordnungsliste",
Self::ClearinglisteBas => "Anforderung Clearingliste BAS",
Self::ClearinglisteDzr => "Anforderung Clearingliste DZR",
Self::Bilanzierungsgebietsclearingliste => {
"Anforderung Bilanzierungsgebietsclearingliste"
}
Self::BkSzrAggregationsebene => "Ab-/Bestellung BK-SZR auf Aggregationsebene",
Self::ClearinglisteUenbDzr => "Anforderung Clearingliste ÜNB-DZR",
Self::LieferantenausfallarbeitsClearingliste => {
"Anforderung Lieferantenausfallarbeitsclearingliste"
}
}
}
#[must_use]
pub fn supports_abonnement(self) -> bool {
!matches!(
self,
Self::ClearinglisteBas | Self::ClearinglisteDzr | Self::ClearinglisteUenbDzr
)
}
#[must_use]
pub fn requester_role(self) -> &'static str {
match self {
Self::NormierteProfile
| Self::Lieferantenclearingliste
| Self::LieferantenausfallarbeitsClearingliste => "LF",
Self::Bilanzkreiszuordnungsliste
| Self::ClearinglisteBas
| Self::BkSzrAggregationsebene => "BKV",
Self::ClearinglisteDzr | Self::Bilanzierungsgebietsclearingliste => "NB",
Self::ClearinglisteUenbDzr => "ÜNB",
}
}
#[must_use]
pub fn target_roles(self) -> &'static [&'static str] {
match self {
Self::NormierteProfile | Self::LieferantenausfallarbeitsClearingliste => &["NB"],
Self::Lieferantenclearingliste | Self::Bilanzkreiszuordnungsliste => &["NB", "ÜNB"],
Self::ClearinglisteBas | Self::ClearinglisteDzr | Self::ClearinglisteUenbDzr => {
&["BIKO"]
}
Self::Bilanzierungsgebietsclearingliste | Self::BkSzrAggregationsebene => &["ÜNB"],
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AbonnementVorgang {
Bestellung,
Abbestellung,
}
pub const ANFORDERUNG_PIDS: &[u32] = &[
17201, 17202, 17203, 17204, 17205, 17206, 17207, 17208, 17210,
];
pub const ABLEHNUNG_PID: u32 = 19204;
#[must_use]
pub fn all_pids() -> Vec<u32> {
let mut v = ANFORDERUNG_PIDS.to_vec();
v.push(ABLEHNUNG_PID);
v.sort_unstable();
v
}
pub const WORKFLOW_NAME: &str = "mabis-anforderung";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct AnforderungData {
pub pruefidentifikator: Pruefidentifikator,
pub kind: AnforderungKind,
pub vorgang: AbonnementVorgang,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub message_ref: MessageRef,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum AnforderungEvent {
AnforderungGesendet {
pruefidentifikator: Pruefidentifikator,
kind: AnforderungKind,
vorgang: AbonnementVorgang,
receiver: MarktpartnerCode,
},
AnforderungErhalten {
pruefidentifikator: Pruefidentifikator,
kind: AnforderungKind,
vorgang: AbonnementVorgang,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
},
AnforderungAbgelehnt {
pruefidentifikator: Pruefidentifikator,
ebd: String,
code: String,
message_ref: MessageRef,
},
ValidationFailed {
reason: String,
},
}
impl EventPayload for AnforderungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnforderungGesendet { .. } => "MabisAnforderungGesendet",
Self::AnforderungErhalten { .. } => "MabisAnforderungErhalten",
Self::AnforderungAbgelehnt { .. } => "MabisAnforderungAbgelehnt",
Self::ValidationFailed { .. } => "MabisAnforderungValidationFailed",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
#[serde(tag = "status", content = "data")]
pub enum AnforderungState {
#[default]
New,
Gesendet {
kind: AnforderungKind,
vorgang: AbonnementVorgang,
},
Erhalten(Box<AnforderungData>),
Abgelehnt {
kind: AnforderungKind,
vorgang: AbonnementVorgang,
ebd: String,
code: String,
},
ValidationFailed {
reason: String,
},
}
impl AnforderungState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Gesendet { .. } => "Gesendet",
Self::Erhalten(_) => "Erhalten",
Self::Abgelehnt { .. } => "Abgelehnt",
Self::ValidationFailed { .. } => "ValidationFailed",
}
}
}
#[derive(Clone)]
pub enum AnforderungCommand {
SendAnforderung {
kind: AnforderungKind,
vorgang: AbonnementVorgang,
receiver: MarktpartnerCode,
},
ReceiveAnforderung {
pid: Pruefidentifikator,
vorgang: AbonnementVorgang,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
ReceiveAblehnung {
pid: Pruefidentifikator,
ebd: String,
code: String,
message_ref: MessageRef,
},
}
impl CommandPayload for AnforderungCommand {}
pub struct MabisAnforderungWorkflow;
impl Workflow for MabisAnforderungWorkflow {
type State = AnforderungState;
type Event = AnforderungEvent;
type Command = AnforderungCommand;
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
AnforderungEvent::AnforderungGesendet { kind, vorgang, .. } => {
AnforderungState::Gesendet {
kind: *kind,
vorgang: *vorgang,
}
}
AnforderungEvent::AnforderungErhalten {
pruefidentifikator,
kind,
vorgang,
sender,
receiver,
message_ref,
} => AnforderungState::Erhalten(Box::new(AnforderungData {
pruefidentifikator: *pruefidentifikator,
kind: *kind,
vorgang: *vorgang,
sender: sender.clone(),
receiver: receiver.clone(),
message_ref: message_ref.clone(),
})),
AnforderungEvent::AnforderungAbgelehnt { ebd, code, .. } => match state {
AnforderungState::Gesendet { kind, vorgang } => AnforderungState::Abgelehnt {
kind,
vorgang,
ebd: ebd.clone(),
code: code.clone(),
},
other => other,
},
AnforderungEvent::ValidationFailed { reason } => AnforderungState::ValidationFailed {
reason: reason.clone(),
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
AnforderungCommand::SendAnforderung {
kind,
vorgang,
receiver,
} => {
if !matches!(state, AnforderungState::New) {
return Ok(vec![].into());
}
if vorgang == AbonnementVorgang::Abbestellung && !kind.supports_abonnement() {
return Err(WorkflowError::rejected(format!(
"{} (PID {}) is a one-shot request — the AHB defines no \
Abonnement to end",
kind.label(),
kind.pid()
)));
}
let pid = Pruefidentifikator::new(kind.pid()).map_err(|e| {
WorkflowError::rejected(format!("invalid PID {}: {e}", kind.pid()))
})?;
let outbox = PendingOutbox::new(
"ORDERS",
receiver.as_str(),
serde_json::json!({
"pid": kind.pid(),
"vorgang": vorgang,
"kind": kind,
}),
);
Ok(WorkflowOutput {
events: vec![AnforderungEvent::AnforderungGesendet {
pruefidentifikator: pid,
kind,
vorgang,
receiver,
}],
outbox: vec![outbox],
deadlines: vec![],
})
}
AnforderungCommand::ReceiveAnforderung {
pid,
vorgang,
sender,
receiver,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, AnforderungState::New) {
return Ok(vec![].into());
}
let Some(kind) = AnforderungKind::from_pid(pid.as_u32()) else {
return Err(WorkflowError::rejected(format!(
"PID {pid} is not a MaBiS Anforderung; expected one of {ANFORDERUNG_PIDS:?}"
)));
};
if !validation_passed {
return Ok(vec![AnforderungEvent::ValidationFailed {
reason: validation_errors.join("; "),
}]
.into());
}
if vorgang == AbonnementVorgang::Abbestellung && !kind.supports_abonnement() {
return Err(WorkflowError::rejected(format!(
"inbound {} (PID {}) claims an Abbestellung, but the AHB \
defines no Abonnement for it",
kind.label(),
kind.pid()
)));
}
Ok(vec![AnforderungEvent::AnforderungErhalten {
pruefidentifikator: pid,
kind,
vorgang,
sender,
receiver,
message_ref,
}]
.into())
}
AnforderungCommand::ReceiveAblehnung {
pid,
ebd,
code,
message_ref,
} => {
let AnforderungState::Gesendet { kind, vorgang } = state else {
return Err(WorkflowError::invalid_state("Gesendet", state.label()));
};
if pid.as_u32() != ABLEHNUNG_PID {
return Err(WorkflowError::rejected(format!(
"PID {pid} ist keine Ablehnung einer MaBiS-Anforderung \
(erwartet {ABLEHNUNG_PID})"
)));
}
let Some((_, erwarteter_ebd)) = kind.ablehnung(*vorgang) else {
return Err(WorkflowError::rejected(format!(
"{} (PID {}) hat keine Ablehnung — der AHB definiert für \
diesen Code keinen ORDRSP",
kind.label(),
kind.pid()
)));
};
if ebd != erwarteter_ebd {
return Err(WorkflowError::rejected(format!(
"Ablehnung nennt EBD {ebd}, für {} ist aber \
{erwarteter_ebd} maßgeblich",
kind.label()
)));
}
if code.trim().is_empty() {
return Err(WorkflowError::rejected(
"Ablehnung ohne Antwortcode ist nicht auswertbar",
));
}
Ok(vec![AnforderungEvent::AnforderungAbgelehnt {
pruefidentifikator: pid,
ebd,
code,
message_ref,
}]
.into())
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn mp(s: &str) -> MarktpartnerCode {
MarktpartnerCode::new(s)
}
fn fold(events: &[AnforderungEvent]) -> AnforderungState {
events.iter().fold(AnforderungState::default(), |s, e| {
MabisAnforderungWorkflow::apply(s, e)
})
}
#[test]
fn the_pid_table_round_trips() {
for &pid in ANFORDERUNG_PIDS {
let kind = AnforderungKind::from_pid(pid).expect("known PID");
assert_eq!(kind.pid(), pid, "from_pid/pid disagree for {pid}");
}
assert!(AnforderungKind::from_pid(17209).is_none());
assert!(AnforderungKind::from_pid(55062).is_none());
}
#[test]
fn one_shot_requests_cannot_be_unsubscribed() {
for kind in [
AnforderungKind::ClearinglisteBas,
AnforderungKind::ClearinglisteDzr,
AnforderungKind::ClearinglisteUenbDzr,
] {
assert!(!kind.supports_abonnement(), "{kind:?}");
let err = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::SendAnforderung {
kind,
vorgang: AbonnementVorgang::Abbestellung,
receiver: mp("9900123456789"),
},
)
.expect_err("must reject");
assert!(format!("{err}").contains("one-shot"), "{err}");
}
}
#[test]
fn subscription_kinds_accept_both_directions_under_the_same_pid() {
let kind = AnforderungKind::BkSzrAggregationsebene;
assert!(kind.supports_abonnement());
for vorgang in [
AbonnementVorgang::Bestellung,
AbonnementVorgang::Abbestellung,
] {
let out = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::SendAnforderung {
kind,
vorgang,
receiver: mp("9900123456789"),
},
)
.expect("accepted");
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].payload["pid"], 17207);
assert_eq!(fold(&out.events).label(), "Gesendet");
}
}
#[test]
fn an_inbound_anforderung_is_recorded() {
let out = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::ReceiveAnforderung {
pid: Pruefidentifikator::new(17205).unwrap(),
vorgang: AbonnementVorgang::Bestellung,
sender: mp("9900123456789"),
receiver: mp("9900987654321"),
message_ref: MessageRef::new("MSG-1"),
validation_passed: true,
validation_errors: vec![],
},
)
.expect("accepted");
assert!(out.outbox.is_empty(), "receiving emits no message");
assert_eq!(fold(&out.events).label(), "Erhalten");
}
#[test]
fn an_inbound_abbestellung_on_a_one_shot_code_is_rejected() {
let err = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::ReceiveAnforderung {
pid: Pruefidentifikator::new(17204).unwrap(),
vorgang: AbonnementVorgang::Abbestellung,
sender: mp("9900123456789"),
receiver: mp("9900987654321"),
message_ref: MessageRef::new("MSG-1"),
validation_passed: true,
validation_errors: vec![],
},
)
.expect_err("must reject");
assert!(format!("{err}").contains("no Abonnement"), "{err}");
}
#[test]
fn validation_failure_is_terminal() {
let out = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::ReceiveAnforderung {
pid: Pruefidentifikator::new(17201).unwrap(),
vorgang: AbonnementVorgang::Bestellung,
sender: mp("9900123456789"),
receiver: mp("9900987654321"),
message_ref: MessageRef::new("MSG-1"),
validation_passed: false,
validation_errors: vec!["BGM missing".to_owned()],
},
)
.expect("accepted");
assert_eq!(fold(&out.events).label(), "ValidationFailed");
}
#[test]
fn an_unknown_orders_pid_is_rejected() {
let err = MabisAnforderungWorkflow::handle(
&AnforderungState::New,
AnforderungCommand::ReceiveAnforderung {
pid: Pruefidentifikator::new(17004).unwrap(),
vorgang: AbonnementVorgang::Bestellung,
sender: mp("9900123456789"),
receiver: mp("9900987654321"),
message_ref: MessageRef::new("MSG-1"),
validation_passed: true,
validation_errors: vec![],
},
)
.expect_err("must reject");
assert!(
format!("{err}").contains("not a MaBiS Anforderung"),
"{err}"
);
}
#[test]
fn roles_match_the_bdew_overview() {
assert_eq!(AnforderungKind::NormierteProfile.requester_role(), "LF");
assert_eq!(AnforderungKind::NormierteProfile.target_roles(), &["NB"]);
assert_eq!(
AnforderungKind::ClearinglisteUenbDzr.requester_role(),
"ÜNB"
);
assert_eq!(
AnforderungKind::ClearinglisteUenbDzr.target_roles(),
&["BIKO"]
);
assert_eq!(
AnforderungKind::Bilanzkreiszuordnungsliste.target_roles(),
&["NB", "ÜNB"]
);
}
}
#[cfg(test)]
mod ablehnung_tests {
use super::*;
fn gesendet(kind: AnforderungKind, vorgang: AbonnementVorgang) -> AnforderungState {
AnforderungState::Gesendet { kind, vorgang }
}
fn ablehnung(ebd: &str, code: &str) -> AnforderungCommand {
AnforderungCommand::ReceiveAblehnung {
pid: Pruefidentifikator::new(ABLEHNUNG_PID).expect("19204"),
ebd: ebd.to_owned(),
code: code.to_owned(),
message_ref: MessageRef::new("ORDRSP-1"),
}
}
#[test]
fn only_17207_can_be_refused() {
for &pid in ANFORDERUNG_PIDS {
let kind = AnforderungKind::from_pid(pid).expect("in the table");
let expected = kind == AnforderungKind::BkSzrAggregationsebene;
for vorgang in [
AbonnementVorgang::Bestellung,
AbonnementVorgang::Abbestellung,
] {
assert_eq!(
kind.ablehnung(vorgang).is_some(),
expected,
"{} ({pid}) / {vorgang:?}",
kind.label()
);
}
}
}
#[test]
fn the_ebd_differs_between_bestellung_and_abbestellung() {
let k = AnforderungKind::BkSzrAggregationsebene;
assert_eq!(
k.ablehnung(AbonnementVorgang::Bestellung),
Some((ABLEHNUNG_PID, "E_0003"))
);
assert_eq!(
k.ablehnung(AbonnementVorgang::Abbestellung),
Some((ABLEHNUNG_PID, "E_0022"))
);
}
#[test]
fn a_refusal_against_the_wrong_tree_is_rejected() {
let state = gesendet(
AnforderungKind::BkSzrAggregationsebene,
AbonnementVorgang::Bestellung,
);
assert!(MabisAnforderungWorkflow::handle(&state, ablehnung("E_0022", "A01")).is_err());
assert!(MabisAnforderungWorkflow::handle(&state, ablehnung("E_0003", "A01")).is_ok());
}
#[test]
fn a_refusal_needs_a_code() {
let state = gesendet(
AnforderungKind::BkSzrAggregationsebene,
AbonnementVorgang::Bestellung,
);
assert!(MabisAnforderungWorkflow::handle(&state, ablehnung("E_0003", " ")).is_err());
}
#[test]
fn a_refusal_of_a_code_that_has_none_is_rejected() {
let state = gesendet(
AnforderungKind::ClearinglisteBas,
AbonnementVorgang::Bestellung,
);
assert!(MabisAnforderungWorkflow::handle(&state, ablehnung("E_0003", "A01")).is_err());
}
#[test]
fn a_refusal_before_the_request_was_sent_is_rejected() {
assert!(
MabisAnforderungWorkflow::handle(&AnforderungState::New, ablehnung("E_0003", "A01"))
.is_err()
);
}
#[test]
fn a_refusal_lands_in_abgelehnt() {
let state = gesendet(
AnforderungKind::BkSzrAggregationsebene,
AbonnementVorgang::Abbestellung,
);
let out =
MabisAnforderungWorkflow::handle(&state, ablehnung("E_0022", "A02")).expect("accepted");
let next = out
.events
.iter()
.fold(state, MabisAnforderungWorkflow::apply);
assert_eq!(next.label(), "Abgelehnt");
}
#[test]
fn the_lieferantenausfallarbeits_clearingliste_is_a_mabis_request() {
let k = AnforderungKind::from_pid(17210).expect("17210 is a MaBiS Anforderung");
assert_eq!(k.requester_role(), "LF");
assert_eq!(k.target_roles(), &["NB"]);
assert!(
k.supports_abonnement(),
"the AHB names a Beendigung des Abonnements"
);
}
#[test]
fn all_pids_includes_the_ablehnung() {
let pids = all_pids();
assert!(pids.contains(&ABLEHNUNG_PID));
assert_eq!(pids.len(), ANFORDERUNG_PIDS.len() + 1);
}
}