Skip to main content

ironflow_api/entities/
step.rs

1//! Step-related DTOs.
2
3use chrono::{DateTime, Utc};
4use ironflow_store::models::{
5    ApprovalRequirement, Assignee, Step, StepApproval, StepKind, StepStatus,
6};
7use rust_decimal::Decimal;
8use serde::{Deserialize, Serialize};
9use serde_json::Value;
10use uuid::Uuid;
11
12use super::ArtifactResponse;
13
14/// Step response DTO — public API representation of a step.
15///
16/// # Examples
17///
18/// ```
19/// use ironflow_store::models::Step;
20/// use ironflow_api::entities::StepResponse;
21/// ```
22#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
23#[derive(Debug, Serialize, Deserialize)]
24pub struct StepResponse {
25    /// Unique step identifier.
26    pub id: Uuid,
27    /// Deterministic trace ID for log correlation.
28    pub trace_id: Uuid,
29    /// Parent run ID.
30    pub run_id: Uuid,
31    /// Step name.
32    pub name: String,
33    /// Step operation type.
34    #[cfg_attr(feature = "openapi", schema(value_type = String))]
35    pub kind: StepKind,
36    /// Execution order (0-based).
37    pub position: u32,
38    /// Current status.
39    pub status: StepStatus,
40    /// Which run attempt produced this step (1-based).
41    ///
42    /// A run retried twice exposes steps with `attempt` 1, 2 and 3. Steps from
43    /// earlier attempts are kept so a failed attempt stays inspectable.
44    pub attempt: u32,
45    /// Input configuration.
46    #[cfg_attr(feature = "openapi", schema(value_type = Option<std::collections::HashMap<String, serde_json::Value>>))]
47    pub input: Option<Value>,
48    /// Step output.
49    #[cfg_attr(feature = "openapi", schema(value_type = Option<std::collections::HashMap<String, serde_json::Value>>))]
50    pub output: Option<Value>,
51    /// Optional error message.
52    pub error: Option<String>,
53    /// Execution duration in milliseconds.
54    pub duration_ms: u64,
55    /// Cost in USD.
56    #[cfg_attr(feature = "openapi", schema(value_type = f64))]
57    pub cost_usd: Decimal,
58    /// Input token count (agent steps).
59    pub input_tokens: Option<u64>,
60    /// Output token count (agent steps).
61    pub output_tokens: Option<u64>,
62    /// When created.
63    pub created_at: DateTime<Utc>,
64    /// When updated.
65    pub updated_at: DateTime<Utc>,
66    /// When execution started.
67    pub started_at: Option<DateTime<Utc>>,
68    /// When execution completed.
69    pub completed_at: Option<DateTime<Utc>>,
70    /// IDs of steps this step depends on (direct dependencies).
71    pub dependencies: Vec<Uuid>,
72    /// Verbose conversation trace for agent steps (thinking blocks, tool
73    /// calls, tool results, per-turn usage). `None` when verbose mode was
74    /// off or the step is not an agent step.
75    #[cfg_attr(feature = "openapi", schema(value_type = Option<serde_json::Value>))]
76    pub debug_messages: Option<Value>,
77    /// Files this step produced, downloadable through the artifact route.
78    ///
79    /// Empty when the step produced none or when artifact storage is not
80    /// configured on the server.
81    #[serde(default)]
82    pub artifacts: Vec<ArtifactResponse>,
83    /// When this approval gate expires, if it carries an SLA deadline.
84    pub approval_deadline_at: Option<DateTime<Utc>>,
85    /// Seconds left before the gate escalates. Clamped at 0, `None` when the
86    /// step has no deadline.
87    pub approval_seconds_remaining: Option<i64>,
88    /// Who the approval is currently assigned to.
89    ///
90    /// Serialized as a prefixed string: `user:{name}` or `group:{name}`.
91    #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
92    pub approval_assignee: Option<Assignee>,
93    /// Approval requirement evaluated from the gate's approval rules when it
94    /// opened. `None` for steps without rules.
95    #[serde(default)]
96    pub approval_requirement: Option<ApprovalRequirement>,
97    /// Votes cast on the approval gate so far, at most one per user.
98    #[serde(default)]
99    pub approvals: Vec<StepApproval>,
100    /// Distinct approvals the gate needs: the requirement's count, `1` for an
101    /// approval step without rules, `None` for any other step kind.
102    #[serde(default)]
103    pub approvals_required: Option<u32>,
104}
105
106impl StepResponse {
107    /// Build a response from a step entity with pre-resolved dependencies.
108    ///
109    /// Artifacts are left empty; use
110    /// [`with_dependencies_and_artifacts`](Self::with_dependencies_and_artifacts)
111    /// when they have been fetched.
112    pub fn with_dependencies(step: Step, dependencies: Vec<Uuid>) -> Self {
113        Self::with_dependencies_and_artifacts(step, dependencies, Vec::new())
114    }
115
116    /// Build a response from a step entity with its dependencies and artifacts.
117    pub fn with_dependencies_and_artifacts(
118        step: Step,
119        dependencies: Vec<Uuid>,
120        artifacts: Vec<ArtifactResponse>,
121    ) -> Self {
122        let approval_seconds_remaining = step
123            .approval_deadline_at
124            .map(|at| (at - Utc::now()).num_seconds().max(0));
125        let approvals_required = match (&step.kind, &step.approval_requirement) {
126            (_, Some(requirement)) => Some(requirement.required_approvers),
127            (StepKind::Approval, None) => Some(1),
128            _ => None,
129        };
130
131        StepResponse {
132            id: step.id,
133            trace_id: step.trace_id,
134            run_id: step.run_id,
135            name: step.name,
136            kind: step.kind,
137            position: step.position,
138            status: step.status.state,
139            attempt: step.attempt,
140            input: step.input,
141            output: step.output,
142            error: step.error,
143            duration_ms: step.duration_ms,
144            cost_usd: step.cost_usd,
145            input_tokens: step.input_tokens,
146            output_tokens: step.output_tokens,
147            created_at: step.created_at,
148            updated_at: step.updated_at,
149            started_at: step.started_at,
150            completed_at: step.completed_at,
151            dependencies,
152            debug_messages: step.debug_messages,
153            artifacts,
154            approval_deadline_at: step.approval_deadline_at,
155            approval_seconds_remaining,
156            approval_assignee: step.approval_assignee,
157            approval_requirement: step.approval_requirement,
158            approvals: step.approvals,
159            approvals_required,
160        }
161    }
162}
163
164impl From<Step> for StepResponse {
165    fn from(step: Step) -> Self {
166        Self::with_dependencies(step, Vec::new())
167    }
168}
169
170#[cfg(test)]
171mod tests {
172    use std::collections::HashMap;
173
174    use chrono::TimeDelta;
175    use ironflow_store::memory::InMemoryStore;
176    use ironflow_store::models::{NewRun, NewStep, StepUpdate, TriggerKind, step_trace_id};
177    use ironflow_store::store::RunStore;
178    use serde_json::json;
179
180    use super::*;
181
182    /// A persisted step -- [`Step`] is `#[non_exhaustive]`, so it can only be
183    /// obtained from a store.
184    async fn step() -> Step {
185        let store = InMemoryStore::new();
186        let run = store
187            .create_run(NewRun {
188                created_by: None,
189                workflow_name: "test".to_string(),
190                trigger: TriggerKind::Manual,
191                payload: json!({}),
192                max_retries: 0,
193                handler_version: None,
194                labels: HashMap::new(),
195                scheduled_at: None,
196                idempotency_key: None,
197                max_cost_usd: None,
198            })
199            .await
200            .expect("create run")
201            .into_run();
202
203        store
204            .create_step(NewStep {
205                run_id: run.id,
206                trace_id: step_trace_id(run.id, "build", 0),
207                name: "build".to_string(),
208                kind: StepKind::Shell,
209                position: 0,
210                input: None,
211                is_error_handler: false,
212            })
213            .await
214            .expect("create step")
215    }
216
217    #[tokio::test]
218    async fn a_step_without_artifacts_exposes_an_empty_list() {
219        let response = StepResponse::from(step().await);
220        assert!(response.artifacts.is_empty());
221    }
222
223    #[tokio::test]
224    async fn artifacts_are_carried_through() {
225        let step = step().await;
226        let artifact = ArtifactResponse {
227            id: Uuid::now_v7(),
228            step_id: step.id,
229            name: "report.html".to_string(),
230            content_type: "text/html".to_string(),
231            size_bytes: 1,
232            sha256: "0".repeat(64),
233            created_at: Utc::now(),
234        };
235
236        let response =
237            StepResponse::with_dependencies_and_artifacts(step, Vec::new(), vec![artifact]);
238
239        assert_eq!(response.artifacts.len(), 1);
240        assert_eq!(response.artifacts[0].name, "report.html");
241    }
242
243    #[tokio::test]
244    async fn artifacts_serialize_as_a_json_array() {
245        let body = serde_json::to_value(StepResponse::from(step().await)).expect("serialize");
246        assert!(body["artifacts"].is_array());
247    }
248
249    /// An approval step whose deadline is `offset_secs` from now.
250    async fn gate_with_deadline(offset_secs: i64) -> Step {
251        let store = InMemoryStore::new();
252        let run = store
253            .create_run(NewRun {
254                created_by: None,
255                workflow_name: "test".to_string(),
256                trigger: TriggerKind::Manual,
257                payload: json!({}),
258                max_retries: 0,
259                handler_version: None,
260                labels: HashMap::new(),
261                scheduled_at: None,
262                idempotency_key: None,
263                max_cost_usd: None,
264            })
265            .await
266            .expect("create run")
267            .into_run();
268
269        let step = store
270            .create_step(NewStep {
271                run_id: run.id,
272                trace_id: step_trace_id(run.id, "prod-gate", 0),
273                name: "prod-gate".to_string(),
274                kind: StepKind::Approval,
275                position: 0,
276                input: None,
277                is_error_handler: false,
278            })
279            .await
280            .expect("create step");
281
282        store
283            .update_step(
284                step.id,
285                StepUpdate {
286                    status: Some(StepStatus::Running),
287                    ..StepUpdate::default()
288                },
289            )
290            .await
291            .expect("to running");
292        store
293            .update_step(
294                step.id,
295                StepUpdate {
296                    status: Some(StepStatus::AwaitingApproval),
297                    approval_deadline_at: Some(Utc::now() + TimeDelta::seconds(offset_secs)),
298                    approval_assignee: Some(Assignee::group("release-managers")),
299                    ..StepUpdate::default()
300                },
301            )
302            .await
303            .expect("arm timer");
304
305        store.get_step(step.id).await.expect("get").expect("exists")
306    }
307
308    #[tokio::test]
309    async fn a_step_without_a_deadline_reports_no_sla() {
310        let response = StepResponse::from(step().await);
311        assert!(response.approval_deadline_at.is_none());
312        assert!(response.approval_seconds_remaining.is_none());
313        assert!(response.approval_assignee.is_none());
314        assert!(response.approval_requirement.is_none());
315        assert!(response.approvals.is_empty());
316    }
317
318    #[tokio::test]
319    async fn a_shell_step_requires_no_approvals() {
320        let response = StepResponse::from(step().await);
321        assert_eq!(response.approvals_required, None);
322    }
323
324    #[tokio::test]
325    async fn a_rule_less_approval_step_requires_one_approval() {
326        let response = StepResponse::from(gate_with_deadline(3600).await);
327        assert!(response.approval_requirement.is_none());
328        assert_eq!(response.approvals_required, Some(1));
329    }
330
331    #[tokio::test]
332    async fn an_approval_requirement_sets_the_required_count() {
333        let mut gate = gate_with_deadline(3600).await;
334        let requirement = ApprovalRequirement {
335            rule_index: Some(0),
336            condition: Some("payload.amount > 10000".to_string()),
337            required_approvers: 3,
338            approver_groups: vec!["finance".to_string()],
339            evaluated: Vec::new(),
340        };
341        gate.approval_requirement = Some(requirement.clone());
342        gate.approvals = vec![StepApproval {
343            user_id: Uuid::now_v7(),
344            approved_by: "alice".to_string(),
345            at: Utc::now(),
346        }];
347
348        let response = StepResponse::from(gate);
349
350        assert_eq!(response.approvals_required, Some(3));
351        assert_eq!(response.approval_requirement, Some(requirement));
352        assert_eq!(response.approvals.len(), 1);
353        assert_eq!(response.approvals[0].approved_by, "alice");
354    }
355
356    #[tokio::test]
357    async fn a_future_deadline_reports_the_remaining_seconds() {
358        let response = StepResponse::from(gate_with_deadline(3600).await);
359
360        assert!(response.approval_deadline_at.is_some());
361        let remaining = response
362            .approval_seconds_remaining
363            .expect("a deadline yields a countdown");
364        assert!(remaining > 0 && remaining <= 3600, "got {remaining}");
365        assert_eq!(
366            response.approval_assignee,
367            Some(Assignee::group("release-managers"))
368        );
369    }
370
371    #[tokio::test]
372    async fn a_past_deadline_clamps_the_countdown_at_zero() {
373        let response = StepResponse::from(gate_with_deadline(-3600).await);
374        assert_eq!(response.approval_seconds_remaining, Some(0));
375    }
376
377    #[tokio::test]
378    async fn trace_id_is_exposed_in_step_response() {
379        let s = step().await;
380        let expected_trace_id = s.trace_id;
381        let response = StepResponse::from(s);
382
383        assert_eq!(response.trace_id, expected_trace_id);
384
385        let body = serde_json::to_value(&response).expect("serialize");
386        assert!(body["trace_id"].is_string());
387    }
388}