Skip to main content

ironflow_api/entities/
step.rs

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