Skip to main content

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(&not_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}