1use chrono::{DateTime, Utc};
4use rust_decimal::Decimal;
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7use uuid::Uuid;
8
9use super::{ApprovalRequirement, Assignee, FsmState, StepApproval, StepKind, StepStatus};
10
11fn default_attempt() -> u32 {
14 1
15}
16
17pub fn step_trace_id(run_id: Uuid, name: &str, position: u32) -> Uuid {
38 Uuid::new_v5(
39 &Uuid::NAMESPACE_OID,
40 format!("{run_id}:{name}:{position}").as_bytes(),
41 )
42}
43
44#[derive(Debug, Clone, Serialize, Deserialize)]
56#[non_exhaustive]
57pub struct Step {
58 pub id: Uuid,
60 pub trace_id: Uuid,
65 pub run_id: Uuid,
67 pub name: String,
69 pub kind: StepKind,
71 pub position: u32,
78 pub status: FsmState<StepStatus>,
80 #[serde(default = "default_attempt")]
86 pub attempt: u32,
87 pub input: Option<Value>,
89 pub output: Option<Value>,
91 pub error: Option<String>,
93 pub duration_ms: u64,
95 pub cost_usd: Decimal,
97 pub input_tokens: Option<u64>,
99 #[serde(default)]
101 pub cache_read_input_tokens: Option<u64>,
102 #[serde(default)]
104 pub cache_creation_input_tokens: Option<u64>,
105 pub output_tokens: Option<u64>,
107 pub created_at: DateTime<Utc>,
109 pub updated_at: DateTime<Utc>,
111 pub started_at: Option<DateTime<Utc>>,
113 pub completed_at: Option<DateTime<Utc>>,
115 pub debug_messages: Option<Value>,
117 #[serde(default)]
119 pub is_error_handler: bool,
120 #[serde(default)]
125 pub approval_deadline_at: Option<DateTime<Utc>>,
126 #[serde(default)]
128 pub approval_stage: u32,
129 #[serde(default)]
131 pub approval_assignee: Option<Assignee>,
132 #[serde(default)]
135 pub approval_requirement: Option<ApprovalRequirement>,
136 #[serde(default)]
138 pub approvals: Vec<StepApproval>,
139 #[serde(default)]
141 pub account_id: Option<Uuid>,
142 #[serde(default)]
145 pub environment_id: Option<String>,
146}
147
148#[derive(Debug, Clone, Serialize, Deserialize)]
169pub struct NewStep {
170 pub run_id: Uuid,
172 pub trace_id: Uuid,
174 pub name: String,
176 pub kind: StepKind,
178 pub position: u32,
180 pub input: Option<Value>,
182 #[serde(default)]
184 pub is_error_handler: bool,
185}
186
187#[derive(Debug, Clone, Default, Serialize, Deserialize)]
204pub struct StepUpdate {
205 pub status: Option<StepStatus>,
207 pub output: Option<Value>,
209 pub error: Option<String>,
211 pub duration_ms: Option<u64>,
213 pub cost_usd: Option<Decimal>,
215 pub input_tokens: Option<u64>,
217 #[serde(default)]
219 pub cache_read_input_tokens: Option<u64>,
220 #[serde(default)]
222 pub cache_creation_input_tokens: Option<u64>,
223 pub output_tokens: Option<u64>,
225 pub started_at: Option<DateTime<Utc>>,
227 pub completed_at: Option<DateTime<Utc>>,
229 pub debug_messages: Option<Value>,
231 #[serde(default)]
234 pub approval_deadline_at: Option<DateTime<Utc>>,
235 #[serde(default)]
237 pub approval_stage: Option<u32>,
238 #[serde(default)]
240 pub approval_assignee: Option<Assignee>,
241 #[serde(default)]
243 pub approval_requirement: Option<ApprovalRequirement>,
244 #[serde(default)]
247 pub clear_approval_deadline: bool,
248 #[serde(default, skip_serializing_if = "Option::is_none")]
250 pub account_id: Option<Uuid>,
251 #[serde(default, skip_serializing_if = "Option::is_none")]
253 pub environment_id: Option<String>,
254}
255
256#[cfg(test)]
257mod tests {
258 use super::*;
259 use serde_json::json;
260
261 #[test]
262 fn newstep_serde_roundtrip() {
263 let new_step = NewStep {
264 run_id: Uuid::nil(),
265 trace_id: step_trace_id(Uuid::nil(), "build", 0),
266 name: "build".to_string(),
267 kind: StepKind::Shell,
268 position: 0,
269 input: Some(json!({"command": "cargo build"})),
270 is_error_handler: false,
271 };
272
273 let json = serde_json::to_string(&new_step).expect("serialize");
274 let back: NewStep = serde_json::from_str(&json).expect("deserialize");
275
276 assert_eq!(back.run_id, new_step.run_id);
277 assert_eq!(back.name, new_step.name);
278 assert_eq!(back.kind, new_step.kind);
279 assert_eq!(back.position, new_step.position);
280 assert_eq!(back.input, new_step.input);
281 }
282
283 #[test]
284 fn step_serde_preserves_all_fields() {
285 use crate::entities::FsmState;
286 use chrono::Utc;
287
288 let now = Utc::now();
289 let run_id = Uuid::now_v7();
290 let step = Step {
291 id: Uuid::now_v7(),
292 trace_id: step_trace_id(run_id, "test-step", 1),
293 run_id,
294 name: "test-step".to_string(),
295 kind: StepKind::Agent,
296 position: 1,
297 status: FsmState::new(StepStatus::Completed, Uuid::now_v7()),
298 attempt: 2,
299 input: Some(json!({"input": "data"})),
300 output: Some(json!({"output": "result"})),
301 error: None,
302 duration_ms: 2500,
303 cost_usd: Decimal::new(150, 2),
304 input_tokens: Some(100),
305 cache_read_input_tokens: Some(4000),
306 cache_creation_input_tokens: Some(300),
307 output_tokens: Some(200),
308 created_at: now,
309 updated_at: now,
310 started_at: Some(now),
311 completed_at: Some(now),
312 debug_messages: None,
313 is_error_handler: false,
314 approval_deadline_at: Some(now),
315 approval_stage: 2,
316 approval_assignee: Some(Assignee::group("sre-oncall")),
317 approval_requirement: Some(ApprovalRequirement {
318 reason: Some("amount > 10k".to_string()),
319 required_approvers: 2,
320 approver_groups: vec!["finance".to_string()],
321 }),
322 approvals: vec![StepApproval {
323 user_id: Uuid::now_v7(),
324 approved_by: "alice".to_string(),
325 at: now,
326 }],
327 account_id: Some(Uuid::now_v7()),
328 environment_id: Some("ironflow-env-0a1b2c".to_string()),
329 };
330
331 let json = serde_json::to_string(&step).expect("serialize");
332 let back: Step = serde_json::from_str(&json).expect("deserialize");
333
334 assert_eq!(back.id, step.id);
335 assert_eq!(back.run_id, step.run_id);
336 assert_eq!(back.name, step.name);
337 assert_eq!(back.kind, step.kind);
338 assert_eq!(back.position, step.position);
339 assert_eq!(back.status.state, step.status.state);
340 assert_eq!(back.attempt, step.attempt);
341 assert_eq!(back.input, step.input);
342 assert_eq!(back.output, step.output);
343 assert_eq!(back.error, step.error);
344 assert_eq!(back.duration_ms, step.duration_ms);
345 assert_eq!(back.cost_usd, step.cost_usd);
346 assert_eq!(back.input_tokens, step.input_tokens);
347 assert_eq!(back.cache_read_input_tokens, step.cache_read_input_tokens);
348 assert_eq!(
349 back.cache_creation_input_tokens,
350 step.cache_creation_input_tokens
351 );
352 assert_eq!(back.output_tokens, step.output_tokens);
353 assert_eq!(back.account_id, step.account_id);
354 assert_eq!(back.environment_id, step.environment_id);
355 assert_eq!(back.approval_deadline_at, step.approval_deadline_at);
356 assert_eq!(back.approval_stage, step.approval_stage);
357 assert_eq!(back.approval_assignee, step.approval_assignee);
358 assert_eq!(back.approval_requirement, step.approval_requirement);
359 assert_eq!(back.approvals, step.approvals);
360 }
361
362 #[test]
363 fn step_serde_defaults_approval_fields_when_absent() {
364 let run_id = Uuid::now_v7();
365 let payload = json!({
366 "id": Uuid::now_v7(),
367 "trace_id": step_trace_id(run_id, "legacy", 0),
368 "run_id": run_id,
369 "name": "legacy",
370 "kind": "shell",
371 "position": 0,
372 "status": {"state": "pending", "state_machine_id": Uuid::now_v7()},
373 "input": null,
374 "output": null,
375 "error": null,
376 "duration_ms": 0,
377 "cost_usd": 0.0,
378 "input_tokens": null,
379 "cache_read_input_tokens": null,
380 "cache_creation_input_tokens": null,
381 "output_tokens": null,
382 "created_at": "2026-09-21T12:00:00Z",
383 "updated_at": "2026-09-21T12:00:00Z",
384 "started_at": null,
385 "completed_at": null,
386 "debug_messages": null
387 });
388
389 let step: Step = serde_json::from_value(payload).expect("deserialize");
390
391 assert_eq!(step.approval_stage, 0);
392 assert!(step.cache_read_input_tokens.is_none());
393 assert!(step.cache_creation_input_tokens.is_none());
394 assert!(step.approval_deadline_at.is_none());
395 assert!(step.approval_assignee.is_none());
396 assert!(step.approval_requirement.is_none());
397 assert!(step.approvals.is_empty());
398 assert!(step.environment_id.is_none());
399 }
400
401 #[test]
402 fn stepupdate_default_is_no_changes() {
403 let update = StepUpdate::default();
404 assert!(update.status.is_none());
405 assert!(update.output.is_none());
406 assert!(update.error.is_none());
407 assert!(update.duration_ms.is_none());
408 assert!(update.cost_usd.is_none());
409 assert!(update.input_tokens.is_none());
410 assert!(update.cache_read_input_tokens.is_none());
411 assert!(update.cache_creation_input_tokens.is_none());
412 assert!(update.output_tokens.is_none());
413 assert!(update.started_at.is_none());
414 assert!(update.completed_at.is_none());
415 assert!(update.debug_messages.is_none());
416 assert!(update.approval_deadline_at.is_none());
417 assert!(update.approval_stage.is_none());
418 assert!(update.approval_assignee.is_none());
419 assert!(update.approval_requirement.is_none());
420 assert!(!update.clear_approval_deadline);
421 assert!(update.environment_id.is_none());
422 }
423
424 #[test]
425 fn stepupdate_serde_roundtrip() {
426 let update = StepUpdate {
427 status: Some(StepStatus::Completed),
428 output: Some(json!({"result": "ok"})),
429 error: None,
430 duration_ms: Some(1000),
431 cost_usd: Some(Decimal::new(50, 2)),
432 input_tokens: Some(50),
433 cache_read_input_tokens: Some(2000),
434 cache_creation_input_tokens: Some(150),
435 output_tokens: Some(75),
436 started_at: None,
437 completed_at: None,
438 debug_messages: None,
439 approval_deadline_at: Some(Utc::now()),
440 approval_stage: Some(1),
441 approval_assignee: Some(Assignee::group("sre-oncall")),
442 approval_requirement: Some(ApprovalRequirement::default()),
443 clear_approval_deadline: false,
444 account_id: Some(Uuid::now_v7()),
445 environment_id: Some("ironflow-env-0a1b2c".to_string()),
446 };
447
448 let json = serde_json::to_string(&update).expect("serialize");
449 let back: StepUpdate = serde_json::from_str(&json).expect("deserialize");
450
451 assert_eq!(back.status, update.status);
452 assert_eq!(back.output, update.output);
453 assert_eq!(back.duration_ms, update.duration_ms);
454 assert_eq!(back.account_id, update.account_id);
455 assert_eq!(back.environment_id, update.environment_id);
456 assert_eq!(back.cost_usd, update.cost_usd);
457 assert_eq!(back.input_tokens, update.input_tokens);
458 assert_eq!(back.cache_read_input_tokens, update.cache_read_input_tokens);
459 assert_eq!(
460 back.cache_creation_input_tokens,
461 update.cache_creation_input_tokens
462 );
463 assert_eq!(back.output_tokens, update.output_tokens);
464 assert_eq!(back.approval_deadline_at, update.approval_deadline_at);
465 assert_eq!(back.approval_stage, update.approval_stage);
466 assert_eq!(back.approval_assignee, update.approval_assignee);
467 assert_eq!(back.approval_requirement, update.approval_requirement);
468 assert_eq!(back.clear_approval_deadline, update.clear_approval_deadline);
469 }
470
471 #[test]
472 fn stepupdate_deserializes_without_cache_fields() {
473 let payload = json!({
474 "status": "completed",
475 "output": null,
476 "error": null,
477 "duration_ms": 10,
478 "cost_usd": null,
479 "input_tokens": 12,
480 "output_tokens": 3,
481 "started_at": null,
482 "completed_at": null,
483 "debug_messages": null
484 });
485
486 let update: StepUpdate = serde_json::from_value(payload).expect("deserialize");
487
488 assert_eq!(update.input_tokens, Some(12));
489 assert!(update.cache_read_input_tokens.is_none());
490 assert!(update.cache_creation_input_tokens.is_none());
491 }
492
493 #[test]
494 fn trace_id_is_deterministic() {
495 let run_id = Uuid::nil();
496 let id1 = step_trace_id(run_id, "build", 0);
497 let id2 = step_trace_id(run_id, "build", 0);
498 assert_eq!(id1, id2);
499 }
500
501 #[test]
502 fn trace_id_differs_for_different_inputs() {
503 let run_id = Uuid::nil();
504 let a = step_trace_id(run_id, "build", 0);
505 let b = step_trace_id(run_id, "test", 0);
506 let c = step_trace_id(run_id, "build", 1);
507 let d = step_trace_id(Uuid::max(), "build", 0);
508
509 assert_ne!(a, b);
510 assert_ne!(a, c);
511 assert_ne!(a, d);
512 }
513
514 #[test]
515 fn trace_id_is_uuid_v5() {
516 let id = step_trace_id(Uuid::nil(), "build", 0);
517 assert_eq!(id.get_version_num(), 5);
518 }
519}