use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MaLo, MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "gpke-allokationsliste";
pub const ORDERS_ANFRAGE_PIDS: &[u32] = &[17110, 17114];
pub const ORDRSP_ABLEHNUNG_PIDS: &[u32] = &[19110, 19115];
pub const MSCONS_RESPONSE_PIDS: &[u32] = &[13014];
pub const ANTWORT_WINDOW_LABEL: &str = "gpke-allokationsliste-antwort";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum AllokationslisteEvent {
AnforderungGesendet {
orders_pid: Pruefidentifikator,
nb_mp_id: MarktpartnerCode,
malo: Option<MaLo>,
message_ref: MessageRef,
},
AnforderungAbgelehnt {
ordrsp_pid: Pruefidentifikator,
reason: Option<String>,
message_ref: MessageRef,
},
DatenGeliefert {
message_ref: MessageRef,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for AllokationslisteEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnforderungGesendet { .. } => "AllokationslisteAnforderungGesendet",
Self::AnforderungAbgelehnt { .. } => "AllokationslisteAnforderungAbgelehnt",
Self::DatenGeliefert { .. } => "AllokationslisteDatenGeliefert",
Self::DeadlineExpired { .. } => "AllokationslisteDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct AnforderungData {
pub orders_pid: Pruefidentifikator,
pub nb_mp_id: MarktpartnerCode,
pub malo: Option<MaLo>,
pub message_ref: MessageRef,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
#[derive(Default)]
pub enum AllokationslisteState {
#[default]
New,
AnforderungGesendet(AnforderungData),
Abgelehnt {
anforderung: AnforderungData,
ordrsp_pid: Pruefidentifikator,
reason: Option<String>,
},
DatenErhalten {
anforderung: AnforderungData,
},
DeadlineExpired {
anforderung: AnforderungData,
},
}
impl AllokationslisteState {
#[must_use]
pub fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::AnforderungGesendet(_) => "AnforderungGesendet",
Self::Abgelehnt { .. } => "Abgelehnt",
Self::DatenErhalten { .. } => "DatenErhalten",
Self::DeadlineExpired { .. } => "DeadlineExpired",
}
}
#[must_use]
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::Abgelehnt { .. } | Self::DatenErhalten { .. } | Self::DeadlineExpired { .. }
)
}
}
#[derive(Clone)]
pub enum AllokationslisteCommand {
SendAnforderung {
orders_pid: Pruefidentifikator,
nb_mp_id: MarktpartnerCode,
malo: Option<MaLo>,
message_ref: MessageRef,
payload: serde_json::Value,
},
ReceiveAblehnung {
ordrsp_pid: Pruefidentifikator,
reason: Option<String>,
message_ref: MessageRef,
},
NotifyDatenGeliefert {
message_ref: MessageRef,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for AllokationslisteCommand {}
pub struct GpkeAllokationslisteWorkflow;
impl Workflow for GpkeAllokationslisteWorkflow {
type State = AllokationslisteState;
type Event = AllokationslisteEvent;
type Command = AllokationslisteCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(ANTWORT_WINDOW_LABEL, AllokationslisteState::AnforderungGesendet(_)) => {
Some(AllokationslisteCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
AllokationslisteEvent::AnforderungGesendet {
orders_pid,
nb_mp_id,
malo,
message_ref,
} => AllokationslisteState::AnforderungGesendet(AnforderungData {
orders_pid: *orders_pid,
nb_mp_id: nb_mp_id.clone(),
malo: malo.clone(),
message_ref: message_ref.clone(),
}),
AllokationslisteEvent::AnforderungAbgelehnt {
ordrsp_pid, reason, ..
} => match state {
AllokationslisteState::AnforderungGesendet(anforderung) => {
AllokationslisteState::Abgelehnt {
anforderung,
ordrsp_pid: *ordrsp_pid,
reason: reason.clone(),
}
}
other => other,
},
AllokationslisteEvent::DatenGeliefert { .. } => match state {
AllokationslisteState::AnforderungGesendet(anforderung) => {
AllokationslisteState::DatenErhalten { anforderung }
}
other => other,
},
AllokationslisteEvent::DeadlineExpired { .. } => match state {
AllokationslisteState::AnforderungGesendet(anforderung) => {
AllokationslisteState::DeadlineExpired { anforderung }
}
other => other,
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
AllokationslisteCommand::SendAnforderung {
orders_pid,
nb_mp_id,
malo,
message_ref,
payload,
} => {
if !matches!(state, AllokationslisteState::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
if !ORDERS_ANFRAGE_PIDS.contains(&orders_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"not a valid Allokationsliste ORDERS PID: {orders_pid}",
)));
}
let event = AllokationslisteEvent::AnforderungGesendet {
orders_pid,
nb_mp_id: nb_mp_id.clone(),
malo: malo.clone(),
message_ref: message_ref.clone(),
};
let outbox = vec![PendingOutbox::new(
"ORDERS",
nb_mp_id.as_str(),
serde_json::json!({
"pid": orders_pid.as_u32(),
"malo": malo.as_ref().map(|m| m.as_str()),
"orders_ref": message_ref.as_str(),
"payload": payload,
}),
)];
Ok(WorkflowOutput::with_outbox(vec![event], outbox))
}
AllokationslisteCommand::ReceiveAblehnung {
ordrsp_pid,
reason,
message_ref,
} => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
if !matches!(state, AllokationslisteState::AnforderungGesendet(_)) {
return Err(WorkflowError::invalid_state(
"AnforderungGesendet",
state.label(),
));
}
if !ORDRSP_ABLEHNUNG_PIDS.contains(&ordrsp_pid.as_u32()) {
return Err(WorkflowError::rejected(format!(
"not a valid Allokationsliste rejection ORDRSP PID: {ordrsp_pid}",
)));
}
Ok(vec![AllokationslisteEvent::AnforderungAbgelehnt {
ordrsp_pid,
reason,
message_ref,
}]
.into())
}
AllokationslisteCommand::NotifyDatenGeliefert { message_ref } => {
if !matches!(state, AllokationslisteState::AnforderungGesendet(_)) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![AllokationslisteEvent::DatenGeliefert { message_ref }].into())
}
AllokationslisteCommand::TimeoutExpired { deadline_id, label } => {
if state.is_terminal() {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![AllokationslisteEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}