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}