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