use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
outbox::PendingOutbox,
types::{MarktpartnerCode, MessageRef, Pruefidentifikator},
workflow::{CommandPayload, EventPayload, PendingDeadline, Workflow, WorkflowOutput},
};
use mako_fristen::{HolidayCalendar, deadline_at_werktage};
use time::OffsetDateTime;
pub const WORKFLOW_NAME: &str = "wim-wertebestellung";
pub const ANFRAGE_PID: Pruefidentifikator = Pruefidentifikator::const_new(35003);
pub const ANGEBOT_PID: Pruefidentifikator = Pruefidentifikator::const_new(15003);
pub const BESTELLUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(17007);
pub const ABBESTELLUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(17008);
pub const STORNIERUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(39002);
pub const BESTAETIGUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(19011);
pub const ABLEHNUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(19012);
pub const STORNO_BESTAETIGUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(19013);
pub const STORNO_ABLEHNUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(19014);
pub const WERTE_UEBERMITTLUNG_PID: Pruefidentifikator = Pruefidentifikator::const_new(13027);
pub const BEENDIGUNG_MSB_PID: Pruefidentifikator = Pruefidentifikator::const_new(21042);
pub const STS_BEENDET: &str = "105";
pub const INBOUND_PIDS: &[Pruefidentifikator] = &[
ANFRAGE_PID,
BESTELLUNG_PID,
ABBESTELLUNG_PID,
STORNIERUNG_PID,
];
pub const ESA_INBOUND_PIDS: &[Pruefidentifikator] = &[
ANGEBOT_PID,
BESTAETIGUNG_PID,
ABLEHNUNG_PID,
STORNO_BESTAETIGUNG_PID,
STORNO_ABLEHNUNG_PID,
BEENDIGUNG_MSB_PID,
];
pub const OUTBOUND_PIDS: &[Pruefidentifikator] = &[
ANGEBOT_PID,
BESTAETIGUNG_PID,
ABLEHNUNG_PID,
STORNO_BESTAETIGUNG_PID,
STORNO_ABLEHNUNG_PID,
BEENDIGUNG_MSB_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,
))
}
}
pub use crate::esa::{Abonnement, Bestellgegenstand, Lokationsebene};
#[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,
gegenstand: Box<Bestellgegenstand>,
message_ref: MessageRef,
quittung: Zustellquittung,
},
AngebotAbgegeben {
message_ref: MessageRef,
bindungsfrist: OffsetDateTime,
},
AnfrageAbgelehnt {
reason: String,
},
BestellungEingegangen {
message_ref: MessageRef,
#[serde(default = "default_start_abo")]
abonnement: Abonnement,
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,
},
AbbestellungAbgelehnt {
reason: String,
},
BeendetDurchMsb {
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
reason: String,
},
WerteUebermittelt {
message_ref: MessageRef,
interval_count: u32,
#[serde(default)]
bis: Option<time::Date>,
},
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::AbbestellungAbgelehnt { .. } => "WertebestellungAbbestellungAbgelehnt",
Self::WerteUebermittelt { .. } => "WertebestellungWerteUebermittelt",
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,
pub gegenstand: Box<Bestellgegenstand>,
#[serde(default)]
pub anfrage_ref: Option<String>,
#[serde(default)]
pub angebot_ref: Option<String>,
#[serde(default)]
pub bestellung_ref: Option<String>,
#[serde(default)]
pub inbound_order_ref: Option<String>,
#[serde(default)]
pub stornierung_ref: Option<String>,
#[serde(default)]
pub offene_antwort_abo: Option<Abonnement>,
#[serde(default)]
pub bindungsfrist: Option<OffsetDateTime>,
#[serde(default)]
pub abo_beginn: Option<time::Date>,
#[serde(default)]
pub juengste_lieferung: Option<time::Date>,
#[serde(default)]
pub bereits_beendet_zum: Option<time::Date>,
#[serde(default)]
pub lieferung_begonnen: bool,
}
#[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(Debug, Clone, PartialEq)]
pub enum WertebestellungCommand {
ReceiveAnfrage {
pid: Pruefidentifikator,
esa: MarktpartnerCode,
msb: MarktpartnerCode,
ebene: Lokationsebene,
lokations_id: String,
gegenstand: Box<Bestellgegenstand>,
message_ref: MessageRef,
quittung: Zustellquittung,
consent_block: Option<String>,
},
SendAngebot {
message_ref: MessageRef,
bindungsfrist: OffsetDateTime,
fruehester_start: Option<OffsetDateTime>,
angebot: Box<crate::esa::Angebot>,
},
RejectAnfrage {
reason: String,
},
ReceiveBestellung {
pid: Pruefidentifikator,
message_ref: MessageRef,
abonnement: Abonnement,
quittung: Zustellquittung,
consent_block: Option<String>,
},
AnswerBestellung {
antwort_code: String,
message_ref: MessageRef,
reason: Option<String>,
},
ReceiveStornierung {
pid: Pruefidentifikator,
message_ref: MessageRef,
quittung: Zustellquittung,
},
AnswerStornierung {
antwort_code: String,
message_ref: MessageRef,
reason: Option<String>,
},
ReceiveAbbestellung {
pid: Pruefidentifikator,
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
quittung: Zustellquittung,
},
AnswerAbbestellung {
antwort_code: String,
message_ref: MessageRef,
reason: Option<String>,
},
LiefereWerte {
message_ref: MessageRef,
reads: serde_json::Value,
},
MarkLieferungBegonnen,
BeendenDurchMsb {
message_ref: MessageRef,
beendigung_zum: OffsetDateTime,
reason: String,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for WertebestellungCommand {}
pub struct WimWertebestellungWorkflow;
const fn default_start_abo() -> Abonnement {
Abonnement::StartAbo
}
fn ccyymmdd(dt: OffsetDateTime) -> String {
format!("{:04}{:02}{:02}", dt.year(), u8::from(dt.month()), dt.day())
}
fn require_pid(
pid: Pruefidentifikator,
expected: Pruefidentifikator,
what: &str,
) -> Result<(), WorkflowError> {
if pid == 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 on_deadline(
deadline: &mako_engine::deadline::Deadline,
_state: &Self::State,
) -> Option<Self::Command> {
let owned = matches!(
deadline.label(),
ANGEBOT_WINDOW_LABEL | ANTWORT_WINDOW_LABEL | BINDUNGSFRIST_LABEL
);
owned.then(|| WertebestellungCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
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,
gegenstand,
message_ref,
..
} => S::AnfrageEingegangen(Box::new(WertebestellungData {
esa: esa.clone(),
msb: msb.clone(),
ebene: *ebene,
lokations_id: lokations_id.clone(),
gegenstand: gegenstand.clone(),
anfrage_ref: Some(message_ref.as_str().to_owned()),
angebot_ref: None,
bestellung_ref: None,
inbound_order_ref: None,
stornierung_ref: None,
offene_antwort_abo: None,
bindungsfrist: None,
abo_beginn: None,
juengste_lieferung: None,
bereits_beendet_zum: None,
lieferung_begonnen: false,
})),
E::AngebotAbgegeben {
bindungsfrist,
message_ref,
} => match state {
S::AnfrageEingegangen(mut data) => {
data.angebot_ref = Some(message_ref.as_str().to_owned());
data.bindungsfrist = Some(*bindungsfrist);
S::AngebotAbgegeben {
data,
bindungsfrist: *bindungsfrist,
}
}
other => other,
},
E::AnfrageAbgelehnt { reason } | E::BestellungAbgelehnt { reason } => S::Abgelehnt {
reason: reason.clone(),
},
E::BestellungEingegangen {
message_ref,
abonnement,
..
} => match state {
S::AngebotAbgegeben { mut data, .. } => {
let r = message_ref.as_str().to_owned();
data.bestellung_ref = Some(r.clone());
data.inbound_order_ref = Some(r);
data.gegenstand.abonnement = *abonnement;
data.offene_antwort_abo = Some(*abonnement);
S::BestellungEingegangen(data)
}
other => other,
},
E::BestellungBestaetigt { .. } => match state {
S::BestellungEingegangen(mut data) => {
data.abo_beginn = Some(data.gegenstand.wunschtermin);
S::BestellungBestaetigt {
data,
lieferung_begonnen: false,
}
}
other => other,
},
E::StornierungEingegangen { message_ref, .. } => match state {
S::BestellungBestaetigt { mut data, .. } => {
data.stornierung_ref = Some(message_ref.as_str().to_owned());
S::StornierungEingegangen(data)
}
other => other,
},
E::StornierungBestaetigt { .. } => match state {
S::StornierungEingegangen(data) => S::Storniert(data),
other => other,
},
E::StornierungAbgelehnt { .. } => match state {
S::StornierungEingegangen(data) => {
let lieferung_begonnen = data.lieferung_begonnen;
S::BestellungBestaetigt {
data,
lieferung_begonnen,
}
}
other => other,
},
E::AbbestellungEingegangen {
beendigung_zum,
message_ref,
..
} => match state {
S::BestellungBestaetigt { mut data, .. } => {
data.inbound_order_ref = Some(message_ref.as_str().to_owned());
data.offene_antwort_abo = Some(Abonnement::EndeAbo);
S::AbbestellungEingegangen {
data,
beendigung_zum: *beendigung_zum,
}
}
other => other,
},
E::AbbestellungBestaetigt { .. } => match state {
S::AbbestellungEingegangen {
mut data,
beendigung_zum,
} => {
data.bereits_beendet_zum = Some(beendigung_zum.date());
S::Beendet {
data,
durch_msb: false,
}
}
other => other,
},
E::AbbestellungAbgelehnt { .. } => match state {
S::AbbestellungEingegangen { mut data, .. } => {
data.offene_antwort_abo = None;
let lieferung_begonnen = data.lieferung_begonnen;
S::BestellungBestaetigt {
data,
lieferung_begonnen,
}
}
other => other,
},
E::LieferungBegonnen => match state {
S::BestellungBestaetigt { mut data, .. } => {
data.lieferung_begonnen = true;
S::BestellungBestaetigt {
data,
lieferung_begonnen: true,
}
}
S::StornierungEingegangen(mut data) => {
data.lieferung_begonnen = true;
S::StornierungEingegangen(data)
}
S::AbbestellungEingegangen {
mut data,
beendigung_zum,
} => {
data.lieferung_begonnen = true;
S::AbbestellungEingegangen {
data,
beendigung_zum,
}
}
other => other,
},
E::WerteUebermittelt { bis, .. } => match state {
S::BestellungBestaetigt { mut data, .. } => {
data.lieferung_begonnen = true;
data.juengste_lieferung = data.juengste_lieferung.max(*bis);
S::BestellungBestaetigt {
data,
lieferung_begonnen: true,
}
}
S::StornierungEingegangen(mut data) => {
data.lieferung_begonnen = true;
data.juengste_lieferung = data.juengste_lieferung.max(*bis);
S::StornierungEingegangen(data)
}
S::AbbestellungEingegangen {
mut data,
beendigung_zum,
} => {
data.lieferung_begonnen = true;
data.juengste_lieferung = data.juengste_lieferung.max(*bis);
S::AbbestellungEingegangen {
data,
beendigung_zum,
}
}
other => other,
},
E::BeendetDurchMsb { beendigung_zum, .. } => match state {
S::BestellungBestaetigt { mut data, .. }
| S::AbbestellungEingegangen { mut data, .. } => {
data.bereits_beendet_zum = Some(beendigung_zum.date());
S::Beendet {
data,
durch_msb: true,
}
}
other => other,
},
E::FristVersaeumt { .. } => state,
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
fn esa_process_initiated(
pid: Pruefidentifikator,
data: &WertebestellungData,
beendigung_zum: Option<OffsetDateTime>,
) -> PendingOutbox {
PendingOutbox::new(
"ProcessInitiated",
data.msb.as_str(),
serde_json::json!({
"pid": pid.as_u32(),
"malo_id": data.lokations_id,
"lokations_id": data.lokations_id,
"esa_mp_id": data.esa.as_str(),
"msb_mp_id": data.msb.as_str(),
"ebene": data.ebene,
"messprodukt": data.gegenstand.messprodukt,
"abonnement": data.gegenstand.abonnement.imd_code(),
"ausfuehrungsdatum": beendigung_zum
.map(|d| d.date())
.unwrap_or(data.gegenstand.wunschtermin)
.to_string(),
"bindungsfrist": data.bindungsfrist.and_then(|b| {
b.format(&time::format_description::well_known::Rfc3339).ok()
}),
"lieferung_begonnen": data.lieferung_begonnen,
"abo_beginn": data.abo_beginn.map(|d| d.to_string()),
"bereits_beendet_zum": data.bereits_beendet_zum.map(|d| d.to_string()),
"juengste_lieferung": data.juengste_lieferung.map(|d| d.to_string()),
}),
)
}
fn esa_answer(
message_type: &'static str,
pid: Pruefidentifikator,
data: &WertebestellungData,
message_ref: &MessageRef,
antwort: Option<(&'static str, &mako_pruefung::codes::AntwortCode)>,
reason: Option<&str>,
) -> PendingOutbox {
let ist_storno_antwort = pid == STORNO_BESTAETIGUNG_PID || pid == STORNO_ABLEHNUNG_PID;
let korrelation_ref = if ist_storno_antwort {
data.stornierung_ref.clone()
} else {
data.inbound_order_ref.clone()
};
let abo = data
.offene_antwort_abo
.unwrap_or(data.gegenstand.abonnement);
PendingOutbox::new(
message_type,
data.esa.as_str(),
serde_json::json!({
"pid": pid,
"sender": data.msb.as_str(),
"receiver": data.esa.as_str(),
"message_ref": message_ref.as_str(),
"korrelation_ref": korrelation_ref,
"abonnement": abo.imd_code(),
"antwort_code": antwort.map(|(_, c)| c.code),
"antwort_codeliste": antwort.map(|(t, _)| t),
"reason": reason,
"messprodukt": data.gegenstand.messprodukt,
}),
)
}
fn resolve_antwort(
tree: &'static str,
antwort_code: &str,
bestaetigung_pid: Pruefidentifikator,
ablehnung_pid: Pruefidentifikator,
) -> Result<
(
Pruefidentifikator,
&'static mako_pruefung::codes::AntwortCode,
bool,
),
WorkflowError,
> {
let code = mako_pruefung::codes::lookup(tree, antwort_code).ok_or_else(|| {
WorkflowError::rejected(format!(
"Antwortcode {antwort_code:?} ist in {tree} nicht veröffentlicht"
))
})?;
let zustimmung = code.ist_zustimmung().ok_or_else(|| {
WorkflowError::rejected(format!(
"{} liegt nicht auf der Zustimmungs-/Ablehnungsachse von {tree}",
code.code
))
})?;
Ok((
if zustimmung {
bestaetigung_pid
} else {
ablehnung_pid
},
code,
zustimmung,
))
}
use WertebestellungCommand as C;
use WertebestellungEvent as E;
use WertebestellungState as S;
match command {
C::ReceiveAnfrage {
pid,
esa,
msb,
ebene,
lokations_id,
gegenstand,
message_ref,
quittung,
consent_block,
} => {
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()
)));
}
if let Some(reason) = consent_block {
let outbox = PendingOutbox::new(
"QUOTES",
esa.as_str(),
serde_json::json!({
"pid": ANGEBOT_PID,
"sender": msb.as_str(),
"receiver": esa.as_str(),
"message_ref": message_ref.as_str(),
"location": lokations_id,
"korrelation_ref": message_ref.as_str(),
"reason": reason.clone(),
}),
);
return Ok(WorkflowOutput::with_outbox(
vec![E::AnfrageAbgelehnt { reason }],
vec![outbox],
));
}
if let Err(e) = gegenstand.validate(ebene) {
let outbox = PendingOutbox::new(
"QUOTES",
esa.as_str(),
serde_json::json!({
"pid": ANGEBOT_PID,
"sender": msb.as_str(),
"receiver": esa.as_str(),
"message_ref": message_ref.as_str(),
"location": lokations_id,
"korrelation_ref": message_ref.as_str(),
"reason": e.to_string(),
}),
);
return Ok(WorkflowOutput::with_outbox(
vec![E::AnfrageAbgelehnt {
reason: e.to_string(),
}],
vec![outbox],
));
}
let due = quittung.frist(ANGEBOT_FRIST_WT)?;
let initiated = esa_process_initiated(
pid,
&WertebestellungData {
esa: esa.clone(),
msb: msb.clone(),
ebene,
lokations_id: lokations_id.clone(),
gegenstand: gegenstand.clone(),
anfrage_ref: Some(message_ref.as_str().to_owned()),
angebot_ref: None,
bestellung_ref: None,
inbound_order_ref: None,
stornierung_ref: None,
offene_antwort_abo: None,
bindungsfrist: None,
abo_beginn: None,
juengste_lieferung: None,
bereits_beendet_zum: None,
lieferung_begonnen: false,
},
None,
);
Ok(WorkflowOutput {
events: vec![E::AnfrageEingegangen {
esa,
msb,
ebene,
lokations_id,
gegenstand,
message_ref,
quittung,
}],
outbox: vec![initiated],
deadlines: vec![PendingDeadline::new(ANGEBOT_WINDOW_LABEL, due)],
})
}
C::SendAngebot {
message_ref,
bindungsfrist,
fruehester_start,
angebot,
} => {
let Some(data) = state
.data()
.filter(|_| matches!(state, S::AnfrageEingegangen(_)))
else {
return Err(WorkflowError::invalid_state(
"AnfrageEingegangen",
state.label(),
));
};
let bindungsfrist_tage = (bindungsfrist - OffsetDateTime::now_utc())
.whole_days()
.max(1);
let start = fruehester_start
.unwrap_or_else(|| data.gegenstand.wunschtermin.midnight().assume_utc());
if angebot.ist_leer() {
return Err(WorkflowError::rejected(
"Angebot ohne Preisangabe — SG31 PRI und die OBIS-Kennzahlen sind Muss \
auf der QUOTES 15003 (QUOTES AHB 1.1a §4.3); eine Anfrage ohne Angebot \
wird mit RejectAnfrage abgelehnt",
));
}
let artikel_ids: Vec<&str> = {
let mut ids: Vec<&str> = Vec::new();
for p in &angebot.preise {
if !ids.contains(&p.artikel_id.as_str()) {
ids.push(p.artikel_id.as_str());
}
}
ids
};
let preise: Vec<serde_json::Value> = angebot
.preise
.iter()
.map(|p| {
serde_json::json!({
"betrag": p.betrag,
"art": p.preistyp.pri_code(),
"einheit": p.einheit,
})
})
.collect();
let outbox = PendingOutbox::new(
"QUOTES",
data.esa.as_str(),
serde_json::json!({
"pid": ANGEBOT_PID,
"sender": data.msb.as_str(),
"receiver": data.esa.as_str(),
"location": data.lokations_id,
"message_ref": message_ref.as_str(),
"korrelation_ref": data.anfrage_ref,
"bindungsfrist_tage": bindungsfrist_tage,
"fruehester_start": ccyymmdd(start),
"messprodukt": data.gegenstand.messprodukt,
"currency": angebot.waehrung.as_deref().unwrap_or("EUR"),
"artikel_ids": artikel_ids,
"obis": angebot.obis_kennzahlen,
"preise": preise,
}),
);
Ok(WorkflowOutput {
events: vec![E::AngebotAbgegeben {
message_ref,
bindungsfrist,
}],
outbox: vec![outbox],
deadlines: vec![PendingDeadline::new(BINDUNGSFRIST_LABEL, bindungsfrist)],
})
}
C::RejectAnfrage { reason } => {
let Some(data) = state
.data()
.filter(|_| matches!(state, S::AnfrageEingegangen(_)))
else {
return Err(WorkflowError::invalid_state(
"AnfrageEingegangen",
state.label(),
));
};
let outbox = PendingOutbox::new(
"QUOTES",
data.esa.as_str(),
serde_json::json!({
"pid": ANGEBOT_PID,
"sender": data.msb.as_str(),
"receiver": data.esa.as_str(),
"location": data.lokations_id,
"korrelation_ref": data.anfrage_ref,
"reason": reason.clone(),
}),
);
Ok(WorkflowOutput::with_outbox(
vec![E::AnfrageAbgelehnt { reason }],
vec![outbox],
))
}
C::ReceiveBestellung {
pid,
message_ref,
abonnement,
quittung,
consent_block,
} => {
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
)));
}
if let Some(reason) = consent_block {
let data = state.data().ok_or_else(|| {
WorkflowError::invalid_state("AngebotAbgegeben", state.label())
})?;
let tree = crate::esa::EBD_ESA_BESTELLUNG;
let code = mako_pruefung::codes::lookup(tree, "A08")
.expect("A08 is published in E_0256");
let outbox = esa_answer(
"ORDRSP",
ABLEHNUNG_PID,
data,
&message_ref,
Some((tree, code)),
Some(reason.as_str()),
);
return Ok(WorkflowOutput::with_outbox(
vec![E::BestellungAbgelehnt { reason }],
vec![outbox],
));
}
let due = quittung.frist(ANTWORT_FRIST_WT)?;
let mut fuer_meldung = state
.data()
.cloned()
.ok_or_else(|| WorkflowError::invalid_state("AngebotAbgegeben", "New"))?;
fuer_meldung.gegenstand.abonnement = abonnement;
let initiated = esa_process_initiated(pid, &fuer_meldung, None);
Ok(WorkflowOutput {
events: vec![E::BestellungEingegangen {
message_ref,
abonnement,
quittung,
}],
outbox: vec![initiated],
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerBestellung {
antwort_code,
message_ref,
reason,
} => {
let Some(data) = state
.data()
.filter(|_| matches!(state, S::BestellungEingegangen(_)))
else {
return Err(WorkflowError::invalid_state(
"BestellungEingegangen",
state.label(),
));
};
let tree = data
.offene_antwort_abo
.unwrap_or(data.gegenstand.abonnement)
.antwort_ebd();
let (pid, code, zustimmung) =
resolve_antwort(tree, &antwort_code, BESTAETIGUNG_PID, ABLEHNUNG_PID)?;
if !zustimmung && reason.is_none() && code.braucht_bemerkung {
return Err(WorkflowError::rejected(format!(
"{tree} {} ({}) verlangt eine schriftliche Erläuterung",
code.code, code.bedeutung
)));
}
let outbox = esa_answer(
"ORDRSP",
pid,
data,
&message_ref,
Some((tree, code)),
reason.as_deref(),
);
let event = if zustimmung {
E::BestellungBestaetigt { message_ref }
} else {
E::BestellungAbgelehnt {
reason: reason.unwrap_or_else(|| code.bedeutung.to_owned()),
}
};
Ok(WorkflowOutput::with_outbox(vec![event], vec![outbox]))
}
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)?;
let initiated = esa_process_initiated(
pid,
state.data().ok_or_else(|| {
WorkflowError::invalid_state("BestellungBestaetigt", "New")
})?,
None,
);
Ok(WorkflowOutput {
events: vec![E::StornierungEingegangen {
message_ref,
quittung,
}],
outbox: vec![initiated],
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerStornierung {
antwort_code,
message_ref,
reason,
} => {
let Some(data) = state
.data()
.filter(|_| matches!(state, S::StornierungEingegangen(_)))
else {
return Err(WorkflowError::invalid_state(
"StornierungEingegangen",
state.label(),
));
};
let tree = crate::esa::EBD_ESA_STORNIERUNG;
let (pid, code, zustimmung) = resolve_antwort(
tree,
&antwort_code,
STORNO_BESTAETIGUNG_PID,
STORNO_ABLEHNUNG_PID,
)?;
let outbox = esa_answer(
"ORDRSP",
pid,
data,
&message_ref,
Some((tree, code)),
reason.as_deref(),
);
let event = if zustimmung {
E::StornierungBestaetigt { message_ref }
} else {
E::StornierungAbgelehnt {
reason: reason.unwrap_or_else(|| code.bedeutung.to_owned()),
}
};
Ok(WorkflowOutput::with_outbox(vec![event], vec![outbox]))
}
C::ReceiveAbbestellung {
pid,
message_ref,
beendigung_zum,
quittung,
} => {
if !matches!(state, S::BestellungBestaetigt { .. }) {
return Err(WorkflowError::invalid_state(
"BestellungBestaetigt",
state.label(),
));
}
require_pid(pid, ABBESTELLUNG_PID, "Abbestellung von Werten")?;
let due = quittung.frist(ANTWORT_FRIST_WT)?;
let mut fuer_meldung = state
.data()
.cloned()
.ok_or_else(|| WorkflowError::invalid_state("BestellungBestaetigt", "New"))?;
fuer_meldung.gegenstand.abonnement = Abonnement::EndeAbo;
let initiated = esa_process_initiated(pid, &fuer_meldung, Some(beendigung_zum));
Ok(WorkflowOutput {
events: vec![E::AbbestellungEingegangen {
message_ref,
beendigung_zum,
quittung,
}],
outbox: vec![initiated],
deadlines: vec![PendingDeadline::new(ANTWORT_WINDOW_LABEL, due)],
})
}
C::AnswerAbbestellung {
antwort_code,
message_ref,
reason,
} => {
let Some(data) = state
.data()
.filter(|_| matches!(state, S::AbbestellungEingegangen { .. }))
else {
return Err(WorkflowError::invalid_state(
"AbbestellungEingegangen",
state.label(),
));
};
let tree = crate::esa::EBD_ESA_BEENDIGUNG;
let (pid, code, zustimmung) =
resolve_antwort(tree, &antwort_code, BESTAETIGUNG_PID, ABLEHNUNG_PID)?;
let outbox = esa_answer(
"ORDRSP",
pid,
data,
&message_ref,
Some((tree, code)),
reason.as_deref(),
);
let event = if zustimmung {
E::AbbestellungBestaetigt { message_ref }
} else {
E::AbbestellungAbgelehnt {
reason: reason.unwrap_or_else(|| code.bedeutung.to_owned()),
}
};
Ok(WorkflowOutput::with_outbox(vec![event], vec![outbox]))
}
C::LiefereWerte { message_ref, reads } => {
let Some(data) = state.data().filter(|_| state.lieferung_erlaubt()) else {
return Err(WorkflowError::invalid_state(
"BestellungBestaetigt|AbbestellungEingegangen",
state.label(),
));
};
let intervals = reads.as_array().filter(|a| !a.is_empty()).ok_or_else(|| {
WorkflowError::rejected(
"Werteübermittlung ohne Intervallwerte — reads muss ein nicht-leeres \
Array sein",
)
})?;
let interval_count = u32::try_from(intervals.len()).unwrap_or(u32::MAX);
let bis = intervals
.iter()
.filter_map(|iv| {
let raw = iv.get("dtm_to").or_else(|| iv.get("bis"))?.as_str()?;
let digits: String =
raw.chars().filter(char::is_ascii_digit).take(8).collect();
(digits.len() == 8).then_some(())?;
time::Date::from_calendar_date(
digits[0..4].parse().ok()?,
time::Month::try_from(digits[4..6].parse::<u8>().ok()?).ok()?,
digits[6..8].parse().ok()?,
)
.ok()
})
.max();
let outbox = PendingOutbox::new(
"MSCONS",
data.esa.as_str(),
serde_json::json!({
"pid": WERTE_UEBERMITTLUNG_PID,
"sender_mp_id": data.msb.as_str(),
"receiver_mp_id": data.esa.as_str(),
"malo_id": data.lokations_id,
"message_ref": message_ref.as_str(),
"korrelation_ref": data.bestellung_ref,
"messprodukt": data.gegenstand.messprodukt,
"reads": reads,
}),
);
Ok(WorkflowOutput::with_outbox(
vec![E::WerteUebermittelt {
message_ref,
interval_count,
bis,
}],
vec![outbox],
))
}
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(),
));
}
let data = state
.data()
.expect("lieferung_erlaubt() implies process data is present");
let outbox = PendingOutbox::new(
"IFTSTA",
data.esa.as_str(),
serde_json::json!({
"pid": BEENDIGUNG_MSB_PID,
"sender": data.msb.as_str(),
"receiver": data.esa.as_str(),
"message_ref": message_ref.as_str(),
"sts_code": STS_BEENDET,
"korrelation_ref": data.bestellung_ref,
"beendigung_zum": beendigung_zum,
"reason": reason,
}),
);
Ok(WorkflowOutput::with_outbox(
vec![E::BeendetDurchMsb {
message_ref,
beendigung_zum,
reason,
}],
vec![outbox],
))
}
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())
}
}
}
}