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            (_, Some(requirement)) => Some(requirement.required_approvers),
131            (StepKind::Approval, None) => Some(1),
132            _ => None,
133        };
134
135        StepResponse {
136            id: step.id,
137            trace_id: step.trace_id,
138            run_id: step.run_id,
139            name: step.name,
140            kind: step.kind,
141            position: step.position,
142            status: step.status.state,
143            attempt: step.attempt,
144            input: step.input,
145            output: step.output,
146            error: step.error,
147            duration_ms: step.duration_ms,
148            cost_usd: step.cost_usd,
149            input_tokens: step.input_tokens,
150            cache_read_input_tokens: step.cache_read_input_tokens,
151            cache_creation_input_tokens: step.cache_creation_input_tokens,
152            output_tokens: step.output_tokens,
153            created_at: step.created_at,
154            updated_at: step.updated_at,
155            started_at: step.started_at,
156            completed_at: step.completed_at,
157            dependencies,
158            debug_messages: step.debug_messages,
159            artifacts,
160            approval_deadline_at: step.approval_deadline_at,
161            approval_seconds_remaining,
162            approval_assignee: step.approval_assignee,
163            approval_requirement: step.approval_requirement,
164            approvals: step.approvals,
165            approvals_required,
166        }
167    }
168}
169
170impl From<Step> for StepResponse {
171    fn from(step: Step) -> Self {
172        Self::with_dependencies(step, Vec::new())
173    }
174}
175
176#[cfg(test)]
177mod tests {
178    use std::collections::HashMap;
179
180    use chrono::TimeDelta;
181    use ironflow_store::memory::InMemoryStore;
182    use ironflow_store::models::{NewRun, NewStep, StepUpdate, TriggerKind, step_trace_id};
183    use ironflow_store::store::RunStore;
184    use serde_json::json;
185
186    use super::*;
187
188    /// A persisted step -- [`Step`] is `#[non_exhaustive]`, so it can only be
189    /// obtained from a store.
190    async fn step() -> Step {
191        let store = InMemoryStore::new();
192        let run = store
193            .create_run(NewRun {
194                created_by: None,
195                workflow_name: "test".to_string(),
196                trigger: TriggerKind::Manual,
197                payload: json!({}),
198                max_retries: 0,
199                handler_version: None,
200                labels: HashMap::new(),
201                scheduled_at: None,
202                idempotency_key: None,
203                max_cost_usd: None,
204            })
205            .await
206            .expect("create run")
207            .into_run();
208
209        store
210            .create_step(NewStep {
211                run_id: run.id,
212                trace_id: step_trace_id(run.id, "build", 0),
213                name: "build".to_string(),
214                kind: StepKind::Shell,
215                position: 0,
216                input: None,
217                is_error_handler: false,
218            })
219            .await
220            .expect("create step")
221    }
222
223    #[tokio::test]
224    async fn a_step_without_artifacts_exposes_an_empty_list() {
225        let response = StepResponse::from(step().await);
226        assert!(response.artifacts.is_empty());
227    }
228
229    #[tokio::test]
230    async fn artifacts_are_carried_through() {
231        let step = step().await;
232        let artifact = ArtifactResponse {
233            id: Uuid::now_v7(),
234            step_id: step.id,
235            name: "report.html".to_string(),
236            content_type: "text/html".to_string(),
237            size_bytes: 1,
238            sha256: "0".repeat(64),
239            created_at: Utc::now(),
240        };
241
242        let response =
243            StepResponse::with_dependencies_and_artifacts(step, Vec::new(), vec![artifact]);
244
245        assert_eq!(response.artifacts.len(), 1);
246        assert_eq!(response.artifacts[0].name, "report.html");
247    }
248
249    #[tokio::test]
250    async fn artifacts_serialize_as_a_json_array() {
251        let body = serde_json::to_value(StepResponse::from(step().await)).expect("serialize");
252        assert!(body["artifacts"].is_array());
253    }
254
255    /// An approval step whose deadline is `offset_secs` from now.
256    async fn gate_with_deadline(offset_secs: i64) -> Step {
257        let store = InMemoryStore::new();
258        let run = store
259            .create_run(NewRun {
260                created_by: None,
261                workflow_name: "test".to_string(),
262                trigger: TriggerKind::Manual,
263                payload: json!({}),
264                max_retries: 0,
265                handler_version: None,
266                labels: HashMap::new(),
267                scheduled_at: None,
268                idempotency_key: None,
269                max_cost_usd: None,
270            })
271            .await
272            .expect("create run")
273            .into_run();
274
275        let step = store
276            .create_step(NewStep {
277                run_id: run.id,
278                trace_id: step_trace_id(run.id, "prod-gate", 0),
279                name: "prod-gate".to_string(),
280                kind: StepKind::Approval,
281                position: 0,
282                input: None,
283                is_error_handler: false,
284            })
285            .await
286            .expect("create step");
287
288        store
289            .update_step(
290                step.id,
291                StepUpdate {
292                    status: Some(StepStatus::Running),
293                    ..StepUpdate::default()
294                },
295            )
296            .await
297            .expect("to running");
298        store
299            .update_step(
300                step.id,
301                StepUpdate {
302                    status: Some(StepStatus::AwaitingApproval),
303                    approval_deadline_at: Some(Utc::now() + TimeDelta::seconds(offset_secs)),
304                    approval_assignee: Some(Assignee::group("release-managers")),
305                    ..StepUpdate::default()
306                },
307            )
308            .await
309            .expect("arm timer");
310
311        store.get_step(step.id).await.expect("get").expect("exists")
312    }
313
314    #[tokio::test]
315    async fn a_step_without_a_deadline_reports_no_sla() {
316        let response = StepResponse::from(step().await);
317        assert!(response.approval_deadline_at.is_none());
318        assert!(response.approval_seconds_remaining.is_none());
319        assert!(response.approval_assignee.is_none());
320        assert!(response.approval_requirement.is_none());
321        assert!(response.approvals.is_empty());
322    }
323
324    #[tokio::test]
325    async fn a_shell_step_requires_no_approvals() {
326        let response = StepResponse::from(step().await);
327        assert_eq!(response.approvals_required, None);
328    }
329
330    #[tokio::test]
331    async fn a_rule_less_approval_step_requires_one_approval() {
332        let response = StepResponse::from(gate_with_deadline(3600).await);
333        assert!(response.approval_requirement.is_none());
334        assert_eq!(response.approvals_required, Some(1));
335    }
336
337    #[tokio::test]
338    async fn an_approval_requirement_sets_the_required_count() {
339        let mut gate = gate_with_deadline(3600).await;
340        let requirement = ApprovalRequirement {
341            reason: Some("amount > 100k".to_string()),
342            required_approvers: 3,
343            approver_groups: vec!["finance".to_string()],
344        };
345        gate.approval_requirement = Some(requirement.clone());
346        gate.approvals = vec![StepApproval {
347            user_id: Uuid::now_v7(),
348            approved_by: "alice".to_string(),
349            at: Utc::now(),
350        }];
351
352        let response = StepResponse::from(gate);
353
354        assert_eq!(response.approvals_required, Some(3));
355        assert_eq!(response.approval_requirement, Some(requirement));
356        assert_eq!(response.approvals.len(), 1);
357        assert_eq!(response.approvals[0].approved_by, "alice");
358    }
359
360    #[tokio::test]
361    async fn a_future_deadline_reports_the_remaining_seconds() {
362        let response = StepResponse::from(gate_with_deadline(3600).await);
363
364        assert!(response.approval_deadline_at.is_some());
365        let remaining = response
366            .approval_seconds_remaining
367            .expect("a deadline yields a countdown");
368        assert!(remaining > 0 && remaining <= 3600, "got {remaining}");
369        assert_eq!(
370            response.approval_assignee,
371            Some(Assignee::group("release-managers"))
372        );
373    }
374
375    #[tokio::test]
376    async fn a_past_deadline_clamps_the_countdown_at_zero() {
377        let response = StepResponse::from(gate_with_deadline(-3600).await);
378        assert_eq!(response.approval_seconds_remaining, Some(0));
379    }
380
381    #[tokio::test]
382    async fn trace_id_is_exposed_in_step_response() {
383        let s = step().await;
384        let expected_trace_id = s.trace_id;
385        let response = StepResponse::from(s);
386
387        assert_eq!(response.trace_id, expected_trace_id);
388
389        let body = serde_json::to_value(&response).expect("serialize");
390        assert!(body["trace_id"].is_string());
391    }
392}