Skip to main content

ironflow_engine/
signal.rs

1//! Signals -- external messages, named and keyed, that resume waiting runs.
2//!
3//! A workflow waits with
4//! [`WorkflowContext::wait_for_signal`](crate::context::WorkflowContext::wait_for_signal);
5//! a producer (a webhook handler, a script, another service) delivers with
6//! [`Engine::send_signal`](crate::engine::Engine::send_signal) or
7//! `POST /api/v1/signals`.
8//!
9//! The **name** says what happened (`"ci.pipeline_finished"`), the **key**
10//! says which occurrence (a commit SHA). Every run waiting on the same
11//! `(name, key)` pair receives the signal.
12//!
13//! # Examples
14//!
15//! ```
16//! use ironflow_engine::signal::Signal;
17//! use schemars::JsonSchema;
18//! use serde::{Deserialize, Serialize};
19//!
20//! #[derive(Serialize, Deserialize, JsonSchema)]
21//! struct PipelineFinished {
22//!     status: String,
23//! }
24//!
25//! impl Signal for PipelineFinished {
26//!     const NAME: &'static str = "ci.pipeline_finished";
27//! }
28//!
29//! assert_eq!(PipelineFinished::NAME, "ci.pipeline_finished");
30//! ```
31
32use jsonschema::validator_for;
33use schemars::JsonSchema;
34use serde::de::DeserializeOwned;
35use serde::{Deserialize, Serialize};
36use serde_json::{Value, json};
37use uuid::Uuid;
38
39use ironflow_store::entities::Signal as StoredSignal;
40
41/// A typed signal a workflow can wait for.
42///
43/// The payload is the type itself: it is serialized when sent, validated
44/// against its JSON schema on delivery, and deserialized for the waiting
45/// handler.
46///
47/// # Examples
48///
49/// ```
50/// use ironflow_engine::signal::Signal;
51/// use schemars::JsonSchema;
52/// use serde::{Deserialize, Serialize};
53///
54/// #[derive(Serialize, Deserialize, JsonSchema)]
55/// struct PipelineFinished {
56///     status: String,
57///     pipeline_id: u64,
58/// }
59///
60/// impl Signal for PipelineFinished {
61///     const NAME: &'static str = "ci.pipeline_finished";
62/// }
63/// ```
64pub trait Signal: DeserializeOwned + Serialize + JsonSchema {
65    /// Signal name, e.g. `"ci.pipeline_finished"`. Must not be empty.
66    const NAME: &'static str;
67}
68
69/// Key of the payload JSON schema in a signal step's input.
70pub const SIGNAL_SCHEMA_KEY: &str = "schema";
71
72/// Key of the timeout flag in a signal step's output.
73pub const SIGNAL_TIMED_OUT_KEY: &str = "timed_out";
74
75/// Result of delivering a signal.
76///
77/// # Examples
78///
79/// ```
80/// use ironflow_engine::signal::SignalDelivery;
81/// use uuid::Uuid;
82///
83/// let delivery = SignalDelivery {
84///     signal_id: Uuid::now_v7(),
85///     duplicate: false,
86///     resumed: Vec::new(),
87///     rejected: Vec::new(),
88/// };
89/// assert!(delivery.resumed.is_empty());
90/// ```
91#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
92pub struct SignalDelivery {
93    /// ID of the stored signal (the pre-existing one on a duplicate).
94    pub signal_id: Uuid,
95    /// `true` when the idempotency ID was already used: nothing was delivered.
96    pub duplicate: bool,
97    /// Waiting steps the signal resolved.
98    pub resumed: Vec<SignalResumed>,
99    /// Waiting steps whose payload schema the signal did not match. These
100    /// steps keep waiting.
101    pub rejected: Vec<SignalRejected>,
102}
103
104/// A waiting step resolved by a signal.
105///
106/// # Examples
107///
108/// ```
109/// use ironflow_engine::signal::SignalResumed;
110/// use uuid::Uuid;
111///
112/// let resumed = SignalResumed { run_id: Uuid::now_v7(), step_id: Uuid::now_v7() };
113/// assert_ne!(resumed.run_id, resumed.step_id);
114/// ```
115#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
116pub struct SignalResumed {
117    /// Run owning the step.
118    pub run_id: Uuid,
119    /// The resolved signal step.
120    pub step_id: Uuid,
121}
122
123/// A waiting step a signal could not resolve.
124///
125/// # Examples
126///
127/// ```
128/// use ironflow_engine::signal::SignalRejected;
129/// use uuid::Uuid;
130///
131/// let rejected = SignalRejected {
132///     run_id: Uuid::now_v7(),
133///     step_id: Uuid::now_v7(),
134///     error: "\"status\" is a required property".to_string(),
135/// };
136/// assert!(rejected.error.contains("status"));
137/// ```
138#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
139pub struct SignalRejected {
140    /// Run owning the step.
141    pub run_id: Uuid,
142    /// The step that keeps waiting.
143    pub step_id: Uuid,
144    /// Why the payload was refused.
145    pub error: String,
146}
147
148/// Output recorded on a signal step resolved by `signal`.
149pub(crate) fn received_output(signal: &StoredSignal) -> Value {
150    json!({
151        SIGNAL_TIMED_OUT_KEY: false,
152        "signal_id": signal.id,
153        "payload": signal.payload,
154    })
155}
156
157/// Output recorded on a signal step whose deadline passed.
158pub(crate) fn timed_out_output() -> Value {
159    json!({ SIGNAL_TIMED_OUT_KEY: true })
160}
161
162/// Validate `payload` against a JSON schema, joining every violation.
163pub(crate) fn validate_payload(schema: &Value, payload: &Value) -> Result<(), String> {
164    let validator = validator_for(schema).map_err(|e| format!("invalid payload schema: {e}"))?;
165    let errors: Vec<String> = validator
166        .iter_errors(payload)
167        .map(|e| e.to_string())
168        .collect();
169    if errors.is_empty() {
170        Ok(())
171    } else {
172        Err(errors.join("; "))
173    }
174}
175
176/// Validate `payload` against the schema a signal step stored in its input.
177pub(crate) fn validate_step_payload(input: Option<&Value>, payload: &Value) -> Result<(), String> {
178    let schema = input
179        .and_then(|i| i.get(SIGNAL_SCHEMA_KEY))
180        .ok_or_else(|| "step has no stored schema".to_string())?;
181    validate_payload(schema, payload)
182}
183
184#[cfg(test)]
185mod tests {
186    use chrono::Utc;
187    use schemars::schema_for;
188    use serde_json::to_value;
189
190    use super::*;
191
192    #[derive(Serialize, Deserialize, JsonSchema)]
193    struct PipelineFinished {
194        status: String,
195    }
196
197    fn schema() -> Value {
198        to_value(schema_for!(PipelineFinished)).unwrap()
199    }
200
201    #[test]
202    fn validate_payload_accepts_a_matching_payload() {
203        assert_eq!(
204            validate_payload(&schema(), &json!({"status": "success"})),
205            Ok(())
206        );
207    }
208
209    #[test]
210    fn validate_payload_rejects_a_mismatching_payload() {
211        let err = validate_payload(&schema(), &json!({"state": 1})).unwrap_err();
212        assert!(err.contains("status"), "got {err}");
213    }
214
215    #[test]
216    fn validate_payload_rejects_an_invalid_schema() {
217        let err = validate_payload(&json!({"type": 12}), &json!({})).unwrap_err();
218        assert!(err.contains("invalid payload schema"), "got {err}");
219    }
220
221    #[test]
222    fn validate_step_payload_requires_a_stored_schema() {
223        let err = validate_step_payload(None, &json!({})).unwrap_err();
224        assert_eq!(err, "step has no stored schema");
225
226        let input = json!({ SIGNAL_SCHEMA_KEY: schema() });
227        assert_eq!(
228            validate_step_payload(Some(&input), &json!({"status": "ok"})),
229            Ok(())
230        );
231    }
232
233    #[test]
234    fn signal_outputs_carry_the_timeout_flag() {
235        let signal = StoredSignal {
236            id: Uuid::now_v7(),
237            name: "ci.done".to_string(),
238            key: "abc".to_string(),
239            payload: json!({"status": "success"}),
240            idempotency_id: None,
241            received_at: Utc::now(),
242        };
243        let received = received_output(&signal);
244        assert_eq!(received[SIGNAL_TIMED_OUT_KEY], json!(false));
245        assert_eq!(received["payload"], signal.payload);
246        assert_eq!(received["signal_id"], json!(signal.id));
247        assert_eq!(timed_out_output(), json!({"timed_out": true}));
248    }
249}