Skip to main content

mako_redispatch/
stammdaten.rs

1//! Stammdatenübermittlung workflow for Redispatch 2.0.
2//!
3//! **Direction:** ANB → VNB → ÜNB\
4//! **Document:** `redispatch_xml::Stammdaten` (Z02 reduced, Z03 enriched,
5//! Z04 NB aggregate, Z14 BKV)
6//!
7//! # Process description
8//!
9//! 1. ANB sends `Stammdaten` to VNB (initial + updates on change).
10//! 2. Receiver sends `AcknowledgementDocument` within **6 wall-clock hours**
11//!    (UTC — see note below).
12//! 3. VNB optionally forwards enriched `Stammdaten` to ÜNB within **1 Werktag**
13//!    of the master-data change (BK6-20-060 §3.2).
14//!
15//! # Clock semantics
16//!
17//! All Redispatch 2.0 fristen use **UTC wall-clock hours**, not German local
18//! time (CET/CEST). The `UtcDateTime` fields in XSD carry explicit `Z` offsets.
19//! This differs from GPKE/WiM deadlines, which use German local time.
20//!
21//! # Regulatory basis
22//!
23//! `BNetzA` BK6-20-059 §4.3 (6h ACK), BK6-20-060 §3.2 (Stammdaten update).
24
25use mako_engine::{
26    deadline::Deadline,
27    error::WorkflowError,
28    ids::DeadlineId,
29    workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
30};
31use serde::{Deserialize, Serialize};
32
33// ── Workflow name ─────────────────────────────────────────────────────────────
34
35/// Stable workflow name — used in `ProcessRegistry` lookups and log output.
36pub const WORKFLOW_NAME: &str = "redispatch-stammdaten";
37
38// ── Deadline labels ───────────────────────────────────────────────────────────
39
40/// 6h UTC window for dispatching `AcknowledgementDocument` (BK6-20-059 §4.3).
41///
42/// Register immediately after [`StammdatenEvent::Received`] is applied.
43pub const ACK_WINDOW_LABEL: &str = "redispatch-stammdaten-ack-window";
44
45/// 1 Werktag forwarding window for VNB→ÜNB enrichment (BK6-20-060 §3.2).
46///
47/// Register after [`StammdatenEvent::Acknowledged`] is applied, when the
48/// deployment role is VNB.
49pub const FORWARD_WINDOW_LABEL: &str = "redispatch-stammdaten-forward-window";
50
51// ── Events ────────────────────────────────────────────────────────────────────
52
53/// Events emitted by the Stammdaten workflow.
54#[derive(Debug, Clone, Serialize, Deserialize)]
55#[serde(tag = "type", content = "data")]
56pub enum StammdatenEvent {
57    /// `Stammdaten` document received from ANB or VNB.
58    Received {
59        /// MRID (UUID) of the received `Stammdaten` document.
60        mrid: String,
61        /// GLN of the sender (ANB or VNB).
62        sender: String,
63        /// GLN of the receiver (VNB or ÜNB).
64        receiver: String,
65        /// Document type code (Z02/Z03/Z04/Z14).
66        doc_type: String,
67        /// Number of resource objects (`Anlagen`) included.
68        anlagen_count: u32,
69        /// UTC receipt timestamp in ISO-8601 format.
70        received_at: String,
71    },
72    /// `AcknowledgementDocument` dispatched within the 6h window.
73    Acknowledged {
74        /// MRID of the outbound `AcknowledgementDocument`.
75        ack_mrid: String,
76    },
77    /// Enriched `Stammdaten` forwarded upstream (VNB→ÜNB, role-conditional).
78    Forwarded {
79        /// MRID of the upstream `Stammdaten` sent to ÜNB.
80        upstream_mrid: String,
81    },
82    /// The 6h acknowledgement window expired without a response.
83    DeadlineExpired {
84        /// Unique ID of the expired deadline.
85        deadline_id: DeadlineId,
86        /// Label identifying the deadline type.
87        label: Box<str>,
88    },
89}
90
91impl EventPayload for StammdatenEvent {
92    fn event_type(&self) -> &'static str {
93        match self {
94            Self::Received { .. } => "StammdatenReceived",
95            Self::Acknowledged { .. } => "StammdatenAcknowledged",
96            Self::Forwarded { .. } => "StammdatenForwarded",
97            Self::DeadlineExpired { .. } => "StammdatenDeadlineExpired",
98        }
99    }
100}
101
102// ── Domain data ───────────────────────────────────────────────────────────────
103
104/// Business data captured when the `Stammdaten` document is first received.
105#[derive(Debug, Clone, Serialize, Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct ReceivedData {
108    /// MRID (UUID) of the received `Stammdaten` document.
109    pub mrid: String,
110    /// GLN of the sender.
111    pub sender: String,
112    /// GLN of the receiver.
113    pub receiver: String,
114    /// Document type code.
115    pub doc_type: String,
116    /// Number of resource objects.
117    pub anlagen_count: u32,
118    /// UTC receipt timestamp.
119    pub received_at: String,
120}
121
122// ── State ─────────────────────────────────────────────────────────────────────
123
124/// Current state of a Stammdaten process stream.
125///
126/// # Lifecycle
127///
128/// ```text
129/// New → Received → Acknowledged → [Forwarded →] Done
130///                ↘ DeadlineExpired (6h window lapsed)
131/// ```
132#[derive(Debug, Clone, Default, Serialize, Deserialize)]
133#[serde(tag = "status", content = "data")]
134pub enum StammdatenState {
135    /// No events yet.
136    #[default]
137    New,
138    /// Document received; `AcknowledgementDocument` not yet sent.
139    Received(ReceivedData),
140    /// `AcknowledgementDocument` sent; forwarding to ÜNB not yet done.
141    Acknowledged(ReceivedData),
142    /// Enriched document forwarded to ÜNB (VNB role only).
143    Forwarded(ReceivedData),
144    /// Process terminated due to a missed deadline.
145    DeadlineExpired {
146        /// Human-readable description of the expired deadline.
147        reason: String,
148    },
149}
150
151impl StammdatenState {
152    /// Stable string label for the current variant.
153    #[must_use]
154    pub fn label(&self) -> &'static str {
155        match self {
156            Self::New => "New",
157            Self::Received(_) => "Received",
158            Self::Acknowledged(_) => "Acknowledged",
159            Self::Forwarded(_) => "Forwarded",
160            Self::DeadlineExpired { .. } => "DeadlineExpired",
161        }
162    }
163}
164
165// ── Commands ──────────────────────────────────────────────────────────────────
166
167/// Commands for the Stammdaten workflow.
168///
169/// All domain values are pre-extracted by the transport layer before
170/// construction. `Workflow::handle` is pure — no I/O.
171#[derive(Clone)]
172pub enum StammdatenCommand {
173    /// Inbound `Stammdaten` document received and parsed by the transport layer.
174    Receive {
175        /// MRID (UUID) of the received document.
176        mrid: String,
177        /// GLN of the sender.
178        sender: String,
179        /// GLN of the receiver.
180        receiver: String,
181        /// Document type code (Z02/Z03/Z04/Z14).
182        doc_type: String,
183        /// Number of resource objects in the document.
184        anlagen_count: u32,
185        /// UTC receipt timestamp (ISO-8601 string).
186        received_at: String,
187    },
188    /// `AcknowledgementDocument` dispatched to the sender.
189    ///
190    /// The caller is responsible for building and enqueuing the outbound XML
191    /// via the outbox before issuing this command.
192    SendAcknowledgement {
193        /// MRID assigned to the outbound `AcknowledgementDocument`.
194        ack_mrid: String,
195    },
196    /// Enriched `Stammdaten` forwarded to ÜNB (VNB role only).
197    ///
198    /// The caller is responsible for building and enqueuing the upstream XML.
199    Forward {
200        /// MRID assigned to the upstream `Stammdaten` document.
201        upstream_mrid: String,
202    },
203    /// A registered deadline fired.
204    TimeoutExpired {
205        /// Unique ID of the expired deadline.
206        deadline_id: DeadlineId,
207        /// Label identifying the deadline type.
208        label: Box<str>,
209    },
210}
211
212impl CommandPayload for StammdatenCommand {}
213
214// ── Workflow ──────────────────────────────────────────────────────────────────
215
216/// Stammdatenübermittlung workflow for Redispatch 2.0.
217///
218/// Handles the reception, acknowledgement, and optional forwarding of
219/// `Stammdaten` documents exchanged between ANB, VNB, and ÜNB.
220///
221/// Spawn via [`mako_engine::process::Process`]:
222/// ```rust,ignore
223/// let process = ctx.spawn::<StammdatenWorkflow>(
224///     tenant_id,
225///     WorkflowId::new(WORKFLOW_NAME, "FV2025-10-01"),
226/// );
227/// ```
228pub struct StammdatenWorkflow;
229
230impl Workflow for StammdatenWorkflow {
231    type State = StammdatenState;
232    type Event = StammdatenEvent;
233    type Command = StammdatenCommand;
234
235    /// Fire deadline commands when the ACK or forward windows expire.
236    fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
237        match (deadline.label(), state) {
238            // 6h ACK window while Received; 1-Werktag forward window (VNB →
239            // ÜNB, BK6-20-060) while Acknowledged — the latter was previously
240            // defined but never enforced, so an acknowledged Stammdaten
241            // document that was never forwarded now expires visibly.
242            (ACK_WINDOW_LABEL, StammdatenState::Received(_))
243            | (FORWARD_WINDOW_LABEL, StammdatenState::Acknowledged { .. }) => {
244                Some(StammdatenCommand::TimeoutExpired {
245                    deadline_id: deadline.deadline_id(),
246                    label: deadline.label().into(),
247                })
248            }
249            _ => None,
250        }
251    }
252
253    fn apply(state: Self::State, event: &Self::Event) -> Self::State {
254        match event {
255            StammdatenEvent::Received {
256                mrid,
257                sender,
258                receiver,
259                doc_type,
260                anlagen_count,
261                received_at,
262            } => StammdatenState::Received(ReceivedData {
263                mrid: mrid.clone(),
264                sender: sender.clone(),
265                receiver: receiver.clone(),
266                doc_type: doc_type.clone(),
267                anlagen_count: *anlagen_count,
268                received_at: received_at.clone(),
269            }),
270
271            StammdatenEvent::Acknowledged { .. } => match state {
272                StammdatenState::Received(data) => StammdatenState::Acknowledged(data),
273                other => other,
274            },
275
276            StammdatenEvent::Forwarded { .. } => match state {
277                StammdatenState::Acknowledged(data) => StammdatenState::Forwarded(data),
278                other => other,
279            },
280
281            StammdatenEvent::DeadlineExpired { label, .. } => StammdatenState::DeadlineExpired {
282                reason: format!("deadline expired: {label}"),
283            },
284        }
285    }
286
287    fn handle(
288        state: &Self::State,
289        command: Self::Command,
290    ) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
291        match command {
292            StammdatenCommand::Receive {
293                mrid,
294                sender,
295                receiver,
296                doc_type,
297                anlagen_count,
298                received_at,
299            } => {
300                if !matches!(state, StammdatenState::New) {
301                    // Idempotent: document already received — this is a retry.
302                    return Ok(vec![].into());
303                }
304                Ok(vec![StammdatenEvent::Received {
305                    mrid,
306                    sender,
307                    receiver,
308                    doc_type,
309                    anlagen_count,
310                    received_at,
311                }]
312                .into())
313            }
314
315            StammdatenCommand::SendAcknowledgement { ack_mrid } => match state {
316                StammdatenState::Received(_) => {
317                    Ok(vec![StammdatenEvent::Acknowledged { ack_mrid }].into())
318                }
319                StammdatenState::Acknowledged(_) | StammdatenState::Forwarded(_) => {
320                    // Idempotent — acknowledgement already sent.
321                    Ok(vec![].into())
322                }
323                other => Err(WorkflowError::rejected(format!(
324                    "SendAcknowledgement not valid in state {}",
325                    other.label()
326                ))),
327            },
328
329            StammdatenCommand::Forward { upstream_mrid } => match state {
330                StammdatenState::Acknowledged(_) => {
331                    Ok(vec![StammdatenEvent::Forwarded { upstream_mrid }].into())
332                }
333                StammdatenState::Forwarded(_) => {
334                    // Idempotent.
335                    Ok(vec![].into())
336                }
337                other => Err(WorkflowError::rejected(format!(
338                    "Forward not valid in state {}",
339                    other.label()
340                ))),
341            },
342
343            StammdatenCommand::TimeoutExpired { deadline_id, label } => {
344                let is_forward_window = &*label == FORWARD_WINDOW_LABEL;
345                match state {
346                    // 1-Werktag forward window: an acknowledged document the
347                    // VNB never forwarded to the ÜNB expires visibly.
348                    StammdatenState::Acknowledged(_) if is_forward_window => {
349                        Ok(vec![StammdatenEvent::DeadlineExpired { deadline_id, label }].into())
350                    }
351                    // Terminal / already-progressed states — no-op.
352                    StammdatenState::Acknowledged(_)
353                    | StammdatenState::Forwarded(_)
354                    | StammdatenState::DeadlineExpired { .. } => Ok(vec![].into()),
355                    _ => Ok(vec![StammdatenEvent::DeadlineExpired { deadline_id, label }].into()),
356                }
357            }
358        }
359    }
360}
361
362#[cfg(test)]
363mod tests {
364    use super::*;
365    use mako_engine::ids::DeadlineId;
366
367    fn received_cmd() -> StammdatenCommand {
368        StammdatenCommand::Receive {
369            mrid: "mrid-001".into(),
370            sender: "4012345000001".into(),
371            receiver: "4012345000002".into(),
372            doc_type: "Z02".into(),
373            anlagen_count: 3,
374            received_at: "2025-10-15T10:00:00Z".into(),
375        }
376    }
377
378    #[test]
379    fn receive_transitions_new_to_received() {
380        let state = StammdatenState::New;
381        let output = StammdatenWorkflow::handle(&state, received_cmd()).unwrap();
382        assert_eq!(output.events.len(), 1);
383        let new_state = StammdatenWorkflow::apply(state, &output.events[0]);
384        assert!(matches!(new_state, StammdatenState::Received(_)));
385    }
386
387    #[test]
388    fn acknowledge_transitions_received_to_acknowledged() {
389        let state = StammdatenState::Received(ReceivedData {
390            mrid: "m".into(),
391            sender: "s".into(),
392            receiver: "r".into(),
393            doc_type: "Z02".into(),
394            anlagen_count: 1,
395            received_at: "2025-10-15T10:00:00Z".into(),
396        });
397        let output = StammdatenWorkflow::handle(
398            &state,
399            StammdatenCommand::SendAcknowledgement {
400                ack_mrid: "ack-001".into(),
401            },
402        )
403        .unwrap();
404        assert_eq!(output.events.len(), 1);
405        let new_state = StammdatenWorkflow::apply(state, &output.events[0]);
406        assert!(matches!(new_state, StammdatenState::Acknowledged(_)));
407    }
408
409    #[test]
410    fn forward_requires_acknowledged_state() {
411        let state = StammdatenState::Received(ReceivedData {
412            mrid: "m".into(),
413            sender: "s".into(),
414            receiver: "r".into(),
415            doc_type: "Z03".into(),
416            anlagen_count: 1,
417            received_at: "2025-10-15T10:00:00Z".into(),
418        });
419        let result = StammdatenWorkflow::handle(
420            &state,
421            StammdatenCommand::Forward {
422                upstream_mrid: "u".into(),
423            },
424        );
425        assert!(result.is_err());
426    }
427
428    #[test]
429    fn timeout_in_received_state_emits_deadline_expired() {
430        let state = StammdatenState::Received(ReceivedData {
431            mrid: "m".into(),
432            sender: "s".into(),
433            receiver: "r".into(),
434            doc_type: "Z02".into(),
435            anlagen_count: 1,
436            received_at: "2025-10-15T10:00:00Z".into(),
437        });
438        let output = StammdatenWorkflow::handle(
439            &state,
440            StammdatenCommand::TimeoutExpired {
441                deadline_id: DeadlineId::new(),
442                label: ACK_WINDOW_LABEL.into(),
443            },
444        )
445        .unwrap();
446        assert!(matches!(
447            output.events.as_slice(),
448            [StammdatenEvent::DeadlineExpired { .. }]
449        ));
450    }
451
452    #[test]
453    fn timeout_in_acknowledged_state_is_noop() {
454        let state = StammdatenState::Acknowledged(ReceivedData {
455            mrid: "m".into(),
456            sender: "s".into(),
457            receiver: "r".into(),
458            doc_type: "Z02".into(),
459            anlagen_count: 1,
460            received_at: "2025-10-15T10:00:00Z".into(),
461        });
462        let output = StammdatenWorkflow::handle(
463            &state,
464            StammdatenCommand::TimeoutExpired {
465                deadline_id: DeadlineId::new(),
466                label: ACK_WINDOW_LABEL.into(),
467            },
468        )
469        .unwrap();
470        assert!(output.events.is_empty());
471    }
472
473    #[test]
474    fn unforwarded_stammdaten_expires_after_the_forward_window() {
475        let data = ReceivedData {
476            mrid: "sd-001".into(),
477            sender: "4012345000001".into(),
478            receiver: "4012345000002".into(),
479            doc_type: "Z01".into(),
480            anlagen_count: 1,
481            received_at: "2025-10-15T09:00:00Z".into(),
482        };
483        let state = StammdatenState::Acknowledged(data);
484        let out = StammdatenWorkflow::handle(
485            &state,
486            StammdatenCommand::TimeoutExpired {
487                deadline_id: DeadlineId::new(),
488                label: FORWARD_WINDOW_LABEL.into(),
489            },
490        )
491        .expect("forward-window timeout handled");
492        assert_eq!(out.events.len(), 1, "1-Werktag forward window must fire");
493        // The 6h ACK label stays a no-op in Acknowledged.
494        let noop = StammdatenWorkflow::handle(
495            &state,
496            StammdatenCommand::TimeoutExpired {
497                deadline_id: DeadlineId::new(),
498                label: ACK_WINDOW_LABEL.into(),
499            },
500        )
501        .expect("ack timeout in acknowledged is noop");
502        assert!(noop.events.is_empty());
503    }
504}