use mako_engine::types::Pruefidentifikator;
use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-insrpt";
pub const INSRPT_ANFRAGE_PIDS: &[u32] = &[23001];
pub const INSRPT_ANTWORT_PIDS: &[u32] = &[23003, 23004, 23008, 23011, 23012];
pub const ANTWORT_WINDOW_LABEL: &str = "wim-insrpt-antwort";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StorungsmeldungData {
pub pruefidentifikator: Pruefidentifikator,
pub msb_mp_id: MarktpartnerCode,
pub document_date: String,
pub message_ref: MessageRef,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum StorungsmeldungEvent {
StorungsmeldungGesendet {
pruefidentifikator: Pruefidentifikator,
msb_mp_id: MarktpartnerCode,
document_date: String,
message_ref: MessageRef,
},
AntwortErhalten {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
is_confirmation: bool,
message_ref: MessageRef,
},
InformationsmeldungErhalten {
pruefidentifikator: Pruefidentifikator,
sender: MarktpartnerCode,
message_ref: MessageRef,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for StorungsmeldungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::StorungsmeldungGesendet { .. } => "InsrptStorungsmeldungGesendet",
Self::AntwortErhalten { .. } => "InsrptAntwortErhalten",
Self::InformationsmeldungErhalten { .. } => "InsrptInformationsmeldungErhalten",
Self::DeadlineExpired { .. } => "InsrptDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum StorungsmeldungState {
#[default]
New,
StorungsmeldungGesendet(StorungsmeldungData),
Bestaetigt(StorungsmeldungData),
Abgelehnt(StorungsmeldungData),
Ergebnisbericht(StorungsmeldungData),
DeadlineExpired {
label: String,
},
}
impl StorungsmeldungState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::StorungsmeldungGesendet(_) => "StorungsmeldungGesendet",
Self::Bestaetigt(_) => "Bestaetigt",
Self::Abgelehnt(_) => "Abgelehnt",
Self::Ergebnisbericht(_) => "Ergebnisbericht",
Self::DeadlineExpired { .. } => "DeadlineExpired",
}
}
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::Bestaetigt(_)
| Self::Abgelehnt(_)
| Self::Ergebnisbericht(_)
| Self::DeadlineExpired { .. }
)
}
}
#[derive(Clone)]
pub enum StorungsmeldungCommand {
SendStorungsmeldung {
pid: Pruefidentifikator,
msb_mp_id: MarktpartnerCode,
document_date: String,
message_ref: MessageRef,
},
ReceiveAntwort {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
message_ref: MessageRef,
},
ReceiveInformationsmeldung {
pid: Pruefidentifikator,
sender: MarktpartnerCode,
message_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for StorungsmeldungCommand {}
pub struct WimInsrptWorkflow;
impl Workflow for WimInsrptWorkflow {
type State = StorungsmeldungState;
type Event = StorungsmeldungEvent;
type Command = StorungsmeldungCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
if deadline.label() == ANTWORT_WINDOW_LABEL && !state.is_terminal() {
return Some(StorungsmeldungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
});
}
None
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
StorungsmeldungEvent::StorungsmeldungGesendet {
pruefidentifikator,
msb_mp_id,
document_date,
message_ref,
} => StorungsmeldungState::StorungsmeldungGesendet(StorungsmeldungData {
pruefidentifikator: *pruefidentifikator,
msb_mp_id: msb_mp_id.clone(),
document_date: document_date.clone(),
message_ref: message_ref.clone(),
}),
StorungsmeldungEvent::AntwortErhalten {
pruefidentifikator,
is_confirmation,
..
} => match state {
StorungsmeldungState::StorungsmeldungGesendet(data) => {
match pruefidentifikator.as_u32() {
23008 | 23009 => StorungsmeldungState::Ergebnisbericht(data),
_ if *is_confirmation => StorungsmeldungState::Bestaetigt(data),
_ => StorungsmeldungState::Abgelehnt(data),
}
}
other => other,
},
StorungsmeldungEvent::InformationsmeldungErhalten { .. } => {
state
}
StorungsmeldungEvent::DeadlineExpired { label, .. } => {
if state.is_terminal() {
state
} else {
StorungsmeldungState::DeadlineExpired {
label: label.to_string(),
}
}
}
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
StorungsmeldungCommand::SendStorungsmeldung {
pid,
msb_mp_id,
document_date,
message_ref,
} => {
if !matches!(state, StorungsmeldungState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !INSRPT_ANFRAGE_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"expected INSRPT PID 23001 for Störungsmeldung, got {pid}",
)));
}
let outbox = PendingOutbox::new(
"INSRPT",
msb_mp_id.as_str(),
serde_json::json!({
"type": "Stoerungsmeldung",
"pid": pid.as_u32(),
"message_ref": message_ref.as_str(),
}),
);
Ok(WorkflowOutput::with_outbox(
vec![StorungsmeldungEvent::StorungsmeldungGesendet {
pruefidentifikator: pid,
msb_mp_id,
document_date,
message_ref,
}],
vec![outbox],
))
}
StorungsmeldungCommand::ReceiveAntwort {
pid,
sender,
message_ref,
} => {
if !INSRPT_ANTWORT_PIDS.contains(&pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"PID {pid} is not a handled INSRPT response PID",
)));
}
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
let is_confirmation = matches!(pid.as_u32(), 23004 | 23008);
Ok(vec![StorungsmeldungEvent::AntwortErhalten {
pruefidentifikator: pid,
sender,
is_confirmation,
message_ref,
}]
.into())
}
StorungsmeldungCommand::ReceiveInformationsmeldung {
pid,
sender,
message_ref,
} => Ok(vec![StorungsmeldungEvent::InformationsmeldungErhalten {
pruefidentifikator: pid,
sender,
message_ref,
}]
.into()),
StorungsmeldungCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![StorungsmeldungEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}