Skip to main content

ironflow_store/entities/
step.rs

1//! [`Step`] entity and related request/update types.
2
3use chrono::{DateTime, Utc};
4use rust_decimal::Decimal;
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7use uuid::Uuid;
8
9use super::{Assignee, FsmState, StepKind, StepStatus};
10
11/// Attempt number assigned to steps deserialized from payloads predating the
12/// `attempt` field.
13fn default_attempt() -> u32 {
14    1
15}
16
17/// Generate a deterministic trace ID for a step.
18///
19/// Uses UUIDv5 with `NAMESPACE_OID` and the input
20/// `"{run_id}:{name}:{position}"`, so the same run replayed with the same
21/// steps always produces the same trace IDs.
22///
23/// # Examples
24///
25/// ```
26/// use ironflow_store::entities::step_trace_id;
27/// use uuid::Uuid;
28///
29/// let run_id = Uuid::nil();
30/// let id1 = step_trace_id(run_id, "build", 0);
31/// let id2 = step_trace_id(run_id, "build", 0);
32/// assert_eq!(id1, id2);
33///
34/// let id3 = step_trace_id(run_id, "test", 1);
35/// assert_ne!(id1, id3);
36/// ```
37pub fn step_trace_id(run_id: Uuid, name: &str, position: u32) -> Uuid {
38    Uuid::new_v5(
39        &Uuid::NAMESPACE_OID,
40        format!("{run_id}:{name}:{position}").as_bytes(),
41    )
42}
43
44/// A single operation within a run.
45///
46/// Steps are executed sequentially in order of [`position`](Step::position).
47///
48/// # Examples
49///
50/// ```
51/// use ironflow_store::entities::Step;
52///
53/// // Steps are created by RunStore::create_step, not directly.
54/// ```
55#[derive(Debug, Clone, Serialize, Deserialize)]
56#[non_exhaustive]
57pub struct Step {
58    /// Unique identifier (UUIDv7).
59    pub id: Uuid,
60    /// Deterministic trace ID for log correlation.
61    ///
62    /// Generated as `UUIDv5(NAMESPACE_OID, "{run_id}:{name}:{position}")` so
63    /// the same run replayed with the same steps always produces the same IDs.
64    pub trace_id: Uuid,
65    /// The run this step belongs to.
66    pub run_id: Uuid,
67    /// Human-readable step name (e.g. "build", "test", "review").
68    pub name: String,
69    /// The type of operation.
70    pub kind: StepKind,
71    /// Execution wave within the run (0-based).
72    ///
73    /// In linear flows, this strictly increases (0, 1, 2, ...).
74    /// In DAGs with parallel execution, steps at the same wave share
75    /// the same position and execute concurrently. Use
76    /// `step_dependencies` to determine the actual execution order.
77    pub position: u32,
78    /// Current FSM status — embeds state + state_machine_id for SQL-side transitions.
79    pub status: FsmState<StepStatus>,
80    /// Which run attempt produced this step (1-based).
81    ///
82    /// A run retried twice holds steps with `attempt` 1, 2 and 3. Steps from
83    /// earlier attempts are kept for inspection and are never replayed.
84    /// Derived by the store from `Run::retry_count` at creation time.
85    #[serde(default = "default_attempt")]
86    pub attempt: u32,
87    /// Serialized operation configuration.
88    pub input: Option<Value>,
89    /// Serialized operation output.
90    pub output: Option<Value>,
91    /// Error message if the step failed.
92    pub error: Option<String>,
93    /// Wall-clock execution duration in milliseconds.
94    pub duration_ms: u64,
95    /// Cost in USD (agent steps only).
96    pub cost_usd: Decimal,
97    /// Input token count (agent steps only).
98    pub input_tokens: Option<u64>,
99    /// Output token count (agent steps only).
100    pub output_tokens: Option<u64>,
101    /// When the step was created.
102    pub created_at: DateTime<Utc>,
103    /// When the step record was last updated.
104    pub updated_at: DateTime<Utc>,
105    /// When step execution started.
106    pub started_at: Option<DateTime<Utc>>,
107    /// When step execution finished.
108    pub completed_at: Option<DateTime<Utc>>,
109    /// Debug messages (verbose conversation trace), stored as JSON.
110    pub debug_messages: Option<Value>,
111    /// Whether this step is an error handler (`on_error`) rather than a normal step.
112    #[serde(default)]
113    pub is_error_handler: bool,
114    /// When the approval gate on this step expires, if it carries an SLA deadline.
115    ///
116    /// Only ever set while the step is [`StepStatus::AwaitingApproval`]. Cleared
117    /// the moment the gate resolves (approved, rejected, or escalated).
118    #[serde(default)]
119    pub approval_deadline_at: Option<DateTime<Utc>>,
120    /// Index of the next escalation policy to apply when the deadline fires.
121    #[serde(default)]
122    pub approval_stage: u32,
123    /// User or group the approval is currently assigned to, after reassignment.
124    #[serde(default)]
125    pub approval_assignee: Option<Assignee>,
126}
127
128/// Request to create a new step.
129///
130/// # Examples
131///
132/// ```
133/// use ironflow_store::entities::{NewStep, StepKind, step_trace_id};
134/// use serde_json::json;
135/// use uuid::Uuid;
136///
137/// let run_id = Uuid::nil();
138/// let req = NewStep {
139///     run_id,
140///     trace_id: step_trace_id(run_id, "build", 0),
141///     name: "build".to_string(),
142///     kind: StepKind::Shell,
143///     position: 0,
144///     input: Some(json!({"command": "cargo build"})),
145///     is_error_handler: false,
146/// };
147/// ```
148#[derive(Debug, Clone, Serialize, Deserialize)]
149pub struct NewStep {
150    /// The run this step belongs to.
151    pub run_id: Uuid,
152    /// Deterministic trace ID for log correlation.
153    pub trace_id: Uuid,
154    /// Step name.
155    pub name: String,
156    /// Operation type.
157    pub kind: StepKind,
158    /// Execution order (0-based).
159    pub position: u32,
160    /// Serialized operation configuration.
161    pub input: Option<Value>,
162    /// Whether this step is an error handler (`on_error`).
163    #[serde(default)]
164    pub is_error_handler: bool,
165}
166
167/// Partial update for a step after execution.
168///
169/// Only `Some` fields are applied; `None` fields are left unchanged.
170///
171/// # Examples
172///
173/// ```
174/// use ironflow_store::entities::{StepUpdate, StepStatus};
175/// use serde_json::json;
176///
177/// let update = StepUpdate {
178///     status: Some(StepStatus::Completed),
179///     output: Some(json!({"stdout": "ok"})),
180///     ..StepUpdate::default()
181/// };
182/// ```
183#[derive(Debug, Clone, Default, Serialize, Deserialize)]
184pub struct StepUpdate {
185    /// New status.
186    pub status: Option<StepStatus>,
187    /// Operation output.
188    pub output: Option<Value>,
189    /// Error message.
190    pub error: Option<String>,
191    /// Execution duration.
192    pub duration_ms: Option<u64>,
193    /// Cost in USD.
194    pub cost_usd: Option<Decimal>,
195    /// Input token count.
196    pub input_tokens: Option<u64>,
197    /// Output token count.
198    pub output_tokens: Option<u64>,
199    /// When execution started.
200    pub started_at: Option<DateTime<Utc>>,
201    /// When execution completed.
202    pub completed_at: Option<DateTime<Utc>>,
203    /// Debug messages (verbose conversation trace), stored as JSON.
204    pub debug_messages: Option<Value>,
205    /// New approval deadline. `None` leaves it unchanged — use
206    /// `clear_approval_deadline` to remove it.
207    #[serde(default)]
208    pub approval_deadline_at: Option<DateTime<Utc>>,
209    /// New escalation stage index.
210    #[serde(default)]
211    pub approval_stage: Option<u32>,
212    /// New approval assignee.
213    #[serde(default)]
214    pub approval_assignee: Option<Assignee>,
215    /// Clear the approval deadline (sets it to `NULL`). Wins over
216    /// `approval_deadline_at` when both are set.
217    #[serde(default)]
218    pub clear_approval_deadline: bool,
219}
220
221#[cfg(test)]
222mod tests {
223    use super::*;
224    use serde_json::json;
225
226    #[test]
227    fn newstep_serde_roundtrip() {
228        let new_step = NewStep {
229            run_id: Uuid::nil(),
230            trace_id: step_trace_id(Uuid::nil(), "build", 0),
231            name: "build".to_string(),
232            kind: StepKind::Shell,
233            position: 0,
234            input: Some(json!({"command": "cargo build"})),
235            is_error_handler: false,
236        };
237
238        let json = serde_json::to_string(&new_step).expect("serialize");
239        let back: NewStep = serde_json::from_str(&json).expect("deserialize");
240
241        assert_eq!(back.run_id, new_step.run_id);
242        assert_eq!(back.name, new_step.name);
243        assert_eq!(back.kind, new_step.kind);
244        assert_eq!(back.position, new_step.position);
245        assert_eq!(back.input, new_step.input);
246    }
247
248    #[test]
249    fn step_serde_preserves_all_fields() {
250        use crate::entities::FsmState;
251        use chrono::Utc;
252
253        let now = Utc::now();
254        let run_id = Uuid::now_v7();
255        let step = Step {
256            id: Uuid::now_v7(),
257            trace_id: step_trace_id(run_id, "test-step", 1),
258            run_id,
259            name: "test-step".to_string(),
260            kind: StepKind::Agent,
261            position: 1,
262            status: FsmState::new(StepStatus::Completed, Uuid::now_v7()),
263            attempt: 2,
264            input: Some(json!({"input": "data"})),
265            output: Some(json!({"output": "result"})),
266            error: None,
267            duration_ms: 2500,
268            cost_usd: Decimal::new(150, 2),
269            input_tokens: Some(100),
270            output_tokens: Some(200),
271            created_at: now,
272            updated_at: now,
273            started_at: Some(now),
274            completed_at: Some(now),
275            debug_messages: None,
276            is_error_handler: false,
277            approval_deadline_at: Some(now),
278            approval_stage: 2,
279            approval_assignee: Some(Assignee::group("sre-oncall")),
280        };
281
282        let json = serde_json::to_string(&step).expect("serialize");
283        let back: Step = serde_json::from_str(&json).expect("deserialize");
284
285        assert_eq!(back.id, step.id);
286        assert_eq!(back.run_id, step.run_id);
287        assert_eq!(back.name, step.name);
288        assert_eq!(back.kind, step.kind);
289        assert_eq!(back.position, step.position);
290        assert_eq!(back.status.state, step.status.state);
291        assert_eq!(back.attempt, step.attempt);
292        assert_eq!(back.input, step.input);
293        assert_eq!(back.output, step.output);
294        assert_eq!(back.error, step.error);
295        assert_eq!(back.duration_ms, step.duration_ms);
296        assert_eq!(back.cost_usd, step.cost_usd);
297        assert_eq!(back.input_tokens, step.input_tokens);
298        assert_eq!(back.output_tokens, step.output_tokens);
299        assert_eq!(back.approval_deadline_at, step.approval_deadline_at);
300        assert_eq!(back.approval_stage, step.approval_stage);
301        assert_eq!(back.approval_assignee, step.approval_assignee);
302    }
303
304    #[test]
305    fn step_serde_defaults_approval_fields_when_absent() {
306        let run_id = Uuid::now_v7();
307        let payload = json!({
308            "id": Uuid::now_v7(),
309            "trace_id": step_trace_id(run_id, "legacy", 0),
310            "run_id": run_id,
311            "name": "legacy",
312            "kind": "shell",
313            "position": 0,
314            "status": {"state": "pending", "state_machine_id": Uuid::now_v7()},
315            "input": null,
316            "output": null,
317            "error": null,
318            "duration_ms": 0,
319            "cost_usd": 0.0,
320            "input_tokens": null,
321            "output_tokens": null,
322            "created_at": "2026-09-21T12:00:00Z",
323            "updated_at": "2026-09-21T12:00:00Z",
324            "started_at": null,
325            "completed_at": null,
326            "debug_messages": null
327        });
328
329        let step: Step = serde_json::from_value(payload).expect("deserialize");
330
331        assert_eq!(step.approval_stage, 0);
332        assert!(step.approval_deadline_at.is_none());
333        assert!(step.approval_assignee.is_none());
334    }
335
336    #[test]
337    fn stepupdate_default_is_no_changes() {
338        let update = StepUpdate::default();
339        assert!(update.status.is_none());
340        assert!(update.output.is_none());
341        assert!(update.error.is_none());
342        assert!(update.duration_ms.is_none());
343        assert!(update.cost_usd.is_none());
344        assert!(update.input_tokens.is_none());
345        assert!(update.output_tokens.is_none());
346        assert!(update.started_at.is_none());
347        assert!(update.completed_at.is_none());
348        assert!(update.debug_messages.is_none());
349        assert!(update.approval_deadline_at.is_none());
350        assert!(update.approval_stage.is_none());
351        assert!(update.approval_assignee.is_none());
352        assert!(!update.clear_approval_deadline);
353    }
354
355    #[test]
356    fn stepupdate_serde_roundtrip() {
357        let update = StepUpdate {
358            status: Some(StepStatus::Completed),
359            output: Some(json!({"result": "ok"})),
360            error: None,
361            duration_ms: Some(1000),
362            cost_usd: Some(Decimal::new(50, 2)),
363            input_tokens: Some(50),
364            output_tokens: Some(75),
365            started_at: None,
366            completed_at: None,
367            debug_messages: None,
368            approval_deadline_at: Some(Utc::now()),
369            approval_stage: Some(1),
370            approval_assignee: Some(Assignee::group("sre-oncall")),
371            clear_approval_deadline: false,
372        };
373
374        let json = serde_json::to_string(&update).expect("serialize");
375        let back: StepUpdate = serde_json::from_str(&json).expect("deserialize");
376
377        assert_eq!(back.status, update.status);
378        assert_eq!(back.output, update.output);
379        assert_eq!(back.duration_ms, update.duration_ms);
380        assert_eq!(back.cost_usd, update.cost_usd);
381        assert_eq!(back.input_tokens, update.input_tokens);
382        assert_eq!(back.output_tokens, update.output_tokens);
383        assert_eq!(back.approval_deadline_at, update.approval_deadline_at);
384        assert_eq!(back.approval_stage, update.approval_stage);
385        assert_eq!(back.approval_assignee, update.approval_assignee);
386        assert_eq!(back.clear_approval_deadline, update.clear_approval_deadline);
387    }
388
389    #[test]
390    fn trace_id_is_deterministic() {
391        let run_id = Uuid::nil();
392        let id1 = step_trace_id(run_id, "build", 0);
393        let id2 = step_trace_id(run_id, "build", 0);
394        assert_eq!(id1, id2);
395    }
396
397    #[test]
398    fn trace_id_differs_for_different_inputs() {
399        let run_id = Uuid::nil();
400        let a = step_trace_id(run_id, "build", 0);
401        let b = step_trace_id(run_id, "test", 0);
402        let c = step_trace_id(run_id, "build", 1);
403        let d = step_trace_id(Uuid::max(), "build", 0);
404
405        assert_ne!(a, b);
406        assert_ne!(a, c);
407        assert_ne!(a, d);
408    }
409
410    #[test]
411    fn trace_id_is_uuid_v5() {
412        let id = step_trace_id(Uuid::nil(), "build", 0);
413        assert_eq!(id.get_version_num(), 5);
414    }
415}