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