ironflow_store/entities/signal.rs
1//! Signal entities -- external messages that resume waiting runs.
2//!
3//! A signal is named (what happened, e.g. `"ci.pipeline_finished"`) and keyed
4//! (which occurrence, e.g. a commit SHA). A run waiting through
5//! `ctx.wait_for_signal` on the same `(name, key)` pair is resumed when a
6//! signal with a valid payload is delivered.
7
8use chrono::{DateTime, Utc};
9use serde::{Deserialize, Serialize};
10use serde_json::Value;
11use uuid::Uuid;
12
13/// A persisted signal.
14///
15/// Signals are kept for a retention period after `received_at`, so a run that
16/// opens its wait step after the signal arrived still finds it.
17///
18/// # Examples
19///
20/// ```
21/// use chrono::Utc;
22/// use ironflow_store::entities::Signal;
23/// use serde_json::json;
24/// use uuid::Uuid;
25///
26/// let signal = Signal {
27/// id: Uuid::now_v7(),
28/// name: "ci.pipeline_finished".to_string(),
29/// key: "4f2a9c1".to_string(),
30/// payload: json!({"status": "success"}),
31/// idempotency_id: Some("delivery-42".to_string()),
32/// received_at: Utc::now(),
33/// };
34/// assert_eq!(signal.key, "4f2a9c1");
35/// ```
36#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
37pub struct Signal {
38 /// Unique signal ID (UUID v7).
39 pub id: Uuid,
40 /// Signal name, e.g. `"ci.pipeline_finished"`.
41 pub name: String,
42 /// Occurrence key, e.g. a commit SHA.
43 pub key: String,
44 /// JSON payload carried by the signal.
45 pub payload: Value,
46 /// Caller-provided deduplication ID. A second signal with the same ID is
47 /// not stored nor delivered again.
48 pub idempotency_id: Option<String>,
49 /// When the signal was received.
50 pub received_at: DateTime<Utc>,
51}
52
53/// Request to store a new signal.
54///
55/// # Examples
56///
57/// ```
58/// use ironflow_store::entities::NewSignal;
59/// use serde_json::json;
60///
61/// let signal = NewSignal {
62/// name: "ci.pipeline_finished".to_string(),
63/// key: "4f2a9c1".to_string(),
64/// payload: json!({"status": "success"}),
65/// idempotency_id: None,
66/// };
67/// assert!(signal.idempotency_id.is_none());
68/// ```
69#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
70pub struct NewSignal {
71 /// Signal name.
72 pub name: String,
73 /// Occurrence key.
74 pub key: String,
75 /// JSON payload.
76 pub payload: Value,
77 /// Optional deduplication ID.
78 pub idempotency_id: Option<String>,
79}
80
81/// Outcome of [`SignalStore::insert_signal`](crate::signal_store::SignalStore::insert_signal).
82///
83/// # Examples
84///
85/// ```
86/// use chrono::Utc;
87/// use ironflow_store::entities::{Signal, SignalInsert};
88/// use serde_json::json;
89/// use uuid::Uuid;
90///
91/// let signal = Signal {
92/// id: Uuid::now_v7(),
93/// name: "demo.done".to_string(),
94/// key: "k1".to_string(),
95/// payload: json!({}),
96/// idempotency_id: Some("d-1".to_string()),
97/// received_at: Utc::now(),
98/// };
99/// let insert = SignalInsert::Duplicate(signal);
100/// assert!(insert.is_duplicate());
101/// assert_eq!(insert.signal().key, "k1");
102/// ```
103#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
104pub enum SignalInsert {
105 /// The signal was stored.
106 Created(Signal),
107 /// A signal with the same idempotency ID already existed; it is returned
108 /// unchanged and nothing was stored.
109 Duplicate(Signal),
110}
111
112impl SignalInsert {
113 /// The stored signal, new or pre-existing.
114 ///
115 /// # Examples
116 ///
117 /// ```
118 /// use chrono::Utc;
119 /// use ironflow_store::entities::{Signal, SignalInsert};
120 /// use serde_json::json;
121 /// use uuid::Uuid;
122 ///
123 /// let signal = Signal {
124 /// id: Uuid::now_v7(),
125 /// name: "demo.done".to_string(),
126 /// key: "k1".to_string(),
127 /// payload: json!({}),
128 /// idempotency_id: None,
129 /// received_at: Utc::now(),
130 /// };
131 /// let insert = SignalInsert::Created(signal.clone());
132 /// assert_eq!(insert.signal(), &signal);
133 /// ```
134 pub fn signal(&self) -> &Signal {
135 match self {
136 SignalInsert::Created(signal) | SignalInsert::Duplicate(signal) => signal,
137 }
138 }
139
140 /// Whether the insert hit an existing idempotency ID.
141 ///
142 /// # Examples
143 ///
144 /// ```
145 /// use chrono::Utc;
146 /// use ironflow_store::entities::{Signal, SignalInsert};
147 /// use serde_json::json;
148 /// use uuid::Uuid;
149 ///
150 /// let signal = Signal {
151 /// id: Uuid::now_v7(),
152 /// name: "demo.done".to_string(),
153 /// key: "k1".to_string(),
154 /// payload: json!({}),
155 /// idempotency_id: None,
156 /// received_at: Utc::now(),
157 /// };
158 /// assert!(!SignalInsert::Created(signal).is_duplicate());
159 /// ```
160 pub fn is_duplicate(&self) -> bool {
161 matches!(self, SignalInsert::Duplicate(_))
162 }
163}
164
165/// Filter criteria for listing signals.
166///
167/// Every field is optional; set fields are combined with AND and match exactly.
168///
169/// # Examples
170///
171/// ```
172/// use ironflow_store::entities::SignalFilter;
173///
174/// let filter = SignalFilter {
175/// name: Some("ci.pipeline_finished".to_string()),
176/// ..SignalFilter::default()
177/// };
178/// assert!(filter.key.is_none());
179/// ```
180#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
181pub struct SignalFilter {
182 /// Exact signal name.
183 pub name: Option<String>,
184 /// Exact occurrence key.
185 pub key: Option<String>,
186}
187
188/// Outcome of [`SignalStore::resolve_signal_step`](crate::signal_store::SignalStore::resolve_signal_step).
189///
190/// # Examples
191///
192/// ```
193/// use ironflow_store::entities::SignalStepResolution;
194/// use uuid::Uuid;
195///
196/// # fn main() -> Result<(), serde_json::Error> {
197/// let resolution = SignalStepResolution::Resolved {
198/// run_id: Uuid::now_v7(),
199/// run_resumed: true,
200/// };
201/// let json = serde_json::to_value(&resolution)?;
202/// assert_eq!(json["outcome"], "resolved");
203/// # Ok(())
204/// # }
205/// ```
206#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
207#[serde(tag = "outcome", rename_all = "snake_case")]
208pub enum SignalStepResolution {
209 /// The step was waiting and is now completed with the given output.
210 Resolved {
211 /// Run owning the step.
212 run_id: Uuid,
213 /// `true` when the run was `Sleeping` and went back to `Pending`.
214 /// `false` when the run was still executing (it will see the step
215 /// completed on its own).
216 run_resumed: bool,
217 },
218 /// The step was no longer waiting: another delivery or a timeout won.
219 NotWaiting {
220 /// Output recorded by whoever resolved the step first.
221 output: Option<Value>,
222 },
223}
224
225#[cfg(test)]
226mod tests {
227 use serde_json::{from_value, json, to_value};
228
229 use super::*;
230
231 fn signal() -> Signal {
232 Signal {
233 id: Uuid::now_v7(),
234 name: "demo.done".to_string(),
235 key: "k1".to_string(),
236 payload: json!({"ok": true}),
237 idempotency_id: None,
238 received_at: Utc::now(),
239 }
240 }
241
242 #[test]
243 fn signal_insert_exposes_the_signal() {
244 let s = signal();
245 assert_eq!(SignalInsert::Created(s.clone()).signal(), &s);
246 assert_eq!(SignalInsert::Duplicate(s.clone()).signal(), &s);
247 assert!(!SignalInsert::Created(s.clone()).is_duplicate());
248 assert!(SignalInsert::Duplicate(s).is_duplicate());
249 }
250
251 #[test]
252 fn signal_step_resolution_serde_roundtrip() {
253 let resolved = SignalStepResolution::Resolved {
254 run_id: Uuid::now_v7(),
255 run_resumed: false,
256 };
257 let json = to_value(&resolved).unwrap();
258 assert_eq!(json["outcome"], "resolved");
259 assert_eq!(from_value::<SignalStepResolution>(json).unwrap(), resolved);
260
261 let not_waiting = SignalStepResolution::NotWaiting {
262 output: Some(json!({"timed_out": true})),
263 };
264 let json = to_value(¬_waiting).unwrap();
265 assert_eq!(json["outcome"], "not_waiting");
266 assert_eq!(
267 from_value::<SignalStepResolution>(json).unwrap(),
268 not_waiting
269 );
270 }
271
272 #[test]
273 fn signal_serde_roundtrip() {
274 let s = signal();
275 let back: Signal = from_value(to_value(&s).unwrap()).unwrap();
276 assert_eq!(back, s);
277 }
278}