1use chrono::{DateTime, Utc};
2use serde::de::DeserializeOwned;
3use serde::{Deserialize, Serialize};
4
5use crate::error::{FlowError, Result};
6
7use super::JsonValue;
8
9#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
14#[non_exhaustive]
15pub struct WorkflowSignal {
16 pub signal_id: String,
18 pub name: String,
20 pub payload: JsonValue,
22}
23
24impl WorkflowSignal {
25 pub fn new(signal_id: impl Into<String>, name: impl Into<String>, payload: JsonValue) -> Self {
27 Self {
28 signal_id: signal_id.into(),
29 name: name.into(),
30 payload,
31 }
32 }
33
34 pub(crate) fn validate(&self) -> Result<()> {
35 if self.signal_id.trim().is_empty() {
36 return Err(FlowError::InvalidTransition(
37 "workflow signal id must not be empty".to_string(),
38 ));
39 }
40 validate_signal_name(&self.name)
41 }
42}
43
44#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
46#[non_exhaustive]
47pub struct WorkflowSignalSnapshot {
48 pub signal_id: String,
50 pub name: String,
52 pub payload: JsonValue,
54 pub received_at: DateTime<Utc>,
56 pub received_sequence: u64,
58 #[serde(default, skip_serializing_if = "Option::is_none")]
60 pub consumed_by: Option<String>,
61}
62
63impl WorkflowSignalSnapshot {
64 pub fn payload_as<T>(&self) -> Result<T>
66 where
67 T: DeserializeOwned,
68 {
69 serde_json::from_value(self.payload.clone()).map_err(FlowError::from)
70 }
71}
72
73#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
75#[non_exhaustive]
76#[serde(rename_all = "snake_case")]
77pub enum SignalWaitStatus {
78 Waiting,
80 Completed,
82 Cancelled,
84}
85
86#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
88#[non_exhaustive]
89pub struct SignalWaitSnapshot {
90 pub wait_id: String,
92 pub signal_name: String,
94 pub status: SignalWaitStatus,
96 pub created_at: DateTime<Utc>,
98 pub created_sequence: u64,
100 #[serde(default, skip_serializing_if = "Option::is_none")]
102 pub signal_id: Option<String>,
103 #[serde(default, skip_serializing_if = "Option::is_none")]
105 pub completed_at: Option<DateTime<Utc>>,
106 #[serde(default, skip_serializing_if = "Option::is_none")]
108 pub completed_sequence: Option<u64>,
109}
110
111pub(crate) fn validate_signal_name(name: &str) -> Result<()> {
112 if name.trim().is_empty() {
113 return Err(FlowError::InvalidTransition(
114 "workflow signal name must not be empty".to_string(),
115 ));
116 }
117 Ok(())
118}
119
120pub(crate) fn validate_signal_wait(wait_id: &str, signal_name: &str) -> Result<()> {
121 if wait_id.trim().is_empty() {
122 return Err(FlowError::InvalidTransition(
123 "workflow signal wait id must not be empty".to_string(),
124 ));
125 }
126 validate_signal_name(signal_name)
127}