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::{ApprovalRequirement, Assignee, FsmState, StepApproval, 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    /// Uncached input token count (agent steps only).
98    pub input_tokens: Option<u64>,
99    /// Input tokens served from the prompt cache (agent steps only).
100    #[serde(default)]
101    pub cache_read_input_tokens: Option<u64>,
102    /// Input tokens written to the prompt cache (agent steps only).
103    #[serde(default)]
104    pub cache_creation_input_tokens: Option<u64>,
105    /// Output token count (agent steps only).
106    pub output_tokens: Option<u64>,
107    /// When the step was created.
108    pub created_at: DateTime<Utc>,
109    /// When the step record was last updated.
110    pub updated_at: DateTime<Utc>,
111    /// When step execution started.
112    pub started_at: Option<DateTime<Utc>>,
113    /// When step execution finished.
114    pub completed_at: Option<DateTime<Utc>>,
115    /// Debug messages (verbose conversation trace), stored as JSON.
116    pub debug_messages: Option<Value>,
117    /// Whether this step is an error handler (`on_error`) rather than a normal step.
118    #[serde(default)]
119    pub is_error_handler: bool,
120    /// When the approval gate on this step expires, if it carries an SLA deadline.
121    ///
122    /// Only ever set while the step is [`StepStatus::AwaitingApproval`]. Cleared
123    /// the moment the gate resolves (approved, rejected, or escalated).
124    #[serde(default)]
125    pub approval_deadline_at: Option<DateTime<Utc>>,
126    /// Index of the next escalation policy to apply when the deadline fires.
127    #[serde(default)]
128    pub approval_stage: u32,
129    /// User or group the approval is currently assigned to, after reassignment.
130    #[serde(default)]
131    pub approval_assignee: Option<Assignee>,
132    /// Approvers the handler required when the gate opened. `None` for a gate
133    /// opened without approvers: one approval resolves it.
134    #[serde(default)]
135    pub approval_requirement: Option<ApprovalRequirement>,
136    /// Votes cast on the approval gate so far, at most one per user.
137    #[serde(default)]
138    pub approvals: Vec<StepApproval>,
139}
140
141/// Request to create a new step.
142///
143/// # Examples
144///
145/// ```
146/// use ironflow_store::entities::{NewStep, StepKind, step_trace_id};
147/// use serde_json::json;
148/// use uuid::Uuid;
149///
150/// let run_id = Uuid::nil();
151/// let req = NewStep {
152///     run_id,
153///     trace_id: step_trace_id(run_id, "build", 0),
154///     name: "build".to_string(),
155///     kind: StepKind::Shell,
156///     position: 0,
157///     input: Some(json!({"command": "cargo build"})),
158///     is_error_handler: false,
159/// };
160/// ```
161#[derive(Debug, Clone, Serialize, Deserialize)]
162pub struct NewStep {
163    /// The run this step belongs to.
164    pub run_id: Uuid,
165    /// Deterministic trace ID for log correlation.
166    pub trace_id: Uuid,
167    /// Step name.
168    pub name: String,
169    /// Operation type.
170    pub kind: StepKind,
171    /// Execution order (0-based).
172    pub position: u32,
173    /// Serialized operation configuration.
174    pub input: Option<Value>,
175    /// Whether this step is an error handler (`on_error`).
176    #[serde(default)]
177    pub is_error_handler: bool,
178}
179
180/// Partial update for a step after execution.
181///
182/// Only `Some` fields are applied; `None` fields are left unchanged.
183///
184/// # Examples
185///
186/// ```
187/// use ironflow_store::entities::{StepUpdate, StepStatus};
188/// use serde_json::json;
189///
190/// let update = StepUpdate {
191///     status: Some(StepStatus::Completed),
192///     output: Some(json!({"stdout": "ok"})),
193///     ..StepUpdate::default()
194/// };
195/// ```
196#[derive(Debug, Clone, Default, Serialize, Deserialize)]
197pub struct StepUpdate {
198    /// New status.
199    pub status: Option<StepStatus>,
200    /// Operation output.
201    pub output: Option<Value>,
202    /// Error message.
203    pub error: Option<String>,
204    /// Execution duration.
205    pub duration_ms: Option<u64>,
206    /// Cost in USD.
207    pub cost_usd: Option<Decimal>,
208    /// Uncached input token count.
209    pub input_tokens: Option<u64>,
210    /// Input tokens served from the prompt cache.
211    #[serde(default)]
212    pub cache_read_input_tokens: Option<u64>,
213    /// Input tokens written to the prompt cache.
214    #[serde(default)]
215    pub cache_creation_input_tokens: Option<u64>,
216    /// Output token count.
217    pub output_tokens: Option<u64>,
218    /// When execution started.
219    pub started_at: Option<DateTime<Utc>>,
220    /// When execution completed.
221    pub completed_at: Option<DateTime<Utc>>,
222    /// Debug messages (verbose conversation trace), stored as JSON.
223    pub debug_messages: Option<Value>,
224    /// New approval deadline. `None` leaves it unchanged — use
225    /// `clear_approval_deadline` to remove it.
226    #[serde(default)]
227    pub approval_deadline_at: Option<DateTime<Utc>>,
228    /// New escalation stage index.
229    #[serde(default)]
230    pub approval_stage: Option<u32>,
231    /// New approval assignee.
232    #[serde(default)]
233    pub approval_assignee: Option<Assignee>,
234    /// Approval requirement recorded when the gate opened.
235    #[serde(default)]
236    pub approval_requirement: Option<ApprovalRequirement>,
237    /// Clear the approval deadline (sets it to `NULL`). Wins over
238    /// `approval_deadline_at` when both are set.
239    #[serde(default)]
240    pub clear_approval_deadline: bool,
241}
242
243#[cfg(test)]
244mod tests {
245    use super::*;
246    use serde_json::json;
247
248    #[test]
249    fn newstep_serde_roundtrip() {
250        let new_step = NewStep {
251            run_id: Uuid::nil(),
252            trace_id: step_trace_id(Uuid::nil(), "build", 0),
253            name: "build".to_string(),
254            kind: StepKind::Shell,
255            position: 0,
256            input: Some(json!({"command": "cargo build"})),
257            is_error_handler: false,
258        };
259
260        let json = serde_json::to_string(&new_step).expect("serialize");
261        let back: NewStep = serde_json::from_str(&json).expect("deserialize");
262
263        assert_eq!(back.run_id, new_step.run_id);
264        assert_eq!(back.name, new_step.name);
265        assert_eq!(back.kind, new_step.kind);
266        assert_eq!(back.position, new_step.position);
267        assert_eq!(back.input, new_step.input);
268    }
269
270    #[test]
271    fn step_serde_preserves_all_fields() {
272        use crate::entities::FsmState;
273        use chrono::Utc;
274
275        let now = Utc::now();
276        let run_id = Uuid::now_v7();
277        let step = Step {
278            id: Uuid::now_v7(),
279            trace_id: step_trace_id(run_id, "test-step", 1),
280            run_id,
281            name: "test-step".to_string(),
282            kind: StepKind::Agent,
283            position: 1,
284            status: FsmState::new(StepStatus::Completed, Uuid::now_v7()),
285            attempt: 2,
286            input: Some(json!({"input": "data"})),
287            output: Some(json!({"output": "result"})),
288            error: None,
289            duration_ms: 2500,
290            cost_usd: Decimal::new(150, 2),
291            input_tokens: Some(100),
292            cache_read_input_tokens: Some(4000),
293            cache_creation_input_tokens: Some(300),
294            output_tokens: Some(200),
295            created_at: now,
296            updated_at: now,
297            started_at: Some(now),
298            completed_at: Some(now),
299            debug_messages: None,
300            is_error_handler: false,
301            approval_deadline_at: Some(now),
302            approval_stage: 2,
303            approval_assignee: Some(Assignee::group("sre-oncall")),
304            approval_requirement: Some(ApprovalRequirement {
305                reason: Some("amount > 10k".to_string()),
306                required_approvers: 2,
307                approver_groups: vec!["finance".to_string()],
308            }),
309            approvals: vec![StepApproval {
310                user_id: Uuid::now_v7(),
311                approved_by: "alice".to_string(),
312                at: now,
313            }],
314        };
315
316        let json = serde_json::to_string(&step).expect("serialize");
317        let back: Step = serde_json::from_str(&json).expect("deserialize");
318
319        assert_eq!(back.id, step.id);
320        assert_eq!(back.run_id, step.run_id);
321        assert_eq!(back.name, step.name);
322        assert_eq!(back.kind, step.kind);
323        assert_eq!(back.position, step.position);
324        assert_eq!(back.status.state, step.status.state);
325        assert_eq!(back.attempt, step.attempt);
326        assert_eq!(back.input, step.input);
327        assert_eq!(back.output, step.output);
328        assert_eq!(back.error, step.error);
329        assert_eq!(back.duration_ms, step.duration_ms);
330        assert_eq!(back.cost_usd, step.cost_usd);
331        assert_eq!(back.input_tokens, step.input_tokens);
332        assert_eq!(back.cache_read_input_tokens, step.cache_read_input_tokens);
333        assert_eq!(
334            back.cache_creation_input_tokens,
335            step.cache_creation_input_tokens
336        );
337        assert_eq!(back.output_tokens, step.output_tokens);
338        assert_eq!(back.approval_deadline_at, step.approval_deadline_at);
339        assert_eq!(back.approval_stage, step.approval_stage);
340        assert_eq!(back.approval_assignee, step.approval_assignee);
341        assert_eq!(back.approval_requirement, step.approval_requirement);
342        assert_eq!(back.approvals, step.approvals);
343    }
344
345    #[test]
346    fn step_serde_defaults_approval_fields_when_absent() {
347        let run_id = Uuid::now_v7();
348        let payload = json!({
349            "id": Uuid::now_v7(),
350            "trace_id": step_trace_id(run_id, "legacy", 0),
351            "run_id": run_id,
352            "name": "legacy",
353            "kind": "shell",
354            "position": 0,
355            "status": {"state": "pending", "state_machine_id": Uuid::now_v7()},
356            "input": null,
357            "output": null,
358            "error": null,
359            "duration_ms": 0,
360            "cost_usd": 0.0,
361            "input_tokens": null,
362            "cache_read_input_tokens": null,
363            "cache_creation_input_tokens": null,
364            "output_tokens": null,
365            "created_at": "2026-09-21T12:00:00Z",
366            "updated_at": "2026-09-21T12:00:00Z",
367            "started_at": null,
368            "completed_at": null,
369            "debug_messages": null
370        });
371
372        let step: Step = serde_json::from_value(payload).expect("deserialize");
373
374        assert_eq!(step.approval_stage, 0);
375        assert!(step.cache_read_input_tokens.is_none());
376        assert!(step.cache_creation_input_tokens.is_none());
377        assert!(step.approval_deadline_at.is_none());
378        assert!(step.approval_assignee.is_none());
379        assert!(step.approval_requirement.is_none());
380        assert!(step.approvals.is_empty());
381    }
382
383    #[test]
384    fn stepupdate_default_is_no_changes() {
385        let update = StepUpdate::default();
386        assert!(update.status.is_none());
387        assert!(update.output.is_none());
388        assert!(update.error.is_none());
389        assert!(update.duration_ms.is_none());
390        assert!(update.cost_usd.is_none());
391        assert!(update.input_tokens.is_none());
392        assert!(update.cache_read_input_tokens.is_none());
393        assert!(update.cache_creation_input_tokens.is_none());
394        assert!(update.output_tokens.is_none());
395        assert!(update.started_at.is_none());
396        assert!(update.completed_at.is_none());
397        assert!(update.debug_messages.is_none());
398        assert!(update.approval_deadline_at.is_none());
399        assert!(update.approval_stage.is_none());
400        assert!(update.approval_assignee.is_none());
401        assert!(update.approval_requirement.is_none());
402        assert!(!update.clear_approval_deadline);
403    }
404
405    #[test]
406    fn stepupdate_serde_roundtrip() {
407        let update = StepUpdate {
408            status: Some(StepStatus::Completed),
409            output: Some(json!({"result": "ok"})),
410            error: None,
411            duration_ms: Some(1000),
412            cost_usd: Some(Decimal::new(50, 2)),
413            input_tokens: Some(50),
414            cache_read_input_tokens: Some(2000),
415            cache_creation_input_tokens: Some(150),
416            output_tokens: Some(75),
417            started_at: None,
418            completed_at: None,
419            debug_messages: None,
420            approval_deadline_at: Some(Utc::now()),
421            approval_stage: Some(1),
422            approval_assignee: Some(Assignee::group("sre-oncall")),
423            approval_requirement: Some(ApprovalRequirement::default()),
424            clear_approval_deadline: false,
425        };
426
427        let json = serde_json::to_string(&update).expect("serialize");
428        let back: StepUpdate = serde_json::from_str(&json).expect("deserialize");
429
430        assert_eq!(back.status, update.status);
431        assert_eq!(back.output, update.output);
432        assert_eq!(back.duration_ms, update.duration_ms);
433        assert_eq!(back.cost_usd, update.cost_usd);
434        assert_eq!(back.input_tokens, update.input_tokens);
435        assert_eq!(back.cache_read_input_tokens, update.cache_read_input_tokens);
436        assert_eq!(
437            back.cache_creation_input_tokens,
438            update.cache_creation_input_tokens
439        );
440        assert_eq!(back.output_tokens, update.output_tokens);
441        assert_eq!(back.approval_deadline_at, update.approval_deadline_at);
442        assert_eq!(back.approval_stage, update.approval_stage);
443        assert_eq!(back.approval_assignee, update.approval_assignee);
444        assert_eq!(back.approval_requirement, update.approval_requirement);
445        assert_eq!(back.clear_approval_deadline, update.clear_approval_deadline);
446    }
447
448    #[test]
449    fn stepupdate_deserializes_without_cache_fields() {
450        let payload = json!({
451            "status": "completed",
452            "output": null,
453            "error": null,
454            "duration_ms": 10,
455            "cost_usd": null,
456            "input_tokens": 12,
457            "output_tokens": 3,
458            "started_at": null,
459            "completed_at": null,
460            "debug_messages": null
461        });
462
463        let update: StepUpdate = serde_json::from_value(payload).expect("deserialize");
464
465        assert_eq!(update.input_tokens, Some(12));
466        assert!(update.cache_read_input_tokens.is_none());
467        assert!(update.cache_creation_input_tokens.is_none());
468    }
469
470    #[test]
471    fn trace_id_is_deterministic() {
472        let run_id = Uuid::nil();
473        let id1 = step_trace_id(run_id, "build", 0);
474        let id2 = step_trace_id(run_id, "build", 0);
475        assert_eq!(id1, id2);
476    }
477
478    #[test]
479    fn trace_id_differs_for_different_inputs() {
480        let run_id = Uuid::nil();
481        let a = step_trace_id(run_id, "build", 0);
482        let b = step_trace_id(run_id, "test", 0);
483        let c = step_trace_id(run_id, "build", 1);
484        let d = step_trace_id(Uuid::max(), "build", 0);
485
486        assert_ne!(a, b);
487        assert_ne!(a, c);
488        assert_ne!(a, d);
489    }
490
491    #[test]
492    fn trace_id_is_uuid_v5() {
493        let id = step_trace_id(Uuid::nil(), "build", 0);
494        assert_eq!(id.get_version_num(), 5);
495    }
496}