1use 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#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
23#[derive(Debug, Serialize, Deserialize)]
24pub struct StepResponse {
25 pub id: Uuid,
27 pub trace_id: Uuid,
29 pub run_id: Uuid,
31 pub name: String,
33 #[cfg_attr(feature = "openapi", schema(value_type = String))]
35 pub kind: StepKind,
36 pub position: u32,
38 pub status: StepStatus,
40 pub attempt: u32,
45 #[cfg_attr(feature = "openapi", schema(value_type = Option<std::collections::HashMap<String, serde_json::Value>>))]
47 pub input: Option<Value>,
48 #[cfg_attr(feature = "openapi", schema(value_type = Option<std::collections::HashMap<String, serde_json::Value>>))]
50 pub output: Option<Value>,
51 pub error: Option<String>,
53 pub duration_ms: u64,
55 #[cfg_attr(feature = "openapi", schema(value_type = f64))]
57 pub cost_usd: Decimal,
58 pub input_tokens: Option<u64>,
60 pub output_tokens: Option<u64>,
62 pub created_at: DateTime<Utc>,
64 pub updated_at: DateTime<Utc>,
66 pub started_at: Option<DateTime<Utc>>,
68 pub completed_at: Option<DateTime<Utc>>,
70 pub dependencies: Vec<Uuid>,
72 #[cfg_attr(feature = "openapi", schema(value_type = Option<serde_json::Value>))]
76 pub debug_messages: Option<Value>,
77 #[serde(default)]
82 pub artifacts: Vec<ArtifactResponse>,
83 pub approval_deadline_at: Option<DateTime<Utc>>,
85 pub approval_seconds_remaining: Option<i64>,
88 #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
92 pub approval_assignee: Option<Assignee>,
93 #[serde(default)]
96 pub approval_requirement: Option<ApprovalRequirement>,
97 #[serde(default)]
99 pub approvals: Vec<StepApproval>,
100 #[serde(default)]
103 pub approvals_required: Option<u32>,
104}
105
106impl StepResponse {
107 pub fn with_dependencies(step: Step, dependencies: Vec<Uuid>) -> Self {
113 Self::with_dependencies_and_artifacts(step, dependencies, Vec::new())
114 }
115
116 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 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 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}