1use mako_engine::{
24 deadline::Deadline,
25 error::WorkflowError,
26 ids::DeadlineId,
27 workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
28};
29use serde::{Deserialize, Serialize};
30
31#[derive(Debug, Clone, Serialize, Deserialize)]
35#[serde(tag = "type", content = "data")]
36pub enum AckForwardEvent {
37 Received {
39 mrid: String,
41 doc_type: String,
43 sender: String,
45 receiver: String,
47 received_at: String,
49 },
50 Acknowledged {
52 ack_mrid: String,
54 },
55 Forwarded {
57 upstream_mrid: String,
59 },
60 DeadlineExpired {
62 deadline_id: DeadlineId,
64 label: Box<str>,
66 },
67}
68
69#[derive(Clone)]
71pub enum AckForwardCommand {
72 Receive {
74 mrid: String,
76 doc_type: String,
78 sender: String,
80 receiver: String,
82 received_at: String,
84 },
85 Acknowledge {
87 ack_mrid: String,
89 },
90 Forward {
92 upstream_mrid: String,
94 },
95 TimeoutExpired {
97 deadline_id: DeadlineId,
99 label: Box<str>,
101 },
102}
103
104impl CommandPayload for AckForwardCommand {}
105
106#[derive(Debug, Clone, Serialize, Deserialize)]
108#[serde(deny_unknown_fields)]
109pub struct ReceivedData {
110 pub mrid: String,
112 pub doc_type: String,
114 pub sender: String,
116 pub receiver: String,
118 pub received_at: String,
120}
121
122#[derive(Debug, Clone, Default, Serialize, Deserialize)]
124#[serde(tag = "status", content = "data")]
125pub enum AckForwardState {
126 #[default]
128 New,
129 Received(ReceivedData),
131 Acknowledged(ReceivedData),
133 Forwarded(ReceivedData),
135 DeadlineExpired {
137 reason: String,
139 },
140}
141
142impl AckForwardState {
143 #[must_use]
145 pub fn label(&self) -> &'static str {
146 match self {
147 Self::New => "New",
148 Self::Received(_) => "Received",
149 Self::Acknowledged(_) => "Acknowledged",
150 Self::Forwarded(_) => "Forwarded",
151 Self::DeadlineExpired { .. } => "DeadlineExpired",
152 }
153 }
154}
155
156impl EventPayload for AckForwardEvent {
157 fn event_type(&self) -> &'static str {
158 match self {
159 Self::Received { .. } => "AckForwardReceived",
160 Self::Acknowledged { .. } => "AckForwardAcknowledged",
161 Self::Forwarded { .. } => "AckForwardForwarded",
162 Self::DeadlineExpired { .. } => "AckForwardDeadlineExpired",
163 }
164 }
165}
166
167macro_rules! define_workflow_event {
179 ($event_type:ident, $prefix:expr) => {
180 #[derive(Debug, Clone, Serialize, Deserialize)]
187 #[serde(transparent)]
188 pub struct $event_type(pub AckForwardEvent);
189
190 impl From<AckForwardEvent> for $event_type {
191 fn from(e: AckForwardEvent) -> Self {
192 Self(e)
193 }
194 }
195
196 impl From<$event_type> for AckForwardEvent {
197 fn from(e: $event_type) -> AckForwardEvent {
198 e.0
199 }
200 }
201
202 impl EventPayload for $event_type {
203 fn event_type(&self) -> &'static str {
204 match &self.0 {
205 AckForwardEvent::Received { .. } => concat!($prefix, "Received"),
206 AckForwardEvent::Acknowledged { .. } => concat!($prefix, "Acknowledged"),
207 AckForwardEvent::Forwarded { .. } => concat!($prefix, "Forwarded"),
208 AckForwardEvent::DeadlineExpired { .. } => {
209 concat!($prefix, "DeadlineExpired")
210 }
211 }
212 }
213 }
214 };
215}
216
217define_workflow_event!(VerfuegbarkeitEvent, "Verfuegbarkeit");
218define_workflow_event!(NetzengpassEvent, "Netzengpass");
219define_workflow_event!(KaskadeEvent, "Kaskade");
220define_workflow_event!(PlanungsdatenEvent, "Planungsdaten");
221define_workflow_event!(StatusanfrageEvent, "Statusanfrage");
222define_workflow_event!(KostenblattEvent, "Kostenblatt");
223
224pub(crate) fn apply(state: AckForwardState, event: &AckForwardEvent) -> AckForwardState {
228 match event {
229 AckForwardEvent::Received {
230 mrid,
231 doc_type,
232 sender,
233 receiver,
234 received_at,
235 } => AckForwardState::Received(ReceivedData {
236 mrid: mrid.clone(),
237 doc_type: doc_type.clone(),
238 sender: sender.clone(),
239 receiver: receiver.clone(),
240 received_at: received_at.clone(),
241 }),
242
243 AckForwardEvent::Acknowledged { .. } => match state {
244 AckForwardState::Received(data) => AckForwardState::Acknowledged(data),
245 other => other,
246 },
247
248 AckForwardEvent::Forwarded { .. } => match state {
249 AckForwardState::Acknowledged(data) => AckForwardState::Forwarded(data),
250 other => other,
251 },
252
253 AckForwardEvent::DeadlineExpired { label, .. } => AckForwardState::DeadlineExpired {
254 reason: format!("deadline expired: {label}"),
255 },
256 }
257}
258
259pub(crate) fn handle(
261 state: &AckForwardState,
262 command: AckForwardCommand,
263 ack_window_label: &str,
264) -> Result<WorkflowOutput<AckForwardEvent>, WorkflowError> {
265 match command {
266 AckForwardCommand::Receive {
267 mrid,
268 doc_type,
269 sender,
270 receiver,
271 received_at,
272 } => {
273 if !matches!(state, AckForwardState::New) {
274 return Ok(vec![].into());
275 }
276 Ok(vec![AckForwardEvent::Received {
277 mrid,
278 doc_type,
279 sender,
280 receiver,
281 received_at,
282 }]
283 .into())
284 }
285
286 AckForwardCommand::Acknowledge { ack_mrid } => match state {
287 AckForwardState::Received(_) => {
288 Ok(vec![AckForwardEvent::Acknowledged { ack_mrid }].into())
289 }
290 AckForwardState::Acknowledged(_) | AckForwardState::Forwarded(_) => Ok(vec![].into()),
291 other => Err(WorkflowError::rejected(format!(
292 "Acknowledge not valid in state {}",
293 other.label()
294 ))),
295 },
296
297 AckForwardCommand::Forward { upstream_mrid } => match state {
298 AckForwardState::Acknowledged(_) => {
299 Ok(vec![AckForwardEvent::Forwarded { upstream_mrid }].into())
300 }
301 AckForwardState::Forwarded(_) => Ok(vec![].into()),
302 other => Err(WorkflowError::rejected(format!(
303 "Forward not valid in state {}",
304 other.label()
305 ))),
306 },
307
308 AckForwardCommand::TimeoutExpired { deadline_id, label } => match state {
309 AckForwardState::Acknowledged(_)
310 | AckForwardState::Forwarded(_)
311 | AckForwardState::DeadlineExpired { .. } => Ok(vec![].into()),
312 _ => {
313 let _ = ack_window_label; Ok(vec![AckForwardEvent::DeadlineExpired { deadline_id, label }].into())
315 }
316 },
317 }
318}
319
320macro_rules! ack_forward_workflow {
323 (
324 $(#[$meta:meta])*
325 $name:ident,
326 $event_newtype:ident,
327 $workflow_name:expr,
328 $ack_label:expr,
329 $event_prefix:expr $(,)?
330 ) => {
331 $(#[$meta])*
332 pub struct $name;
333
334 impl Workflow for $name {
335 type State = AckForwardState;
336 type Event = $event_newtype;
337 type Command = AckForwardCommand;
338
339 fn on_deadline(
340 deadline: &Deadline,
341 state: &Self::State,
342 ) -> Option<Self::Command> {
343 if deadline.label() == $ack_label {
344 if matches!(state, AckForwardState::Received(_)) {
345 return Some(AckForwardCommand::TimeoutExpired {
346 deadline_id: deadline.deadline_id(),
347 label: deadline.label().into(),
348 });
349 }
350 }
351 None
352 }
353
354 fn apply(state: Self::State, event: &Self::Event) -> Self::State {
355 crate::ack_forward::apply(state, &event.0)
356 }
357
358 fn handle(
359 state: &Self::State,
360 command: Self::Command,
361 ) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
362 let output = crate::ack_forward::handle(state, command, $ack_label)?;
363 Ok(WorkflowOutput::with_outbox(
364 output.events.into_iter().map($event_newtype).collect(),
365 output.outbox,
366 ))
367 }
368 }
369
370 impl $name {
371 #[must_use]
373 pub fn event_prefix() -> &'static str {
374 $event_prefix
375 }
376 }
377 };
378}
379
380ack_forward_workflow!(
381 VerfuegbarkeitWorkflow,
386 VerfuegbarkeitEvent,
387 "redispatch-verfuegbarkeit",
388 "redispatch-verfuegbarkeit-ack-window",
389 "Verfuegbarkeit",
390);
391
392ack_forward_workflow!(
393 NetzengpassWorkflow,
398 NetzengpassEvent,
399 "redispatch-netzengpass",
400 "redispatch-netzengpass-ack-window",
401 "Netzengpass",
402);
403
404ack_forward_workflow!(
405 KaskadeWorkflow,
411 KaskadeEvent,
412 "redispatch-kaskade",
413 "redispatch-kaskade-ack-window",
414 "Kaskade",
415);
416
417ack_forward_workflow!(
418 PlanungsdatenWorkflow,
423 PlanungsdatenEvent,
424 "redispatch-planungsdaten",
425 "redispatch-planungsdaten-ack-window",
426 "Planungsdaten",
427);
428
429ack_forward_workflow!(
430 StatusanfrageWorkflow,
440 StatusanfrageEvent,
441 "redispatch-statusanfrage",
442 "redispatch-statusanfrage-response-window",
443 "Statusanfrage",
444);
445
446ack_forward_workflow!(
447 KostenblattWorkflow,
456 KostenblattEvent,
457 "redispatch-kostenblatt",
458 "redispatch-kostenblatt-ack-window",
459 "Kostenblatt",
460);
461
462pub mod names {
465 pub const VERFUEGBARKEIT: &str = "redispatch-verfuegbarkeit";
467 pub const NETZENGPASS: &str = "redispatch-netzengpass";
469 pub const KASKADE: &str = "redispatch-kaskade";
471 pub const PLANUNGSDATEN: &str = "redispatch-planungsdaten";
473 pub const STATUSANFRAGE: &str = "redispatch-statusanfrage";
475 pub const KOSTENBLATT: &str = "redispatch-kostenblatt";
477}
478
479#[cfg(test)]
480mod tests {
481 use super::*;
482 use mako_engine::workflow::EventPayload;
483
484 #[test]
485 fn verfuegbarkeit_receive_to_acknowledged() {
486 let state = AckForwardState::New;
487 let output = VerfuegbarkeitWorkflow::handle(
488 &state,
489 AckForwardCommand::Receive {
490 mrid: "m1".into(),
491 doc_type: "Unavailability".into(),
492 sender: "s".into(),
493 receiver: "r".into(),
494 received_at: "2025-10-15T10:00:00Z".into(),
495 },
496 )
497 .unwrap();
498 assert_eq!(output.events.len(), 1);
499
500 let state2 = VerfuegbarkeitWorkflow::apply(state, &output.events[0]);
501 assert!(matches!(state2, AckForwardState::Received(_)));
502
503 let output2 = VerfuegbarkeitWorkflow::handle(
504 &state2,
505 AckForwardCommand::Acknowledge {
506 ack_mrid: "ack-1".into(),
507 },
508 )
509 .unwrap();
510 let state3 = VerfuegbarkeitWorkflow::apply(state2, &output2.events[0]);
511 assert!(matches!(state3, AckForwardState::Acknowledged(_)));
512 }
513
514 #[test]
515 fn kaskade_forward_requires_acknowledged_state() {
516 let state = AckForwardState::Received(ReceivedData {
517 mrid: "m".into(),
518 doc_type: "Kaskade".into(),
519 sender: "s".into(),
520 receiver: "r".into(),
521 received_at: "2025-10-15T10:00:00Z".into(),
522 });
523 let result = KaskadeWorkflow::handle(
524 &state,
525 AckForwardCommand::Forward {
526 upstream_mrid: "u".into(),
527 },
528 );
529 assert!(result.is_err());
530 }
531
532 #[test]
534 fn event_types_are_unique_per_workflow() {
535 let inner = AckForwardEvent::Received {
536 mrid: "m".into(),
537 doc_type: "X".into(),
538 sender: "s".into(),
539 receiver: "r".into(),
540 received_at: "t".into(),
541 };
542
543 let types: Vec<&'static str> = vec![
544 VerfuegbarkeitEvent(inner.clone()).event_type(),
545 NetzengpassEvent(inner.clone()).event_type(),
546 KaskadeEvent(inner.clone()).event_type(),
547 PlanungsdatenEvent(inner.clone()).event_type(),
548 StatusanfrageEvent(inner.clone()).event_type(),
549 KostenblattEvent(inner.clone()).event_type(),
550 ];
551
552 let unique: std::collections::HashSet<_> = types.iter().collect();
554 assert_eq!(
555 unique.len(),
556 types.len(),
557 "event_type() strings must be unique across all ack-forward workflows: {types:?}"
558 );
559
560 for t in &types {
562 assert!(
563 !t.starts_with("AckForward"),
564 "event_type '{t}' must not use the generic AckForward prefix"
565 );
566 }
567 }
568}