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-stornierung";
pub const STORNIERUNG_APERAK_WINDOW_LABEL: &str = "gpke-stornierung-aperak-24h";
pub const STORNIERUNG_PIDS: &[u32] = &[55022, 55023, 55024];
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum GpkeStornierungEvent {
StornierungReceived {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
vorgang_id: MaLo,
document_date: String,
message_ref: MessageRef,
},
ValidationPassed {
message_ref: MessageRef,
},
ValidationFailed {
errors: Vec<String>,
},
AperakDispatched {
positive: bool,
reason: Option<String>,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for GpkeStornierungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::StornierungReceived { .. } => "GpkeStornierungReceived",
Self::ValidationPassed { .. } => "GpkeStornierungValidationPassed",
Self::ValidationFailed { .. } => "GpkeStornierungValidationFailed",
Self::AperakDispatched { .. } => "GpkeStornierungAperakDispatched",
Self::DeadlineExpired { .. } => "GpkeStornierungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct GpkeStornierungData {
pub pruefidentifikator: Pruefidentifikator,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub vorgang_id: MaLo,
pub document_date: String,
pub message_ref: Option<MessageRef>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum GpkeStornierungState {
New,
Initiated(GpkeStornierungData),
ValidationPassed(GpkeStornierungData),
AperakSent(GpkeStornierungData),
Completed(GpkeStornierungData),
Rejected {
reason: String,
},
}
impl Default for GpkeStornierungState {
fn default() -> Self {
Self::New
}
}
impl GpkeStornierungState {
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::Initiated(_) => "Initiated",
Self::ValidationPassed(_) => "ValidationPassed",
Self::AperakSent(_) => "AperakSent",
Self::Completed(_) => "Completed",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum GpkeStornierungCommand {
ReceiveUtilmd {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
vorgang_id: MaLo,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
DispatchAperak {
positive: bool,
reason: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for GpkeStornierungCommand {}
pub struct GpkeStornierungWorkflow;
impl Workflow for GpkeStornierungWorkflow {
type State = GpkeStornierungState;
type Event = GpkeStornierungEvent;
type Command = GpkeStornierungCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
STORNIERUNG_APERAK_WINDOW_LABEL,
GpkeStornierungState::Initiated(_)
| GpkeStornierungState::ValidationPassed(_)
| GpkeStornierungState::AperakSent(_),
) => Some(GpkeStornierungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
GpkeStornierungEvent::StornierungReceived {
pruefidentifikator,
sender,
receiver,
vorgang_id,
document_date,
message_ref,
} => GpkeStornierungState::Initiated(GpkeStornierungData {
pruefidentifikator: *pruefidentifikator,
sender: sender.clone(),
receiver: receiver.clone(),
vorgang_id: vorgang_id.clone(),
document_date: document_date.clone(),
message_ref: Some(message_ref.clone()),
}),
GpkeStornierungEvent::ValidationPassed { .. } => {
if let GpkeStornierungState::Initiated(data) = state {
GpkeStornierungState::ValidationPassed(data)
} else {
state
}
}
GpkeStornierungEvent::ValidationFailed { errors } => GpkeStornierungState::Rejected {
reason: errors.join("; "),
},
GpkeStornierungEvent::AperakDispatched { positive, reason } => match state {
GpkeStornierungState::ValidationPassed(data) => {
if *positive {
GpkeStornierungState::Completed(data)
} else {
GpkeStornierungState::Rejected {
reason: reason
.clone()
.unwrap_or_else(|| "negative APERAK".to_owned()),
}
}
}
_ => state,
},
GpkeStornierungEvent::DeadlineExpired { label, .. } => match state {
GpkeStornierungState::Completed(_) | GpkeStornierungState::Rejected { .. } => state,
_ => GpkeStornierungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
GpkeStornierungCommand::ReceiveUtilmd {
pid,
sender,
receiver,
vorgang_id,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, GpkeStornierungState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if !STORNIERUNG_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"unsupported GPKE Stornierung PID {pid} (expected one of: {STORNIERUNG_PIDS:?})",
)));
}
let mut events = vec![GpkeStornierungEvent::StornierungReceived {
pruefidentifikator: pid,
sender,
receiver,
vorgang_id,
document_date,
message_ref: message_ref.clone(),
}];
if validation_passed {
events.push(GpkeStornierungEvent::ValidationPassed { message_ref });
} else {
events.push(GpkeStornierungEvent::ValidationFailed {
errors: validation_errors,
});
}
Ok(WorkflowOutput::events(events))
}
GpkeStornierungCommand::DispatchAperak { positive, reason } => {
let data = match state {
GpkeStornierungState::ValidationPassed(d) => d,
_ => {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
};
let mut aperak_payload = serde_json::json!({
"sender": data.receiver.as_str(),
"receiver": data.sender.as_str(),
"pid": 29001_u32,
"document_code": if positive { "312" } else { "313" },
});
if !positive {
aperak_payload["error_code"] = serde_json::Value::String("Z29".to_owned());
}
if let Some(ref r) = reason {
aperak_payload["reason"] = serde_json::Value::String(r.clone());
}
let outbox = vec![
PendingOutbox::new("APERAK", data.sender.as_str(), aperak_payload).caused_by(0),
];
Ok(WorkflowOutput::with_outbox(
vec![GpkeStornierungEvent::AperakDispatched { positive, reason }],
outbox,
))
}
GpkeStornierungCommand::TimeoutExpired { deadline_id, label } => match state {
GpkeStornierungState::Completed(_) | GpkeStornierungState::Rejected { .. } => {
Ok(WorkflowOutput::events(vec![]))
}
_ => Ok(WorkflowOutput::events(vec![
GpkeStornierungEvent::DeadlineExpired { deadline_id, label },
])),
},
}
}
}
#[cfg(test)]
mod tests {
use mako_engine::ids::DeadlineId;
use super::*;
fn pid(n: u32) -> Pruefidentifikator {
Pruefidentifikator::new(n).expect("valid PID")
}
fn stornierung_cmd(p: u32, validation_passed: bool) -> GpkeStornierungCommand {
GpkeStornierungCommand::ReceiveUtilmd {
pid: pid(p),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9907317000007"),
vorgang_id: MaLo::new("STORNO0000A"),
document_date: "20251001000000+00".to_owned(),
message_ref: MessageRef::new("00001"),
validation_passed,
validation_errors: if validation_passed {
vec![]
} else {
vec!["missing DTM".to_owned()]
},
}
}
#[test]
fn receive_valid_55022_transitions_to_validation_passed() {
let state = GpkeStornierungState::New;
let output = GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55022, true)).unwrap();
assert_eq!(
output.events.len(),
2,
"StornierungReceived + ValidationPassed"
);
assert!(matches!(
output.events[0],
GpkeStornierungEvent::StornierungReceived { .. }
));
assert!(matches!(
output.events[1],
GpkeStornierungEvent::ValidationPassed { .. }
));
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(matches!(state, GpkeStornierungState::ValidationPassed(_)));
}
#[test]
fn receive_invalid_55022_transitions_to_rejected() {
let state = GpkeStornierungState::New;
let output =
GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55022, false)).unwrap();
assert_eq!(
output.events.len(),
2,
"StornierungReceived + ValidationFailed"
);
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(matches!(state, GpkeStornierungState::Rejected { .. }));
}
#[test]
fn dispatch_positive_aperak_completes_process() {
let state = GpkeStornierungState::New;
let output = GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55022, true)).unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(matches!(state, GpkeStornierungState::ValidationPassed(_)));
let output = GpkeStornierungWorkflow::handle(
&state,
GpkeStornierungCommand::DispatchAperak {
positive: true,
reason: None,
},
)
.unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(
matches!(state, GpkeStornierungState::Completed(_)),
"positive APERAK must complete stornierung; got: {state:?}"
);
}
#[test]
fn dispatch_negative_aperak_rejects_process() {
let state = GpkeStornierungState::New;
let output = GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55022, true)).unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
let output = GpkeStornierungWorkflow::handle(
&state,
GpkeStornierungCommand::DispatchAperak {
positive: false,
reason: Some("Vorgang nicht stornierbar".to_owned()),
},
)
.unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(matches!(state, GpkeStornierungState::Rejected { .. }));
}
#[test]
fn deadline_expired_rejects_initiated_process() {
let state = GpkeStornierungState::New;
let output = GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55022, true)).unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
let output = GpkeStornierungWorkflow::handle(
&state,
GpkeStornierungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: STORNIERUNG_APERAK_WINDOW_LABEL.into(),
},
)
.unwrap();
let state = output
.events
.iter()
.fold(state, GpkeStornierungWorkflow::apply);
assert!(matches!(state, GpkeStornierungState::Rejected { .. }));
}
#[test]
fn reject_unsupported_pid() {
let state = GpkeStornierungState::New;
let result = GpkeStornierungWorkflow::handle(&state, stornierung_cmd(55021, true));
assert!(result.is_err(), "PID 55021 must be rejected");
}
#[test]
fn deadline_is_noop_for_completed_process() {
let data = GpkeStornierungData {
pruefidentifikator: pid(55022),
sender: MarktpartnerCode::new("4012345000023"),
receiver: MarktpartnerCode::new("9907317000007"),
vorgang_id: MaLo::new("STORNO0000A"),
document_date: "20251001000000+00".to_owned(),
message_ref: None,
};
let state = GpkeStornierungState::Completed(data);
let output = GpkeStornierungWorkflow::handle(
&state,
GpkeStornierungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: STORNIERUNG_APERAK_WINDOW_LABEL.into(),
},
)
.unwrap();
assert!(
output.events.is_empty(),
"deadline must be no-op for Completed state"
);
}
}