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 cache_read_input_tokens: Option<u64>,
62 pub cache_creation_input_tokens: Option<u64>,
64 pub output_tokens: Option<u64>,
66 pub created_at: DateTime<Utc>,
68 pub updated_at: DateTime<Utc>,
70 pub started_at: Option<DateTime<Utc>>,
72 pub completed_at: Option<DateTime<Utc>>,
74 pub dependencies: Vec<Uuid>,
76 #[cfg_attr(feature = "openapi", schema(value_type = Option<serde_json::Value>))]
80 pub debug_messages: Option<Value>,
81 #[serde(default)]
86 pub artifacts: Vec<ArtifactResponse>,
87 pub approval_deadline_at: Option<DateTime<Utc>>,
89 pub approval_seconds_remaining: Option<i64>,
92 #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
96 pub approval_assignee: Option<Assignee>,
97 #[serde(default)]
100 pub approval_requirement: Option<ApprovalRequirement>,
101 #[serde(default)]
103 pub approvals: Vec<StepApproval>,
104 #[serde(default)]
107 pub approvals_required: Option<u32>,
108}
109
110impl StepResponse {
111 pub fn with_dependencies(step: Step, dependencies: Vec<Uuid>) -> Self {
117 Self::with_dependencies_and_artifacts(step, dependencies, Vec::new())
118 }
119
120 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 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 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}