1use mako_engine::{
26 deadline::Deadline,
27 error::WorkflowError,
28 ids::DeadlineId,
29 workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
30};
31use serde::{Deserialize, Serialize};
32
33pub const WORKFLOW_NAME: &str = "redispatch-stammdaten";
37
38pub const ACK_WINDOW_LABEL: &str = "redispatch-stammdaten-ack-window";
44
45pub const FORWARD_WINDOW_LABEL: &str = "redispatch-stammdaten-forward-window";
50
51#[derive(Debug, Clone, Serialize, Deserialize)]
55#[serde(tag = "type", content = "data")]
56pub enum StammdatenEvent {
57 Received {
59 mrid: String,
61 sender: String,
63 receiver: String,
65 doc_type: String,
67 anlagen_count: u32,
69 received_at: String,
71 },
72 Acknowledged {
74 ack_mrid: String,
76 },
77 Forwarded {
79 upstream_mrid: String,
81 },
82 DeadlineExpired {
84 deadline_id: DeadlineId,
86 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#[derive(Debug, Clone, Serialize, Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct ReceivedData {
108 pub mrid: String,
110 pub sender: String,
112 pub receiver: String,
114 pub doc_type: String,
116 pub anlagen_count: u32,
118 pub received_at: String,
120}
121
122#[derive(Debug, Clone, Default, Serialize, Deserialize)]
133#[serde(tag = "status", content = "data")]
134pub enum StammdatenState {
135 #[default]
137 New,
138 Received(ReceivedData),
140 Acknowledged(ReceivedData),
142 Forwarded(ReceivedData),
144 DeadlineExpired {
146 reason: String,
148 },
149}
150
151impl StammdatenState {
152 #[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#[derive(Clone)]
172pub enum StammdatenCommand {
173 Receive {
175 mrid: String,
177 sender: String,
179 receiver: String,
181 doc_type: String,
183 anlagen_count: u32,
185 received_at: String,
187 },
188 SendAcknowledgement {
193 ack_mrid: String,
195 },
196 Forward {
200 upstream_mrid: String,
202 },
203 TimeoutExpired {
205 deadline_id: DeadlineId,
207 label: Box<str>,
209 },
210}
211
212impl CommandPayload for StammdatenCommand {}
213
214pub struct StammdatenWorkflow;
229
230impl Workflow for StammdatenWorkflow {
231 type State = StammdatenState;
232 type Event = StammdatenEvent;
233 type Command = StammdatenCommand;
234
235 fn on_deadline(deadline: &Deadline, state: &Self::State) -> Option<Self::Command> {
237 match (deadline.label(), state) {
238 (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 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 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 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 StammdatenState::Acknowledged(_) if is_forward_window => {
349 Ok(vec![StammdatenEvent::DeadlineExpired { deadline_id, label }].into())
350 }
351 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 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}