use mako_engine::types::Pruefidentifikator;
use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "gpke-lf-anmeldung";
pub const ANFRAGE_PIDS_LF: &[u32] = &[
55001, 55004, 55016, 55077, ];
pub const ANMELDUNG_PIDS_MIT_BILANZKREIS: &[u32] = &[
55001, 55077, ];
pub const ANMELDUNG_PIDS_MIT_GESCHAEFTSVORFALL: &[u32] = &[
55077, ];
pub const GESCHAEFTSVORFALL_1: &str = "ZW0";
pub const GESCHAEFTSVORFALL_2: &str = "ZW1";
pub const GESCHAEFTSVORFALL_3: &str = "ZW2";
pub const GESCHAEFTSVORFAELLE: &[&str] = &[
GESCHAEFTSVORFALL_1,
GESCHAEFTSVORFALL_2,
GESCHAEFTSVORFALL_3,
];
pub const ANTWORT_PIDS_LF: &[u32] = &[
55002, 55003, 55005, 55006, 55017, 55018, 55078, 55080, ];
pub const NB_RESPONSE_WINDOW_LABEL: &str = "nb-response-window";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum LfAnmeldungEvent {
Initiated {
pruefidentifikator: Pruefidentifikator,
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
process_date: String,
},
AntwortReceived {
response_pid: Pruefidentifikator,
accepted: bool,
reason: Option<String>,
response_ref: MessageRef,
},
Activated,
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for LfAnmeldungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::Initiated { .. } => "LfAnmeldungInitiated",
Self::AntwortReceived { .. } => "LfAnmeldungAntwortReceived",
Self::Activated => "LfAnmeldungActivated",
Self::DeadlineExpired { .. } => "LfAnmeldungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct LfAnmeldungData {
pub pruefidentifikator: Pruefidentifikator,
pub location_id: MaLo,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub process_date: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
pub enum LfAnmeldungState {
#[default]
New,
Pending(LfAnmeldungData),
Active(LfAnmeldungData),
Rejected {
reason: String,
},
}
impl mako_engine::workflow::OccupiesBusinessKey for LfAnmeldungState {
fn occupies_business_key(&self) -> bool {
match self {
Self::Pending(_) | Self::Active(_) => true,
Self::New | Self::Rejected { .. } => false,
}
}
}
impl LfAnmeldungState {
fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Pending(_) => "Pending",
Self::Active(_) => "Active",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum LfAnmeldungCommand {
InitiateAnmeldung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
process_date: String,
transaktionsgrund: Option<String>,
bilanzkreis: Option<String>,
transaktionsgrund_ergaenzung: Option<String>,
tranchengroesse: Option<serde_json::Value>,
},
HandleAntwort {
response_pid: Pruefidentifikator,
accepted: bool,
reason: Option<String>,
response_ref: MessageRef,
},
Activate,
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for LfAnmeldungCommand {}
pub struct GpkeLfAnmeldungWorkflow;
impl Workflow for GpkeLfAnmeldungWorkflow {
type State = LfAnmeldungState;
type Event = LfAnmeldungEvent;
type Command = LfAnmeldungCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(NB_RESPONSE_WINDOW_LABEL, LfAnmeldungState::Pending(_)) => {
Some(LfAnmeldungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
LfAnmeldungEvent::Initiated {
pruefidentifikator,
location_id,
sender,
receiver,
process_date,
} => LfAnmeldungState::Pending(LfAnmeldungData {
pruefidentifikator: *pruefidentifikator,
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
process_date: process_date.clone(),
}),
LfAnmeldungEvent::AntwortReceived {
accepted, reason, ..
} => {
if *accepted {
match state {
LfAnmeldungState::Pending(data) => LfAnmeldungState::Active(data),
other => other,
}
} else {
LfAnmeldungState::Rejected {
reason: reason.clone().unwrap_or_else(|| "Ablehnung".to_owned()),
}
}
}
LfAnmeldungEvent::Activated => match state {
LfAnmeldungState::Active(data) => LfAnmeldungState::Active(data),
other => other,
},
LfAnmeldungEvent::DeadlineExpired { label, .. } => match state {
LfAnmeldungState::Active(_) | LfAnmeldungState::Rejected { .. } => state,
_ => LfAnmeldungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
LfAnmeldungCommand::InitiateAnmeldung {
pid,
sender,
receiver,
location_id,
process_date,
transaktionsgrund,
bilanzkreis,
transaktionsgrund_ergaenzung,
tranchengroesse,
} => {
if !matches!(state, LfAnmeldungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !ANFRAGE_PIDS_LF.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected an LF Anfrage PID (55001, 55004, 55016, 55077), got {pid}",
)));
}
if ANMELDUNG_PIDS_MIT_BILANZKREIS.contains(&pid.as_u32())
&& bilanzkreis.as_deref().is_none_or(str::is_empty)
{
return Err(WorkflowError::rejected(format!(
"Anmeldung {pid} ohne Bilanzkreis: die UTILMD AHB Strom 2.2 \
Kap. 5.3 macht das Produktpaket SG8 SEQ+Z79 mit dem Produkt-Code \
9991000002082 (Bilanzkreis) zur Muss-Angabe — ohne einen für den \
LF gültigen Bilanzkreis kann der NB die Zuordnung nicht vornehmen.",
)));
}
if ANMELDUNG_PIDS_MIT_GESCHAEFTSVORFALL.contains(&pid.as_u32()) {
let erg = transaktionsgrund_ergaenzung.as_deref().unwrap_or("");
if !GESCHAEFTSVORFAELLE.contains(&erg) {
return Err(WorkflowError::rejected(format!(
"Anmeldung {pid} ohne Geschäftsvorfall: SG4 STS+7 DE 9013 \
lässt für die erzeugende Marktlokation nur {} zu — \
ZW0 (100%ige Zuordnung), ZW1 (bestehende Tranche) und \
ZW2 (neu zu bildende Tranche) sind verschiedene Vorgänge.",
GESCHAEFTSVORFAELLE.join("/"),
)));
}
if erg == GESCHAEFTSVORFALL_3 && tranchengroesse.is_none() {
return Err(WorkflowError::rejected(format!(
"Anmeldung {pid} im Geschäftsvorfall 3 ohne Tranchengröße: \
die Codeliste der Konfigurationen 1.4 Kap. 6.1.1 macht den \
Produkt-Code 9991000002090 bei STS+7++xxx+ZW2 zur \
Muss-Angabe, weil die neu zu bildende Tranche sonst keine \
Größe hat.",
)));
}
}
let event = LfAnmeldungEvent::Initiated {
pruefidentifikator: pid,
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
process_date: process_date.clone(),
};
let outbox = PendingOutbox::new(
"UTILMD",
receiver.as_str(),
serde_json::json!({
"direction": "outbound",
"pid": pid.as_u32(),
"sender": sender.as_str(),
"receiver": receiver.as_str(),
"malo": location_id.as_str(),
"process_date": process_date,
"transaktionsgrund": transaktionsgrund,
"bilanzkreis": bilanzkreis,
"transaktionsgrund_ergaenzung": transaktionsgrund_ergaenzung,
"tranchengroesse": tranchengroesse,
}),
);
Ok(WorkflowOutput::with_outbox(vec![event], vec![outbox]))
}
LfAnmeldungCommand::HandleAntwort {
response_pid,
accepted,
reason,
response_ref,
} => {
if !matches!(state, LfAnmeldungState::Pending(_)) {
return Err(WorkflowError::invalid_state("Pending", state.label()));
}
if !ANTWORT_PIDS_LF.contains(&response_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected an LF Antwort PID (55002/55003, 55005/55006, 55017, 55018, 55078, 55080), got {response_pid}",
)));
}
let malo_id_str = match state {
LfAnmeldungState::Pending(data) => data.location_id.as_str().to_owned(),
_ => String::new(),
};
let outbox = vec![PendingOutbox::new(
"ProcessCompleted",
"",
serde_json::json!({
"pid": response_pid.as_u32(),
"malo_id": malo_id_str,
"accepted": accepted,
"outcome": if accepted { "accepted" } else { "rejected" },
}),
)];
Ok(WorkflowOutput::with_outbox(
vec![LfAnmeldungEvent::AntwortReceived {
response_pid,
accepted,
reason,
response_ref,
}],
outbox,
))
}
LfAnmeldungCommand::Activate => {
if !matches!(state, LfAnmeldungState::Active(_)) {
return Err(WorkflowError::invalid_state("Active", state.label()));
}
Ok(vec![LfAnmeldungEvent::Activated].into())
}
LfAnmeldungCommand::TimeoutExpired { deadline_id, label } => {
if matches!(
state,
LfAnmeldungState::Active(_) | LfAnmeldungState::Rejected { .. }
) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![LfAnmeldungEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}
#[cfg(test)]
mod tests {
use mako_engine::{
types::{MaLo, MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::Workflow,
};
use super::*;
fn make_initiate(pid: u32) -> LfAnmeldungCommand {
LfAnmeldungCommand::InitiateAnmeldung {
transaktionsgrund_ergaenzung: None,
tranchengroesse: None,
pid: Pruefidentifikator::new(pid).unwrap(),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
location_id: MaLo::new("10001234558"),
process_date: "2026-10-01".to_owned(),
transaktionsgrund: None,
bilanzkreis: Some("11XBK-LF-------9".to_owned()),
}
}
fn make_erzeugende(erg: &str, groesse: Option<serde_json::Value>) -> LfAnmeldungCommand {
let LfAnmeldungCommand::InitiateAnmeldung {
pid,
sender,
receiver,
location_id,
process_date,
transaktionsgrund,
bilanzkreis,
..
} = make_initiate(55077)
else {
unreachable!("make_initiate builds an InitiateAnmeldung")
};
LfAnmeldungCommand::InitiateAnmeldung {
pid,
sender,
receiver,
location_id,
process_date,
transaktionsgrund,
bilanzkreis,
transaktionsgrund_ergaenzung: Some(erg.to_owned()),
tranchengroesse: groesse,
}
}
#[test]
fn an_erzeugende_anmeldung_states_its_geschaeftsvorfall() {
let err = GpkeLfAnmeldungWorkflow::handle(&LfAnmeldungState::New, make_initiate(55077))
.expect_err("55077 without a Geschäftsvorfall must be refused");
assert!(format!("{err}").contains("Geschäftsvorfall"), "{err}");
for erg in GESCHAEFTSVORFAELLE {
let groesse = (*erg == GESCHAEFTSVORFALL_3).then(|| serde_json::json!("33.33"));
let out = GpkeLfAnmeldungWorkflow::handle(
&LfAnmeldungState::New,
make_erzeugende(erg, groesse),
)
.unwrap_or_else(|e| panic!("{erg} is a lawful Geschäftsvorfall: {e}"));
assert_eq!(out.outbox[0].payload["transaktionsgrund_ergaenzung"], *erg);
}
}
#[test]
fn geschaeftsvorfall_3_carries_a_tranchengroesse() {
let err = GpkeLfAnmeldungWorkflow::handle(
&LfAnmeldungState::New,
make_erzeugende(GESCHAEFTSVORFALL_3, None),
)
.expect_err("ZW2 without a Tranchengröße must be refused");
assert!(format!("{err}").contains("9991000002090"), "{err}");
let out = GpkeLfAnmeldungWorkflow::handle(
&LfAnmeldungState::New,
make_erzeugende(GESCHAEFTSVORFALL_3, Some(serde_json::json!("33.33"))),
)
.expect("a stated Tranchengröße is accepted");
assert_eq!(out.outbox[0].payload["tranchengroesse"], "33.33");
for erg in [GESCHAEFTSVORFALL_1, GESCHAEFTSVORFALL_2] {
let out =
GpkeLfAnmeldungWorkflow::handle(&LfAnmeldungState::New, make_erzeugende(erg, None))
.expect("no Tranchengröße is owed");
assert!(out.outbox[0].payload["tranchengroesse"].is_null());
}
}
#[test]
fn initiate_lieferbeginn_transitions_to_pending() {
let state = LfAnmeldungState::New;
let out = GpkeLfAnmeldungWorkflow::handle(&state, make_initiate(55001)).unwrap();
assert_eq!(out.events.len(), 1);
assert_eq!(out.outbox.len(), 1, "must enqueue UTILMD outbox entry");
let new_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(new_state, LfAnmeldungState::Pending(_)));
}
#[test]
fn initiate_abmeldung_transitions_to_pending() {
let state = LfAnmeldungState::New;
let out = GpkeLfAnmeldungWorkflow::handle(&state, make_initiate(55004)).unwrap();
let new_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(new_state, LfAnmeldungState::Pending(_)));
}
#[test]
fn initiate_kuendigung_transitions_to_pending() {
let state = LfAnmeldungState::New;
let out = GpkeLfAnmeldungWorkflow::handle(&state, make_initiate(55016)).unwrap();
let new_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(new_state, LfAnmeldungState::Pending(_)));
}
#[test]
fn nb_acceptance_transitions_to_active() {
let initiated_event = LfAnmeldungEvent::Initiated {
pruefidentifikator: Pruefidentifikator::new(55001).unwrap(),
location_id: MaLo::new("10001234558"),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
process_date: "2026-10-01".to_owned(),
};
let state = GpkeLfAnmeldungWorkflow::apply(LfAnmeldungState::New, &initiated_event);
let cmd = LfAnmeldungCommand::HandleAntwort {
response_pid: Pruefidentifikator::new(55002).unwrap(), accepted: true,
reason: None,
response_ref: MessageRef::new("NB-RESP-001"),
};
let out = GpkeLfAnmeldungWorkflow::handle(&state, cmd).unwrap();
assert_eq!(out.events.len(), 1);
let final_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(final_state, LfAnmeldungState::Active(_)));
}
#[test]
fn handle_antwort_emits_process_completed_for_marktd() {
let initiated_event = LfAnmeldungEvent::Initiated {
pruefidentifikator: Pruefidentifikator::new(55001).unwrap(),
location_id: MaLo::new("10001234558"),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
process_date: "2026-10-01".to_owned(),
};
let state = GpkeLfAnmeldungWorkflow::apply(LfAnmeldungState::New, &initiated_event);
for (response_pid, accepted, outcome) in
[(55002u32, true, "accepted"), (55003, false, "rejected")]
{
let out = GpkeLfAnmeldungWorkflow::handle(
&state,
LfAnmeldungCommand::HandleAntwort {
response_pid: Pruefidentifikator::new(response_pid).unwrap(),
accepted,
reason: None,
response_ref: MessageRef::new("NB-RESP-001"),
},
)
.unwrap();
assert_eq!(out.outbox.len(), 1, "{response_pid}: one outbox entry");
let entry = &out.outbox[0];
assert_eq!(entry.message_type.as_ref(), "ProcessCompleted");
assert_eq!(
entry.payload["pid"].as_u64().unwrap(),
u64::from(response_pid)
);
assert_eq!(
entry.payload["malo_id"].as_str().unwrap(),
"10001234558",
"{response_pid}: marktd resolves the MaLo from the payload, not the CE subject",
);
assert_eq!(entry.payload["outcome"].as_str().unwrap(), outcome);
}
}
#[test]
fn nb_rejection_transitions_to_rejected() {
let initiated_event = LfAnmeldungEvent::Initiated {
pruefidentifikator: Pruefidentifikator::new(55001).unwrap(),
location_id: MaLo::new("10001234558"),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
process_date: "2026-10-01".to_owned(),
};
let state = GpkeLfAnmeldungWorkflow::apply(LfAnmeldungState::New, &initiated_event);
let cmd = LfAnmeldungCommand::HandleAntwort {
response_pid: Pruefidentifikator::new(55003).unwrap(), accepted: false,
reason: Some("MaLo nicht in Netzgebiet".to_owned()),
response_ref: MessageRef::new("NB-RESP-002"),
};
let out = GpkeLfAnmeldungWorkflow::handle(&state, cmd).unwrap();
let final_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(final_state, LfAnmeldungState::Rejected { .. }));
}
#[test]
fn invalid_pid_is_rejected() {
let state = LfAnmeldungState::New;
let err = GpkeLfAnmeldungWorkflow::handle(&state, make_initiate(55003));
assert!(err.is_err());
}
#[test]
fn timeout_on_pending_transitions_to_rejected() {
use mako_engine::ids::DeadlineId;
let initiated_event = LfAnmeldungEvent::Initiated {
pruefidentifikator: Pruefidentifikator::new(55001).unwrap(),
location_id: MaLo::new("10001234558"),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
process_date: "2026-10-01".to_owned(),
};
let state = GpkeLfAnmeldungWorkflow::apply(LfAnmeldungState::New, &initiated_event);
let cmd = LfAnmeldungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "nb-response-window".into(),
};
let out = GpkeLfAnmeldungWorkflow::handle(&state, cmd).unwrap();
let final_state = GpkeLfAnmeldungWorkflow::apply(state, &out.events[0]);
assert!(matches!(final_state, LfAnmeldungState::Rejected { .. }));
}
#[test]
fn timeout_on_active_is_noop() {
use mako_engine::ids::DeadlineId;
let initiated_event = LfAnmeldungEvent::Initiated {
pruefidentifikator: Pruefidentifikator::new(55001).unwrap(),
location_id: MaLo::new("10001234558"),
sender: MarktpartnerCode::new("4012345000009"),
receiver: MarktpartnerCode::new("9900123456789"),
process_date: "2026-10-01".to_owned(),
};
let state = GpkeLfAnmeldungWorkflow::apply(LfAnmeldungState::New, &initiated_event);
let accepted_event = LfAnmeldungEvent::AntwortReceived {
response_pid: Pruefidentifikator::new(55003).unwrap(),
accepted: true,
reason: None,
response_ref: MessageRef::new("REF-002"),
};
let state = GpkeLfAnmeldungWorkflow::apply(state, &accepted_event);
assert!(matches!(state, LfAnmeldungState::Active(_)));
let cmd = LfAnmeldungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: "nb-response-window".into(),
};
let out = GpkeLfAnmeldungWorkflow::handle(&state, cmd).unwrap();
assert_eq!(out.events.len(), 0, "timeout is no-op on Active");
}
}