use mako_engine::types::Pruefidentifikator;
use mako_engine::{
deadline::Deadline,
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, PendingDeadline, Workflow, WorkflowOutput},
};
use mako_fristen::{APERAK_STROM_WINDOW_LABEL, aperak_strom_due_at};
pub const WORKFLOW_NAME: &str = "gpke-eog";
pub const EOG_ANMELDUNG_PID: u32 = 55013;
pub const EOG_PIDS: &[u32] = &[55013, 55014, 55015];
pub const EOG_ANTWORT_PIDS: &[u32] = &[55014, 55015];
pub const EOG_RESPONSE_WINDOW_LABEL: &str = "gpke-eog-response-window";
#[must_use]
pub fn eog_response_pid(accepted: bool) -> u32 {
if accepted { 55014 } else { 55015 }
}
#[must_use]
pub fn eog_antwort_due_at(
pid: u32,
received_at: time::OffsetDateTime,
) -> Option<time::OffsetDateTime> {
mako_fristen::antwort::antwort_deadline(pid, received_at)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum Versorgungsart {
Ersatzversorgung,
Grundversorgung,
Ersatzbelieferung,
}
pub const UEBERGANGSVERSORGUNG: &str = "ZZD";
impl Versorgungsart {
#[must_use]
pub fn code(self) -> &'static str {
match self {
Self::Ersatzversorgung => "ZC9",
Self::Grundversorgung => "ZD0",
Self::Ersatzbelieferung => "ZE3",
}
}
#[must_use]
pub fn from_code(code: &str) -> Option<Self> {
match code {
"ZC9" => Some(Self::Ersatzversorgung),
"ZD0" => Some(Self::Grundversorgung),
"ZE3" => Some(Self::Ersatzbelieferung),
_ => None,
}
}
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Self::Ersatzversorgung => "ERSATZVERSORGUNG",
Self::Grundversorgung => "GRUNDVERSORGUNG",
Self::Ersatzbelieferung => "ERSATZBELIEFERUNG",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum EogEvent {
Angemeldet {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
process_date: String,
pruefidentifikator: Pruefidentifikator,
transaktionsgrund: String,
haushaltskunde: Option<bool>,
},
AntwortErhalten {
response_pid: Pruefidentifikator,
accepted: bool,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<String>,
reason: Option<String>,
},
ZugeordnetOhneAntwort {
deadline_id: DeadlineId,
},
AnmeldungErhalten {
location_id: MaLo,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
document_date: String,
process_date: String,
message_ref: MessageRef,
pruefidentifikator: Pruefidentifikator,
transaktionsgrund: String,
haushaltskunde: Option<bool>,
},
ValidationPassed {
message_ref: MessageRef,
},
AntwortGesendet {
response_pid: Pruefidentifikator,
accepted: bool,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<String>,
reason: Option<String>,
},
Rejected {
reason: String,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for EogEvent {
fn event_type(&self) -> &'static str {
match self {
Self::Angemeldet { .. } => "EogAngemeldet",
Self::AntwortErhalten { .. } => "EogAntwortErhalten",
Self::ZugeordnetOhneAntwort { .. } => "EogZugeordnetOhneAntwort",
Self::AnmeldungErhalten { .. } => "EogAnmeldungErhalten",
Self::ValidationPassed { .. } => "EogValidationPassed",
Self::AntwortGesendet { .. } => "EogAntwortGesendet",
Self::Rejected { .. } => "EogRejected",
Self::DeadlineExpired { .. } => "EogDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct EogData {
pub location_id: MaLo,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub process_date: String,
pub pruefidentifikator: Pruefidentifikator,
pub transaktionsgrund: String,
pub haushaltskunde: Option<bool>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum EogState {
#[default]
New,
Angemeldet(EogData),
Zugeordnet {
data: EogData,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<String>,
ohne_antwort: bool,
},
Abgelehnt {
reason: String,
},
Eingegangen(EogData),
ValidationPassed(EogData),
AntwortGesendet {
data: EogData,
response_pid: Pruefidentifikator,
accepted: bool,
versorgungsart: Option<Versorgungsart>,
},
Rejected {
reason: String,
},
}
impl mako_engine::workflow::OccupiesBusinessKey for EogState {
fn occupies_business_key(&self) -> bool {
match self {
Self::Angemeldet(_) | Self::Zugeordnet { .. } => true,
Self::Eingegangen(_) | Self::ValidationPassed(_) | Self::AntwortGesendet { .. } => true,
Self::New | Self::Abgelehnt { .. } | Self::Rejected { .. } => false,
}
}
}
impl EogState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::Angemeldet(_) => "Angemeldet",
Self::Zugeordnet { .. } => "Zugeordnet",
Self::Abgelehnt { .. } => "Abgelehnt",
Self::Eingegangen(_) => "Eingegangen",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AntwortGesendet { .. } => "AntwortGesendet",
Self::Rejected { .. } => "Rejected",
}
}
#[must_use]
pub fn data(&self) -> Option<&EogData> {
match self {
Self::Angemeldet(d) | Self::Eingegangen(d) | Self::ValidationPassed(d) => Some(d),
Self::Zugeordnet { data, .. } | Self::AntwortGesendet { data, .. } => Some(data),
Self::New | Self::Abgelehnt { .. } | Self::Rejected { .. } => None,
}
}
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::Zugeordnet { .. }
| Self::Abgelehnt { .. }
| Self::AntwortGesendet { .. }
| Self::Rejected { .. }
)
}
}
#[derive(Clone)]
pub enum EogCommand {
Anmelden {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
process_date: String,
transaktionsgrund: String,
haushaltskunde: Option<bool>,
},
ReceiveAntwort {
response_pid: Pruefidentifikator,
accepted: bool,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<String>,
reason: Option<String>,
},
ReceiveAnmeldung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
location_id: MaLo,
document_date: String,
process_date: String,
message_ref: MessageRef,
transaktionsgrund: String,
haushaltskunde: Option<bool>,
validation_passed: bool,
validation_errors: Vec<String>,
received_at: time::OffsetDateTime,
},
SendAntwort {
accepted: bool,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<String>,
reason: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for EogCommand {}
fn process_completed_outbox(
data: &EogData,
versorgungsart: Option<Versorgungsart>,
bilanzkreis: Option<&str>,
ohne_antwort: bool,
) -> PendingOutbox {
let art = versorgungsart.unwrap_or(Versorgungsart::Ersatzversorgung);
PendingOutbox::new(
"ProcessCompleted",
"",
serde_json::json!({
"pid": EOG_ANMELDUNG_PID,
"malo_id": data.location_id.as_str(),
"new_supplier": data.receiver.as_str(),
"grid_operator": data.sender.as_str(),
"process_date": data.process_date,
"eog_art": art.as_str(),
"transaktionsgrund": data.transaktionsgrund,
"haushaltskunde": data.haushaltskunde,
"bilanzkreis": bilanzkreis,
"ohne_antwort": ohne_antwort,
}),
)
}
pub struct GpkeEogWorkflow;
impl Workflow for GpkeEogWorkflow {
type State = EogState;
type Event = EogEvent;
type Command = EogCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(
EOG_RESPONSE_WINDOW_LABEL | APERAK_STROM_WINDOW_LABEL,
EogState::Angemeldet(_) | EogState::Eingegangen(_) | EogState::ValidationPassed(_),
) => Some(EogCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
EogEvent::Angemeldet {
location_id,
sender,
receiver,
process_date,
pruefidentifikator,
transaktionsgrund,
haushaltskunde,
} => EogState::Angemeldet(EogData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
process_date: process_date.clone(),
pruefidentifikator: *pruefidentifikator,
transaktionsgrund: transaktionsgrund.clone(),
haushaltskunde: *haushaltskunde,
}),
EogEvent::AntwortErhalten {
accepted,
versorgungsart,
bilanzkreis,
reason,
..
} => match state {
EogState::Angemeldet(data) => {
if *accepted {
EogState::Zugeordnet {
data,
versorgungsart: *versorgungsart,
bilanzkreis: bilanzkreis.clone(),
ohne_antwort: false,
}
} else {
EogState::Abgelehnt {
reason: reason
.clone()
.unwrap_or_else(|| "EoG Zuordnung abgelehnt".to_owned()),
}
}
}
other => other,
},
EogEvent::ZugeordnetOhneAntwort { .. } => match state {
EogState::Angemeldet(data) => EogState::Zugeordnet {
data,
versorgungsart: None,
bilanzkreis: None,
ohne_antwort: true,
},
other => other,
},
EogEvent::AnmeldungErhalten {
location_id,
sender,
receiver,
process_date,
pruefidentifikator,
transaktionsgrund,
haushaltskunde,
..
} => EogState::Eingegangen(EogData {
location_id: location_id.clone(),
sender: sender.clone(),
receiver: receiver.clone(),
process_date: process_date.clone(),
pruefidentifikator: *pruefidentifikator,
transaktionsgrund: transaktionsgrund.clone(),
haushaltskunde: *haushaltskunde,
}),
EogEvent::ValidationPassed { .. } => match state {
EogState::Eingegangen(data) => EogState::ValidationPassed(data),
other => other,
},
EogEvent::AntwortGesendet {
response_pid,
accepted,
versorgungsart,
..
} => match state {
EogState::ValidationPassed(data) => EogState::AntwortGesendet {
data,
response_pid: *response_pid,
accepted: *accepted,
versorgungsart: *versorgungsart,
},
other => other,
},
EogEvent::Rejected { reason } => EogState::Rejected {
reason: reason.clone(),
},
EogEvent::DeadlineExpired { label, .. } => {
if state.is_terminal() {
state
} else {
EogState::Rejected {
reason: format!("deadline expired: {label}"),
}
}
}
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
EogCommand::Anmelden {
pid,
sender,
receiver,
location_id,
process_date,
transaktionsgrund,
haushaltskunde,
} => {
if !matches!(state, EogState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if pid.as_u32() != EOG_ANMELDUNG_PID {
return Err(WorkflowError::rejected(format!(
"expected EoG Anmeldung PID ({EOG_ANMELDUNG_PID}), got {pid}",
)));
}
if transaktionsgrund.trim().is_empty() {
return Err(WorkflowError::rejected(
"EoG Anmeldung requires a Transaktionsgrund (SG4 STS DE9013)".to_owned(),
));
}
let utilmd = 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,
}),
);
let initiated = PendingOutbox::new(
"ProcessInitiated",
receiver.as_str(),
serde_json::json!({
"pid": pid.as_u32(),
"malo_id": location_id.as_str(),
"new_supplier": receiver.as_str(),
"grid_operator": sender.as_str(),
"process_date": process_date,
"transaktionsgrund": transaktionsgrund,
}),
);
let event = EogEvent::Angemeldet {
location_id,
sender,
receiver,
process_date,
pruefidentifikator: pid,
transaktionsgrund,
haushaltskunde,
};
Ok(WorkflowOutput::with_outbox(
vec![event],
vec![utilmd, initiated],
))
}
EogCommand::ReceiveAntwort {
response_pid,
accepted,
versorgungsart,
bilanzkreis,
reason,
} => {
let data = match state {
EogState::Angemeldet(d) => d,
_ => {
return Err(WorkflowError::invalid_state("Angemeldet", state.label()));
}
};
if !EOG_ANTWORT_PIDS.contains(&response_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected EoG Antwort PID (55014/55015), got {response_pid}",
)));
}
let mut outbox = Vec::new();
if accepted {
outbox.push(process_completed_outbox(
data,
versorgungsart,
bilanzkreis.as_deref(),
false,
));
}
Ok(WorkflowOutput::with_outbox(
vec![EogEvent::AntwortErhalten {
response_pid,
accepted,
versorgungsart,
bilanzkreis,
reason,
}],
outbox,
))
}
EogCommand::ReceiveAnmeldung {
pid,
sender,
receiver,
location_id,
document_date,
process_date,
message_ref,
transaktionsgrund,
haushaltskunde,
validation_passed,
validation_errors,
received_at,
} => {
if !matches!(state, EogState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if pid.as_u32() != EOG_ANMELDUNG_PID {
return Err(WorkflowError::rejected(format!(
"expected EoG Anmeldung PID ({EOG_ANMELDUNG_PID}), 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 grund = transaktionsgrund.clone();
let mut events = vec![EogEvent::AnmeldungErhalten {
location_id,
sender,
receiver,
document_date,
process_date,
message_ref: message_ref.clone(),
pruefidentifikator: pid,
transaktionsgrund,
haushaltskunde,
}];
if validation_passed {
events.push(EogEvent::ValidationPassed {
message_ref: message_ref.clone(),
});
let outbox = vec![
crate::LfVorgangsdaten {
transaktionsgrund: Some(grund.clone()),
..crate::LfVorgangsdaten::default()
}
.process_initiated(
pid,
¬ify_malo,
&sender_mp_id,
&receiver_gln,
¬ify_termin,
&serde_json::json!({ "haushaltskunde": haushaltskunde }),
)
.caused_by(1),
PendingOutbox::aperak_anerkennung(
receiver_gln.as_str(),
sender_mp_id.as_str(),
message_ref.as_str(),
)
.caused_by(1),
];
let deadlines: Vec<PendingDeadline> = core::iter::once(PendingDeadline::new(
APERAK_STROM_WINDOW_LABEL,
aperak_strom_due_at(received_at),
))
.chain(
eog_antwort_due_at(pid.as_u32(), received_at)
.map(|due| PendingDeadline::new(EOG_RESPONSE_WINDOW_LABEL, due)),
)
.collect();
Ok(WorkflowOutput::with_outbox_and_deadlines(
events, outbox, deadlines,
))
} else {
let reason = if validation_errors.is_empty() {
"AHB validation failed".to_owned()
} else {
validation_errors.join("; ")
};
events.push(EogEvent::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))
}
}
EogCommand::SendAntwort {
accepted,
versorgungsart,
bilanzkreis,
reason,
} => {
let data = match state {
EogState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.label(),
));
}
};
if accepted && versorgungsart.is_none() {
return Err(WorkflowError::rejected(
"Bestätigung EOG requires the Versorgungsart (CCI+Z36: ZC9/ZD0/ZE3)"
.to_owned(),
));
}
if !accepted && reason.is_none() {
return Err(WorkflowError::rejected(
"Ablehnung EOG requires a reason (EBD E_0615: A02/A04/A05)".to_owned(),
));
}
let response_pid = Pruefidentifikator::new(eog_response_pid(accepted))
.map_err(|e| WorkflowError::rejected(e.clone()))?;
let mut outbox = vec![PendingOutbox::new(
"UTILMD",
data.sender.as_str(),
serde_json::json!({
"direction": "outbound",
"pid": response_pid.as_u32(),
"sender": data.receiver.as_str(),
"receiver": data.sender.as_str(),
"malo": data.location_id.as_str(),
"process_date": data.process_date,
"versorgungsart": versorgungsart.map(Versorgungsart::code),
"bilanzkreis": bilanzkreis,
"reason": reason,
}),
)];
if accepted {
outbox.push(process_completed_outbox(
data,
versorgungsart,
bilanzkreis.as_deref(),
false,
));
}
Ok(WorkflowOutput::with_outbox(
vec![EogEvent::AntwortGesendet {
response_pid,
accepted,
versorgungsart,
bilanzkreis,
reason,
}],
outbox,
))
}
EogCommand::TimeoutExpired { deadline_id, label } => match state {
EogState::Angemeldet(data) if label.as_ref() == EOG_RESPONSE_WINDOW_LABEL => {
Ok(WorkflowOutput::with_outbox(
vec![EogEvent::ZugeordnetOhneAntwort { deadline_id }],
vec![process_completed_outbox(data, None, None, true)],
))
}
s if s.is_terminal() => Ok(vec![].into()),
_ => Ok(vec![EogEvent::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 now() -> time::OffsetDateTime {
time::macros::datetime!(2026-07-01 10:00:00 UTC)
}
fn anmelden_cmd() -> EogCommand {
EogCommand::Anmelden {
pid: pid(55013),
sender: mcod("9900357000004"),
receiver: mcod("9900357000011"),
location_id: malo("51238696781"),
process_date: "20260615".to_owned(), transaktionsgrund: "ZT7".to_owned(), haushaltskunde: Some(true),
}
}
fn receive_anmeldung_cmd(ok: bool) -> EogCommand {
EogCommand::ReceiveAnmeldung {
pid: pid(55013),
sender: mcod("9900357000004"),
receiver: mcod("9900357000011"),
location_id: malo("51238696781"),
document_date: "20260701".to_owned(),
process_date: "20260615".to_owned(),
message_ref: mref("EOG-001"),
transaktionsgrund: "ZT7".to_owned(),
haushaltskunde: Some(true),
validation_passed: ok,
validation_errors: if ok {
vec![]
} else {
vec!["missing mandatory segment".to_owned()]
},
received_at: now(),
}
}
fn apply_all(init: EogState, events: &[EogEvent]) -> EogState {
events.iter().fold(init, GpkeEogWorkflow::apply)
}
#[test]
fn initiator_happy_path_zugeordnet() {
let out = GpkeEogWorkflow::handle(&EogState::New, anmelden_cmd()).unwrap();
assert_eq!(out.events.len(), 1);
assert_eq!(out.outbox.len(), 2);
assert_eq!(out.outbox[0].message_type.as_ref(), "UTILMD");
assert_eq!(out.outbox[0].payload["transaktionsgrund"], "ZT7");
assert_eq!(out.outbox[1].message_type.as_ref(), "ProcessInitiated");
let state = apply_all(EogState::New, &out.events);
assert!(matches!(state, EogState::Angemeldet(_)));
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::ReceiveAntwort {
response_pid: pid(55014),
accepted: true,
versorgungsart: Some(Versorgungsart::Ersatzversorgung),
bilanzkreis: Some("11XGRUNDV-BK--I".to_owned()),
reason: None,
},
)
.unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "ProcessCompleted");
let payload = &out.outbox[0].payload;
assert_eq!(payload["pid"], 55013);
assert_eq!(payload["eog_art"], "ERSATZVERSORGUNG");
assert_eq!(payload["process_date"], "20260615");
assert_eq!(payload["ohne_antwort"], false);
let state = apply_all(state, &out.events);
assert!(matches!(
state,
EogState::Zugeordnet {
versorgungsart: Some(Versorgungsart::Ersatzversorgung),
ohne_antwort: false,
..
}
));
}
#[test]
fn initiator_grundversorgung_classification_from_antwort() {
let out = GpkeEogWorkflow::handle(&EogState::New, anmelden_cmd()).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::ReceiveAntwort {
response_pid: pid(55014),
accepted: true,
versorgungsart: Some(Versorgungsart::Grundversorgung),
bilanzkreis: Some("11XGRUNDV-BK--I".to_owned()),
reason: None,
},
)
.unwrap();
assert_eq!(out.outbox[0].payload["eog_art"], "GRUNDVERSORGUNG");
}
#[test]
fn initiator_ablehnung() {
let out = GpkeEogWorkflow::handle(&EogState::New, anmelden_cmd()).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::ReceiveAntwort {
response_pid: pid(55015),
accepted: false,
versorgungsart: None,
bilanzkreis: None,
reason: Some("A02".to_owned()),
},
)
.unwrap();
assert!(out.outbox.is_empty());
let state = apply_all(state, &out.events);
assert!(matches!(state, EogState::Abgelehnt { .. }));
}
#[test]
fn initiator_timeout_assigns_with_default_bk() {
let out = GpkeEogWorkflow::handle(&EogState::New, anmelden_cmd()).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: EOG_RESPONSE_WINDOW_LABEL.into(),
},
)
.unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "ProcessCompleted");
assert_eq!(out.outbox[0].payload["eog_art"], "ERSATZVERSORGUNG");
assert_eq!(out.outbox[0].payload["ohne_antwort"], true);
let state = apply_all(state, &out.events);
assert!(matches!(
state,
EogState::Zugeordnet {
versorgungsart: None,
ohne_antwort: true,
..
}
));
}
#[test]
fn initiator_requires_transaktionsgrund() {
let result = GpkeEogWorkflow::handle(
&EogState::New,
EogCommand::Anmelden {
pid: pid(55013),
sender: mcod("9900357000004"),
receiver: mcod("9900357000011"),
location_id: malo("51238696781"),
process_date: "20260701".to_owned(),
transaktionsgrund: " ".to_owned(),
haushaltskunde: None,
},
);
assert!(result.is_err());
}
#[test]
fn initiator_wrong_pid_rejected() {
let result = GpkeEogWorkflow::handle(
&EogState::New,
EogCommand::Anmelden {
pid: pid(55001),
sender: mcod("9900357000004"),
receiver: mcod("9900357000011"),
location_id: malo("51238696781"),
process_date: "20260701".to_owned(),
transaktionsgrund: "ZT7".to_owned(),
haushaltskunde: None,
},
);
assert!(result.is_err());
}
#[test]
fn responder_happy_path_bestaetigung() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
assert_eq!(out.events.len(), 2); assert_eq!(out.deadlines.len(), 2); let state = apply_all(EogState::New, &out.events);
assert!(matches!(state, EogState::ValidationPassed(_)));
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::SendAntwort {
accepted: true,
versorgungsart: Some(Versorgungsart::Grundversorgung),
bilanzkreis: Some("11XGRUNDV-BK--I".to_owned()),
reason: None,
},
)
.unwrap();
assert_eq!(out.outbox.len(), 2);
assert_eq!(out.outbox[0].message_type.as_ref(), "UTILMD");
assert_eq!(out.outbox[0].payload["pid"], 55014);
assert_eq!(out.outbox[0].payload["versorgungsart"], "ZD0");
assert_eq!(out.outbox[1].message_type.as_ref(), "ProcessCompleted");
assert_eq!(out.outbox[1].payload["eog_art"], "GRUNDVERSORGUNG");
let state = apply_all(state, &out.events);
assert!(
matches!(state, EogState::AntwortGesendet { response_pid, accepted: true, .. }
if response_pid.as_u32() == 55014)
);
}
#[test]
fn responder_bestaetigung_requires_versorgungsart() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
let state = apply_all(EogState::New, &out.events);
let result = GpkeEogWorkflow::handle(
&state,
EogCommand::SendAntwort {
accepted: true,
versorgungsart: None,
bilanzkreis: None,
reason: None,
},
);
assert!(result.is_err());
}
#[test]
fn responder_ablehnung_requires_reason() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
let state = apply_all(EogState::New, &out.events);
let result = GpkeEogWorkflow::handle(
&state,
EogCommand::SendAntwort {
accepted: false,
versorgungsart: None,
bilanzkreis: None,
reason: None,
},
);
assert!(result.is_err());
}
#[test]
fn responder_ablehnung() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::SendAntwort {
accepted: false,
versorgungsart: None,
bilanzkreis: None,
reason: Some("A05".to_owned()),
},
)
.unwrap();
assert_eq!(out.outbox.len(), 1); assert_eq!(out.outbox[0].payload["pid"], 55015);
let state = apply_all(state, &out.events);
assert!(
matches!(state, EogState::AntwortGesendet { response_pid, accepted: false, .. }
if response_pid.as_u32() == 55015)
);
}
#[test]
fn responder_validation_failure_rejects_with_aperak() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(false)).unwrap();
assert_eq!(out.outbox.len(), 1);
assert_eq!(out.outbox[0].message_type.as_ref(), "APERAK");
let state = apply_all(EogState::New, &out.events);
assert!(matches!(state, EogState::Rejected { .. }));
}
#[test]
fn responder_timeout_rejects() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: EOG_RESPONSE_WINDOW_LABEL.into(),
},
)
.unwrap();
let state = apply_all(state, &out.events);
assert!(matches!(state, EogState::Rejected { .. }));
}
#[test]
fn an_inbound_anmeldung_notifies_the_eog_supplier() {
let out = GpkeEogWorkflow::handle(&EogState::New, receive_anmeldung_cmd(true)).unwrap();
let notification = out
.outbox
.iter()
.find(|o| &*o.message_type == "ProcessInitiated")
.expect("the E/G is notified of the Zuordnung");
assert_eq!(notification.payload["pid"], EOG_ANMELDUNG_PID);
assert_eq!(notification.payload["malo_id"], "51238696781");
assert!(notification.payload.get("transaktionsgrund").is_some());
assert_eq!(&*notification.recipient, "9900357000011");
}
#[test]
fn versorgungsart_code_roundtrip() {
for art in [
Versorgungsart::Ersatzversorgung,
Versorgungsart::Grundversorgung,
Versorgungsart::Ersatzbelieferung,
] {
assert_eq!(Versorgungsart::from_code(art.code()), Some(art));
}
assert_eq!(Versorgungsart::from_code("E06"), None);
}
#[test]
fn zzd_is_not_a_versorgungsart() {
assert_eq!(Versorgungsart::from_code(UEBERGANGSVERSORGUNG), None);
assert!(
[
Versorgungsart::Ersatzversorgung,
Versorgungsart::Grundversorgung,
Versorgungsart::Ersatzbelieferung,
]
.iter()
.all(|a| a.code() != UEBERGANGSVERSORGUNG)
);
}
#[test]
fn timeout_in_terminal_state_is_noop() {
let out = GpkeEogWorkflow::handle(&EogState::New, anmelden_cmd()).unwrap();
let state = apply_all(EogState::New, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::ReceiveAntwort {
response_pid: pid(55014),
accepted: true,
versorgungsart: Some(Versorgungsart::Ersatzversorgung),
bilanzkreis: None,
reason: None,
},
)
.unwrap();
let state = apply_all(state, &out.events);
let out = GpkeEogWorkflow::handle(
&state,
EogCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: EOG_RESPONSE_WINDOW_LABEL.into(),
},
)
.unwrap();
assert!(out.events.is_empty());
}
#[test]
fn the_code_matches_the_wire_table() {
assert_eq!(
super::UEBERGANGSVERSORGUNG,
edi_energy::utilmd_codes::transaktionsgrund::UEBERGANGSVERSORGUNG
);
}
}