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