use mako_engine::{
error::WorkflowError,
fristen::{HolidayCalendar, deadline_at_werktage},
ids::DeadlineId,
types::{MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, PendingDeadline, Workflow, WorkflowOutput},
};
use time::OffsetDateTime;
pub const WORKFLOW_NAME: &str = "wim-wertebestellung";
pub const ANFRAGE_PID: u32 = 35002;
pub const ANGEBOT_PID: u32 = 15003;
pub const BESTELLUNG_PID: u32 = 17007;
pub const STORNIERUNG_PID: u32 = 39002;
pub const BESTAETIGUNG_PID: u32 = 19011;
pub const ABLEHNUNG_PID: u32 = 19012;
pub const STORNO_BESTAETIGUNG_PID: u32 = 19013;
pub const STORNO_ABLEHNUNG_PID: u32 = 19014;
pub const INBOUND_PIDS: &[u32] = &[ANFRAGE_PID, BESTELLUNG_PID, STORNIERUNG_PID];
pub const ESA_INBOUND_PIDS: &[u32] = &[
ANGEBOT_PID,
BESTAETIGUNG_PID,
ABLEHNUNG_PID,
STORNO_BESTAETIGUNG_PID,
STORNO_ABLEHNUNG_PID,
];
pub const OUTBOUND_PIDS: &[u32] = &[
ANGEBOT_PID,
BESTAETIGUNG_PID,
ABLEHNUNG_PID,
STORNO_BESTAETIGUNG_PID,
STORNO_ABLEHNUNG_PID,
];
pub const ANGEBOT_FRIST_WT: u32 = 5;
pub const ANTWORT_FRIST_WT: u32 = 2;
pub const ANGEBOT_WINDOW_LABEL: &str = "wim-wertebestellung-angebot";
pub const BINDUNGSFRIST_LABEL: &str = "wim-wertebestellung-bindungsfrist";
pub const ANTWORT_WINDOW_LABEL: &str = "wim-wertebestellung-antwort";
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Zustellquittung {
pub received_at: OffsetDateTime,
pub positive: bool,
}
impl Zustellquittung {
#[must_use]
pub const fn positive(received_at: OffsetDateTime) -> Self {
Self {
received_at,
positive: true,
}
}
#[must_use]
pub const fn negative(received_at: OffsetDateTime) -> Self {
Self {
received_at,
positive: false,
}
}
pub fn frist(&self, werktage: u32) -> Result<OffsetDateTime, WorkflowError> {
if !self.positive {
return Err(WorkflowError::rejected(
"Frist cannot start from a negative AS4-Zustellquittung — GPKE Teil 1 \
admits only a positive Zustellquittung for Fristberechnung",
));
}
Ok(deadline_at_werktage(
self.received_at,
werktage,
HolidayCalendar::BdewMaKo,
))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Lokationsebene {
Marktlokation,
Messlokation,
Netzlokation,
}
impl Lokationsebene {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Marktlokation => "Marktlokation",
Self::Messlokation => "Messlokation",
Self::Netzlokation => "Netzlokation",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum WertebestellungEvent {
AnfrageEingegangen {
esa: MarktpartnerCode,
msb: MarktpartnerCode,
ebene: Lokationsebene,
lokations_id: String,
message_ref: MessageRef,
quittung: Zustellquittung,
},
AngebotAbgegeben {
message_ref: MessageRef,
bindungsfrist: OffsetDateTime,
},
AnfrageAbgelehnt {
reason: String,
},
BestellungEingegangen {
message_ref: MessageRef,
quittung: Zustellquittung,
},
BestellungBestaetigt {
message_ref: MessageRef,
},
BestellungAbgelehnt {
reason: String,
},
StornierungEingegangen {
message_ref: MessageRef,
quittung: Zustellquittung,
},
StornierungBestaetigt {
message_ref: MessageRef,
},
StornierungAbgelehnt {
reason: String,
},
AbbestellungEingegangen {
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
quittung: Zustellquittung,
},
AbbestellungBestaetigt {
message_ref: MessageRef,
},
BeendetDurchMsb {
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
reason: String,
},
LieferungBegonnen,
FristVersaeumt {
label: String,
},
}
impl EventPayload for WertebestellungEvent {
fn event_type(&self) -> &'static str {
match self {
Self::AnfrageEingegangen { .. } => "WertebestellungAnfrageEingegangen",
Self::AngebotAbgegeben { .. } => "WertebestellungAngebotAbgegeben",
Self::AnfrageAbgelehnt { .. } => "WertebestellungAnfrageAbgelehnt",
Self::BestellungEingegangen { .. } => "WertebestellungBestellungEingegangen",
Self::BestellungBestaetigt { .. } => "WertebestellungBestellungBestaetigt",
Self::BestellungAbgelehnt { .. } => "WertebestellungBestellungAbgelehnt",
Self::StornierungEingegangen { .. } => "WertebestellungStornierungEingegangen",
Self::StornierungBestaetigt { .. } => "WertebestellungStornierungBestaetigt",
Self::StornierungAbgelehnt { .. } => "WertebestellungStornierungAbgelehnt",
Self::AbbestellungEingegangen { .. } => "WertebestellungAbbestellungEingegangen",
Self::AbbestellungBestaetigt { .. } => "WertebestellungAbbestellungBestaetigt",
Self::LieferungBegonnen => "WertebestellungLieferungBegonnen",
Self::BeendetDurchMsb { .. } => "WertebestellungBeendetDurchMsb",
Self::FristVersaeumt { .. } => "WertebestellungFristVersaeumt",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct WertebestellungData {
pub esa: MarktpartnerCode,
pub msb: MarktpartnerCode,
pub ebene: Lokationsebene,
pub lokations_id: String,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum WertebestellungState {
#[default]
New,
AnfrageEingegangen(Box<WertebestellungData>),
AngebotAbgegeben {
data: Box<WertebestellungData>,
bindungsfrist: OffsetDateTime,
},
BestellungEingegangen(Box<WertebestellungData>),
BestellungBestaetigt {
data: Box<WertebestellungData>,
lieferung_begonnen: bool,
},
StornierungEingegangen(Box<WertebestellungData>),
AbbestellungEingegangen {
data: Box<WertebestellungData>,
beendigung_zum: OffsetDateTime,
},
Storniert(Box<WertebestellungData>),
Beendet {
data: Box<WertebestellungData>,
durch_msb: bool,
},
Abgelehnt {
reason: String,
},
}
impl WertebestellungState {
#[must_use]
pub const fn label(&self) -> &'static str {
match self {
Self::New => "New",
Self::AnfrageEingegangen(_) => "AnfrageEingegangen",
Self::AngebotAbgegeben { .. } => "AngebotAbgegeben",
Self::BestellungEingegangen(_) => "BestellungEingegangen",
Self::BestellungBestaetigt { .. } => "BestellungBestaetigt",
Self::StornierungEingegangen(_) => "StornierungEingegangen",
Self::AbbestellungEingegangen { .. } => "AbbestellungEingegangen",
Self::Storniert(_) => "Storniert",
Self::Beendet { .. } => "Beendet",
Self::Abgelehnt { .. } => "Abgelehnt",
}
}
#[must_use]
pub const fn lieferung_erlaubt(&self) -> bool {
matches!(
self,
Self::BestellungBestaetigt { .. } | Self::AbbestellungEingegangen { .. }
)
}
#[must_use]
pub const fn data(&self) -> Option<&WertebestellungData> {
match self {
Self::AnfrageEingegangen(d)
| Self::BestellungEingegangen(d)
| Self::StornierungEingegangen(d)
| Self::Storniert(d) => Some(d),
Self::AngebotAbgegeben { data, .. }
| Self::BestellungBestaetigt { data, .. }
| Self::AbbestellungEingegangen { data, .. }
| Self::Beendet { data, .. } => Some(data),
Self::New | Self::Abgelehnt { .. } => None,
}
}
}
#[derive(Clone)]
pub enum WertebestellungCommand {
ReceiveAnfrage {
pid: Pruefidentifikator,
esa: MarktpartnerCode,
msb: MarktpartnerCode,
ebene: Lokationsebene,
lokations_id: String,
message_ref: MessageRef,
quittung: Zustellquittung,
},
SendAngebot {
message_ref: MessageRef,
bindungsfrist: OffsetDateTime,
},
RejectAnfrage {
reason: String,
},
ReceiveBestellung {
pid: Pruefidentifikator,
message_ref: MessageRef,
quittung: Zustellquittung,
},
AnswerBestellung {
accept: bool,
message_ref: MessageRef,
reason: Option<String>,
},
ReceiveStornierung {
pid: Pruefidentifikator,
message_ref: MessageRef,
quittung: Zustellquittung,
},
AnswerStornierung {
accept: bool,
message_ref: MessageRef,
reason: Option<String>,
},
ReceiveAbbestellung {
pid: Pruefidentifikator,
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
quittung: Zustellquittung,
},
AnswerAbbestellung {
message_ref: MessageRef,
},
MarkLieferungBegonnen,
BeendenDurchMsb {
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
reason: String,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for WertebestellungCommand {}
pub struct WimWertebestellungWorkflow;
fn require_pid(pid: Pruefidentifikator, expected: u32, what: &str) -> Result<(), WorkflowError> {
if pid.as_u32() == expected {
Ok(())
} else {
Err(WorkflowError::rejected(format!(
"{what} expects PID {expected}, got {pid}"
)))
}
}
impl Workflow for WimWertebestellungWorkflow {
type State = WertebestellungState;
type Event = WertebestellungEvent;
type Command = WertebestellungCommand;
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
use WertebestellungEvent as E;
use WertebestellungState as S;
match event {
E::AnfrageEingegangen {
esa,
msb,
ebene,
lokations_id,
..
} => S::AnfrageEingegangen(Box::new(WertebestellungData {
esa: esa.clone(),
msb: msb.clone(),
ebene: *ebene,
lokations_id: lokations_id.clone(),
})),
E::AngebotAbgegeben { bindungsfrist, .. } => match state {
S::AnfrageEingegangen(data) => S::AngebotAbgegeben {
data,
bindungsfrist: *bindungsfrist,
},
other => other,
},
E::AnfrageAbgelehnt { reason } | E::BestellungAbgelehnt { reason } => S::Abgelehnt {
reason: reason.clone(),
},
E::BestellungEingegangen { .. } => match state {
S::AngebotAbgegeben { data, .. } => S::BestellungEingegangen(data),
other => other,
},
E::BestellungBestaetigt { .. } => match state {
S::BestellungEingegangen(data) => S::BestellungBestaetigt {
data,
lieferung_begonnen: false,
},
other => other,
},
E::StornierungEingegangen { .. } => match state {
S::BestellungBestaetigt { data, .. } => S::StornierungEingegangen(data),
other => other,
},
E::StornierungBestaetigt { .. } => match state {
S::StornierungEingegangen(data) => S::Storniert(data),
other => other,
},
E::StornierungAbgelehnt { .. } => match state {
S::StornierungEingegangen(data) => S::BestellungBestaetigt {
data,
lieferung_begonnen: false,
},
other => other,
},
E::AbbestellungEingegangen { beendigung_zum, .. } => match state {
S::BestellungBestaetigt { data, .. } => S::AbbestellungEingegangen {
data,
beendigung_zum: *beendigung_zum,
},
other => other,
},
E::AbbestellungBestaetigt { .. } => match state {
S::AbbestellungEingegangen { data, .. } => S::Beendet {
data,
durch_msb: false,
},
other => other,
},
E::LieferungBegonnen => match state {
S::BestellungBestaetigt { data, .. } => S::BestellungBestaetigt {
data,
lieferung_begonnen: true,
},
other => other,
},
E::BeendetDurchMsb { .. } => match state {
S::BestellungBestaetigt { data, .. } | S::AbbestellungEingegangen { data, .. } => {
S::Beendet {
data,
durch_msb: true,
}
}
other => other,
},
E::FristVersaeumt { .. } => state,
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
use WertebestellungCommand as C;
use WertebestellungEvent as E;
use WertebestellungState as S;
match command {
C::ReceiveAnfrage {
pid,
esa,
msb,
ebene,
lokations_id,
message_ref,
quittung,
} => {
if !matches!(state, S::New) {
return Err(WorkflowError::invalid_state("New", state.label()));
}
require_pid(pid, ANFRAGE_PID, "Anfrage von Werten")?;
if lokations_id.trim().is_empty() {
return Err(WorkflowError::rejected(format!(
"Anfrage auf Ebene {} ohne Lokations-ID",
ebene.as_str()
)));
}
let due = quittung.frist(ANGEBOT_FRIST_WT)?;
Ok(WorkflowOutput {
events: vec![E::AnfrageEingegangen {
esa,
msb,
ebene,
lokations_id,
message_ref,
quittung,
}],
outbox: Vec::new(),
deadlines: vec![PendingDeadline::new(ANGEBOT_WINDOW_LABEL, due)],
})
}
C::SendAngebot {
message_ref,
bindungsfrist,
} => {
if !matches!(state, S::AnfrageEingegangen(_)) {
return Err(WorkflowError::invalid_state(
"AnfrageEingegangen",
state.label(),
));
}
Ok(WorkflowOutput {
events: vec![E::AngebotAbgegeben {
message_ref,
bindungsfrist,
}],
outbox: Vec::new(),
deadlines: vec![PendingDeadline::new(BINDUNGSFRIST_LABEL, bindungsfrist)],
})
}
C::RejectAnfrage { reason } => {
if !matches!(state, S::AnfrageEingegangen(_)) {
return Err(WorkflowError::invalid_state(
"AnfrageEingegangen",
state.label(),
));
}
Ok(vec![E::AnfrageAbgelehnt { reason }].into())
}
C::ReceiveBestellung {
pid,
message_ref,
quittung,
} => {
let S::AngebotAbgegeben { bindungsfrist, .. } = state else {
return Err(WorkflowError::invalid_state(
"AngebotAbgegeben",
state.label(),
));
};
require_pid(pid, BESTELLUNG_PID, "Bestellung von Werten")?;
if quittung.received_at > *bindungsfrist {
return Err(WorkflowError::rejected(format!(
"Bestellung ging am {} ein, die Bindungsfrist des Angebots endete am {}",
quittung.received_at, bindungsfrist
)));
}
let due = quittung.frist(ANTWORT_FRIST_WT)?;
Ok(WorkflowOutput {
events: vec![E::BestellungEingegangen {
message_ref,
quittung,
}],
outbox: Vec::new(),
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerBestellung {
accept,
message_ref,
reason,
} => {
if !matches!(state, S::BestellungEingegangen(_)) {
return Err(WorkflowError::invalid_state(
"BestellungEingegangen",
state.label(),
));
}
if accept {
Ok(vec![E::BestellungBestaetigt { message_ref }].into())
} else {
Ok(vec![E::BestellungAbgelehnt {
reason: reason.ok_or_else(|| {
WorkflowError::rejected(
"Ablehnung der Bestellung erfordert eine Begründung \
(UC 4.1 Nr. 4: \"informiert der MSB den ESA über die Gründe\")",
)
})?,
}]
.into())
}
}
C::ReceiveStornierung {
pid,
message_ref,
quittung,
} => {
let S::BestellungBestaetigt {
lieferung_begonnen, ..
} = state
else {
return Err(WorkflowError::invalid_state(
"BestellungBestaetigt",
state.label(),
));
};
require_pid(pid, STORNIERUNG_PID, "Stornierung einer Bestellung")?;
if *lieferung_begonnen {
return Err(WorkflowError::rejected(
"Stornierung nicht mehr möglich — die Übermittlung von Werten hat \
bereits begonnen; die Beendigung erfolgt über die Abbestellung \
(WiM Teil 2, UC 4.3)",
));
}
let due = quittung.frist(ANTWORT_FRIST_WT)?;
Ok(WorkflowOutput {
events: vec![E::StornierungEingegangen {
message_ref,
quittung,
}],
outbox: Vec::new(),
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerStornierung {
accept,
message_ref,
reason,
} => {
if !matches!(state, S::StornierungEingegangen(_)) {
return Err(WorkflowError::invalid_state(
"StornierungEingegangen",
state.label(),
));
}
if accept {
Ok(vec![E::StornierungBestaetigt { message_ref }].into())
} else {
Ok(vec![E::StornierungAbgelehnt {
reason: reason.ok_or_else(|| {
WorkflowError::rejected(
"Ablehnung der Stornierung erfordert eine Begründung",
)
})?,
}]
.into())
}
}
C::ReceiveAbbestellung {
pid,
message_ref,
beendigung_zum,
quittung,
} => {
if !matches!(state, S::BestellungBestaetigt { .. }) {
return Err(WorkflowError::invalid_state(
"BestellungBestaetigt",
state.label(),
));
}
require_pid(pid, BESTELLUNG_PID, "Abbestellung von Werten")?;
let due = quittung.frist(ANTWORT_FRIST_WT)?;
Ok(WorkflowOutput {
events: vec![E::AbbestellungEingegangen {
message_ref,
beendigung_zum,
quittung,
}],
outbox: Vec::new(),
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerAbbestellung { message_ref } => {
if !matches!(state, S::AbbestellungEingegangen { .. }) {
return Err(WorkflowError::invalid_state(
"AbbestellungEingegangen",
state.label(),
));
}
Ok(vec![E::AbbestellungBestaetigt { message_ref }].into())
}
C::MarkLieferungBegonnen => match state {
S::BestellungBestaetigt {
lieferung_begonnen: true,
..
} => Ok(Vec::new().into()),
S::BestellungBestaetigt { .. } => Ok(vec![E::LieferungBegonnen].into()),
other => Err(WorkflowError::invalid_state(
"BestellungBestaetigt",
other.label(),
)),
},
C::BeendenDurchMsb {
message_ref,
beendigung_zum,
reason,
} => {
if !state.lieferung_erlaubt() {
return Err(WorkflowError::invalid_state(
"BestellungBestaetigt",
state.label(),
));
}
Ok(vec![E::BeendetDurchMsb {
message_ref,
beendigung_zum,
reason,
}]
.into())
}
C::TimeoutExpired { label, .. } => {
let outstanding = matches!(
(state, label.as_ref()),
(S::AnfrageEingegangen(_), ANGEBOT_WINDOW_LABEL)
| (S::BestellungEingegangen(_), ANTWORT_WINDOW_LABEL)
| (S::StornierungEingegangen(_), ANTWORT_WINDOW_LABEL)
| (S::AbbestellungEingegangen { .. }, ANTWORT_WINDOW_LABEL)
);
if outstanding {
return Ok(vec![E::FristVersaeumt {
label: label.to_string(),
}]
.into());
}
Ok(Vec::new().into())
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReqoteKind {
EsaWerteanfrage,
Preisanfrage,
}
#[must_use]
pub fn classify_reqote(sender_is_esa: bool, has_messprodukt: bool) -> ReqoteKind {
if sender_is_esa || has_messprodukt {
ReqoteKind::EsaWerteanfrage
} else {
ReqoteKind::Preisanfrage
}
}
#[must_use]
pub fn has_messprodukt<'a>(pia_codes: impl IntoIterator<Item = &'a str>) -> bool {
pia_codes.into_iter().any(|c| !c.trim().is_empty())
}