use mako_engine::{
deadline::Deadline,
error::WorkflowError,
ids::DeadlineId,
types::{MaLo, MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const ANFRAGE_PID: u32 = 55555;
pub const WORKFLOW_NAME: &str = "gpke-anfrage-bestellung";
pub const ANFRAGE_WINDOW_LABEL: &str = "gpke-anfrage-bestellung-24h";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum AnfrageBestellungEvent {
AnfrageErhalten {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
vorgang_id: MaLo,
bearbeitungsstatus: String,
document_date: String,
message_ref: MessageRef,
},
ValidationPassed {
message_ref: MessageRef,
},
ValidationFailed {
errors: Vec<String>,
},
ResponseDispatched {
data_provided: bool,
reason: Option<String>,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for AnfrageBestellungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnfrageErhalten { .. } => "AnfrageBestellungErhalten",
Self::ValidationPassed { .. } => "AnfrageBestellungValidationPassed",
Self::ValidationFailed { .. } => "AnfrageBestellungValidationFailed",
Self::ResponseDispatched { .. } => "AnfrageBestellungResponseDispatched",
Self::DeadlineExpired { .. } => "AnfrageBestellungDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct AnfrageData {
pub pruefidentifikator: Pruefidentifikator,
pub sender: MarktpartnerCode,
pub receiver: MarktpartnerCode,
pub vorgang_id: MaLo,
pub bearbeitungsstatus: String,
pub document_date: String,
pub message_ref: Option<MessageRef>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum AnfrageBestellungState {
New,
Initiated(AnfrageData),
ValidationPassed(AnfrageData),
ResponseDispatched(AnfrageData),
Rejected {
reason: String,
},
}
impl Default for AnfrageBestellungState {
fn default() -> Self {
Self::New
}
}
impl AnfrageBestellungState {
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::Initiated(_) => "Initiated",
Self::ValidationPassed(_) => "ValidationPassed",
Self::ResponseDispatched(_) => "ResponseDispatched",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum AnfrageBestellungCommand {
ReceiveAnfrage {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
receiver: MarktpartnerCode,
vorgang_id: MaLo,
bearbeitungsstatus: String,
document_date: String,
message_ref: MessageRef,
validation_passed: bool,
validation_errors: Vec<String>,
},
DispatchResponse {
data_provided: bool,
reason: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for AnfrageBestellungCommand {}
pub struct GpkeAnfrageBestellungWorkflow;
impl Workflow for GpkeAnfrageBestellungWorkflow {
type State = AnfrageBestellungState;
type Event = AnfrageBestellungEvent;
type Command = AnfrageBestellungCommand;
fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
match (deadline.label(), state) {
(
ANFRAGE_WINDOW_LABEL,
AnfrageBestellungState::Initiated(_) | AnfrageBestellungState::ValidationPassed(_),
) => Some(AnfrageBestellungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
AnfrageBestellungEvent::AnfrageErhalten {
pruefidentifikator,
sender,
receiver,
vorgang_id,
bearbeitungsstatus,
document_date,
message_ref,
} => AnfrageBestellungState::Initiated(AnfrageData {
pruefidentifikator: *pruefidentifikator,
sender: sender.clone(),
receiver: receiver.clone(),
vorgang_id: vorgang_id.clone(),
bearbeitungsstatus: bearbeitungsstatus.clone(),
document_date: document_date.clone(),
message_ref: Some(message_ref.clone()),
}),
AnfrageBestellungEvent::ValidationPassed { .. } => match state {
AnfrageBestellungState::Initiated(data) => {
AnfrageBestellungState::ValidationPassed(data)
}
other => other,
},
AnfrageBestellungEvent::ValidationFailed { errors } => {
AnfrageBestellungState::Rejected {
reason: errors.join("; "),
}
}
AnfrageBestellungEvent::ResponseDispatched {
data_provided,
reason,
} => match state {
AnfrageBestellungState::ValidationPassed(data) => {
if *data_provided {
AnfrageBestellungState::ResponseDispatched(data)
} else {
AnfrageBestellungState::Rejected {
reason: reason
.clone()
.unwrap_or_else(|| "Anfrage abgelehnt".to_owned()),
}
}
}
other => other,
},
AnfrageBestellungEvent::DeadlineExpired { label, .. } => match state {
AnfrageBestellungState::ResponseDispatched(_)
| AnfrageBestellungState::Rejected { .. } => state,
_ => AnfrageBestellungState::Rejected {
reason: format!("deadline expired: {label}"),
},
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
AnfrageBestellungCommand::ReceiveAnfrage {
pid,
sender,
receiver,
vorgang_id,
bearbeitungsstatus,
document_date,
message_ref,
validation_passed,
validation_errors,
} => {
if !matches!(state, AnfrageBestellungState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
if pid.as_u32() != ANFRAGE_PID {
return Err(WorkflowError::rejected(format!(
"unsupported PID {pid} for GpkeAnfrageBestellungWorkflow \
(expected {ANFRAGE_PID})",
)));
}
let mut events = vec![AnfrageBestellungEvent::AnfrageErhalten {
pruefidentifikator: pid,
sender,
receiver,
vorgang_id,
bearbeitungsstatus,
document_date,
message_ref: message_ref.clone(),
}];
if validation_passed {
events.push(AnfrageBestellungEvent::ValidationPassed { message_ref });
} else {
events.push(AnfrageBestellungEvent::ValidationFailed {
errors: validation_errors,
});
}
Ok(WorkflowOutput::events(events))
}
AnfrageBestellungCommand::DispatchResponse {
data_provided,
reason,
} => {
if !matches!(state, AnfrageBestellungState::ValidationPassed(_)) {
return Err(WorkflowError::invalid_state(
"ValidationPassed",
state.status_str(),
));
}
Ok(WorkflowOutput::events(vec![
AnfrageBestellungEvent::ResponseDispatched {
data_provided,
reason,
},
]))
}
AnfrageBestellungCommand::TimeoutExpired { deadline_id, label } => match state {
AnfrageBestellungState::ResponseDispatched(_)
| AnfrageBestellungState::Rejected { .. } => Ok(WorkflowOutput::events(vec![])),
_ => Ok(WorkflowOutput::events(vec![
AnfrageBestellungEvent::DeadlineExpired { deadline_id, label },
])),
},
}
}
}
#[cfg(test)]
mod tests {
use mako_engine::{
ids::DeadlineId,
types::{MaLo, MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{Workflow, WorkflowOutput},
};
use super::*;
fn pid() -> Pruefidentifikator {
Pruefidentifikator::new(55555).unwrap()
}
fn sender() -> MarktpartnerCode {
MarktpartnerCode::new("4012345000023")
}
fn receiver() -> MarktpartnerCode {
MarktpartnerCode::new("9900357000004")
}
fn vorgang_id() -> MaLo {
MaLo::new("DE0000000000000000000000000012345")
}
fn msg_ref() -> MessageRef {
MessageRef::new("MSG-ANFR-001")
}
fn receive_valid() -> AnfrageBestellungCommand {
AnfrageBestellungCommand::ReceiveAnfrage {
pid: pid(),
sender: sender(),
receiver: receiver(),
vorgang_id: vorgang_id(),
bearbeitungsstatus: "E07".to_owned(),
document_date: "20250115".to_owned(),
message_ref: msg_ref(),
validation_passed: true,
validation_errors: vec![],
}
}
fn receive_invalid() -> AnfrageBestellungCommand {
AnfrageBestellungCommand::ReceiveAnfrage {
pid: pid(),
sender: sender(),
receiver: receiver(),
vorgang_id: vorgang_id(),
bearbeitungsstatus: "E08".to_owned(),
document_date: "20250115".to_owned(),
message_ref: msg_ref(),
validation_passed: false,
validation_errors: vec!["segment BGM missing qualifier E03".to_owned()],
}
}
#[test]
fn receive_valid_transitions_to_validation_passed() {
let state = AnfrageBestellungState::New;
let result = GpkeAnfrageBestellungWorkflow::handle(&state, receive_valid()).unwrap();
let WorkflowOutput { events, .. } = result;
assert_eq!(events.len(), 2);
assert!(matches!(
events[0],
AnfrageBestellungEvent::AnfrageErhalten { .. }
));
assert!(matches!(
events[1],
AnfrageBestellungEvent::ValidationPassed { .. }
));
let final_state = events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply);
assert!(matches!(
final_state,
AnfrageBestellungState::ValidationPassed(_)
));
}
#[test]
fn receive_invalid_transitions_to_rejected() {
let state = AnfrageBestellungState::New;
let result = GpkeAnfrageBestellungWorkflow::handle(&state, receive_invalid()).unwrap();
let WorkflowOutput { events, .. } = result;
assert_eq!(events.len(), 2);
assert!(matches!(
events[0],
AnfrageBestellungEvent::AnfrageErhalten { .. }
));
assert!(matches!(
events[1],
AnfrageBestellungEvent::ValidationFailed { .. }
));
let final_state = events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply);
assert!(matches!(
final_state,
AnfrageBestellungState::Rejected { .. }
));
}
#[test]
fn receive_wrong_pid_is_rejected() {
let bad_pid = Pruefidentifikator::new(55001).unwrap();
let cmd = AnfrageBestellungCommand::ReceiveAnfrage {
pid: bad_pid,
sender: sender(),
receiver: receiver(),
vorgang_id: vorgang_id(),
bearbeitungsstatus: "E07".to_owned(),
document_date: "20250115".to_owned(),
message_ref: msg_ref(),
validation_passed: true,
validation_errors: vec![],
};
let err =
GpkeAnfrageBestellungWorkflow::handle(&AnfrageBestellungState::New, cmd).unwrap_err();
assert!(err.to_string().contains("55001"));
}
#[test]
fn receive_in_non_new_state_is_invalid_state_error() {
let init_state = AnfrageBestellungState::New;
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&init_state, receive_valid()).unwrap();
let initiated = events
.iter()
.fold(init_state, GpkeAnfrageBestellungWorkflow::apply);
let err = GpkeAnfrageBestellungWorkflow::handle(&initiated, receive_valid()).unwrap_err();
assert!(err.to_string().contains("New"));
}
fn make_validation_passed_state() -> AnfrageBestellungState {
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&AnfrageBestellungState::New, receive_valid())
.unwrap();
events.iter().fold(
AnfrageBestellungState::New,
GpkeAnfrageBestellungWorkflow::apply,
)
}
#[test]
fn dispatch_data_provided_transitions_to_response_dispatched() {
let state = make_validation_passed_state();
let cmd = AnfrageBestellungCommand::DispatchResponse {
data_provided: true,
reason: None,
};
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&state, cmd).unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
AnfrageBestellungEvent::ResponseDispatched {
data_provided: true,
..
}
));
let final_state = events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply);
assert!(matches!(
final_state,
AnfrageBestellungState::ResponseDispatched(_)
));
}
#[test]
fn dispatch_rejection_transitions_to_rejected() {
let state = make_validation_passed_state();
let cmd = AnfrageBestellungCommand::DispatchResponse {
data_provided: false,
reason: Some("Vorgang nicht gefunden".to_owned()),
};
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&state, cmd).unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
AnfrageBestellungEvent::ResponseDispatched {
data_provided: false,
..
}
));
let final_state = events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply);
assert!(matches!(
final_state,
AnfrageBestellungState::Rejected { .. }
));
}
#[test]
fn dispatch_response_from_wrong_state_returns_error() {
let initiated = {
let WorkflowOutput { events, .. } = GpkeAnfrageBestellungWorkflow::handle(
&AnfrageBestellungState::New,
receive_invalid(),
)
.unwrap();
events.iter().fold(
AnfrageBestellungState::New,
GpkeAnfrageBestellungWorkflow::apply,
)
};
let err = GpkeAnfrageBestellungWorkflow::handle(
&initiated,
AnfrageBestellungCommand::DispatchResponse {
data_provided: true,
reason: None,
},
)
.unwrap_err();
assert!(err.to_string().contains("ValidationPassed"));
}
#[test]
fn timeout_in_validation_passed_transitions_to_rejected() {
let state = make_validation_passed_state();
let cmd = AnfrageBestellungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: ANFRAGE_WINDOW_LABEL.into(),
};
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&state, cmd).unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
AnfrageBestellungEvent::DeadlineExpired { .. }
));
let final_state = events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply);
assert!(matches!(
final_state,
AnfrageBestellungState::Rejected { reason } if reason.contains("deadline")
));
}
#[test]
fn timeout_in_terminal_state_is_noop() {
let dispatched = {
let state = make_validation_passed_state();
let WorkflowOutput { events, .. } = GpkeAnfrageBestellungWorkflow::handle(
&state,
AnfrageBestellungCommand::DispatchResponse {
data_provided: true,
reason: None,
},
)
.unwrap();
events
.iter()
.fold(state, GpkeAnfrageBestellungWorkflow::apply)
};
let cmd = AnfrageBestellungCommand::TimeoutExpired {
deadline_id: DeadlineId::new(),
label: ANFRAGE_WINDOW_LABEL.into(),
};
let WorkflowOutput { events, .. } =
GpkeAnfrageBestellungWorkflow::handle(&dispatched, cmd).unwrap();
assert!(events.is_empty(), "no events emitted in terminal state");
}
#[test]
fn on_deadline_fires_for_initiated_and_validation_passed() {
use mako_engine::deadline::Deadline;
use mako_engine::ids::{ProcessId, StreamId, TenantId};
use mako_engine::version::WorkflowId;
use time::OffsetDateTime;
let stream_id = StreamId::new("test-stream-anfrage-001");
let due = OffsetDateTime::now_utc();
let dl = Deadline::new(
stream_id,
ProcessId::new(),
TenantId::from_party_id("9900357000004"),
WorkflowId::new(WORKFLOW_NAME, "FV2025-10-01"),
ANFRAGE_WINDOW_LABEL,
due,
);
let initiated = {
let e0 = AnfrageBestellungEvent::AnfrageErhalten {
pruefidentifikator: pid(),
sender: sender(),
receiver: receiver(),
vorgang_id: vorgang_id(),
bearbeitungsstatus: "E07".to_owned(),
document_date: "20250115".to_owned(),
message_ref: msg_ref(),
};
GpkeAnfrageBestellungWorkflow::apply(AnfrageBestellungState::New, &e0)
};
assert!(
GpkeAnfrageBestellungWorkflow::on_deadline(&dl, &initiated).is_some(),
"Initiated state must trigger TimeoutExpired"
);
let vp = make_validation_passed_state();
assert!(
GpkeAnfrageBestellungWorkflow::on_deadline(&dl, &vp).is_some(),
"ValidationPassed state must trigger TimeoutExpired"
);
let WorkflowOutput { events, .. } = GpkeAnfrageBestellungWorkflow::handle(
&vp,
AnfrageBestellungCommand::DispatchResponse {
data_provided: true,
reason: None,
},
)
.unwrap();
let dispatched = events.iter().fold(
make_validation_passed_state(),
GpkeAnfrageBestellungWorkflow::apply,
);
assert!(
GpkeAnfrageBestellungWorkflow::on_deadline(&dl, &dispatched).is_none(),
"terminal state must not fire deadline"
);
}
}