ironflow_engine/context/steps/
delay.rs1use chrono::{Duration, Utc};
7use serde_json::{json, to_value};
8use tracing::info;
9
10use ironflow_store::models::{NewStep, StepKind, StepStatus, StepUpdate, step_trace_id};
11
12use crate::config::delay::DelayConfig;
13use crate::context::WorkflowContext;
14use crate::error::EngineError;
15use crate::plan::lock_plan;
16
17impl WorkflowContext {
18 pub async fn delay(&mut self, name: &str, config: DelayConfig) -> Result<(), EngineError> {
45 if let Some(plan) = self.plan().cloned() {
49 self.position += 1;
50 let mut recorder = lock_plan(&plan);
51 recorder.seed_estimate(name, config.duration());
52 if recorder.record(
53 name,
54 StepKind::Custom("delay".to_string()),
55 &self.workflow_name,
56 None,
57 ) {
58 recorder.set_last(vec![name.to_string()]);
59 }
60 return Ok(());
61 }
62
63 let position = self.next_position();
64
65 if let Some(existing) = self.replay_steps().get(&position)
66 && existing.kind == StepKind::Custom("delay".to_string())
67 && existing.status.state == StepStatus::Completed
68 {
69 self.set_last_step_ids(vec![existing.id]);
70 info!(
71 run_id = %self.run_id(),
72 step = %name,
73 position,
74 "delay step replayed (already completed)"
75 );
76 return Ok(());
77 }
78
79 let trace_id = step_trace_id(self.run_id(), name, position);
80 let step = self
81 .store()
82 .create_step(NewStep {
83 run_id: self.run_id(),
84 trace_id,
85 name: name.to_string(),
86 kind: StepKind::Custom("delay".to_string()),
87 position,
88 input: Some(to_value(&config)?),
89 is_error_handler: false,
90 })
91 .await?;
92
93 let now = Utc::now();
94 self.start_step(step.id, now).await?;
95
96 if config.is_zero() {
97 self.store()
98 .update_step(
99 step.id,
100 StepUpdate {
101 status: Some(StepStatus::Completed),
102 completed_at: Some(now),
103 ..StepUpdate::default()
104 },
105 )
106 .await?;
107 self.set_last_step_ids(vec![step.id]);
108 info!(run_id = %self.run_id(), step = %name, "delay(0) completed immediately");
109 return Ok(());
110 }
111
112 let wake_at = now + Duration::seconds(config.duration_secs() as i64);
113
114 self.store()
115 .update_step(
116 step.id,
117 StepUpdate {
118 status: Some(StepStatus::Completed),
119 output: Some(json!({"wake_at": wake_at.to_rfc3339()})),
120 completed_at: Some(Utc::now()),
121 ..StepUpdate::default()
122 },
123 )
124 .await?;
125
126 self.set_last_step_ids(vec![step.id]);
127
128 Err(EngineError::DelaySleeping {
129 run_id: self.run_id(),
130 step_id: step.id,
131 wake_at,
132 })
133 }
134}