ironflow_engine/context/steps/delay.rs
1//! Control-flow step implementations for [`WorkflowContext`].
2//!
3//! Adds the [`delay`](WorkflowContext::delay) method for persistent
4//! timed pauses that survive server restarts.
5
6use 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::context::lifecycle::check_replay_identity;
15use crate::error::EngineError;
16use crate::plan::lock_plan;
17
18impl WorkflowContext {
19 /// Execute a delay (timed pause) step.
20 ///
21 /// A zero-duration delay completes immediately. Otherwise, the
22 /// delay step is marked completed and the method returns
23 /// [`EngineError::DelaySleeping`] so the engine transitions the
24 /// run to [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping).
25 ///
26 /// On resume (after the worker picks up the re-queued run), the
27 /// delay step is replayed as completed via the replay mechanism.
28 ///
29 /// # Errors
30 ///
31 /// Returns [`EngineError::DelaySleeping`] to suspend the run. Returns
32 /// [`EngineError::ReplayDivergence`] when the step recorded at this
33 /// position has a different name or kind.
34 ///
35 /// # Examples
36 ///
37 /// ```no_run
38 /// use ironflow_engine::context::WorkflowContext;
39 /// use ironflow_engine::config::delay::DelayConfig;
40 /// use ironflow_engine::error::EngineError;
41 ///
42 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
43 /// ctx.delay("cooldown", DelayConfig::from_secs(300)).await?;
44 /// # Ok(())
45 /// # }
46 /// ```
47 pub async fn delay(&mut self, name: &str, config: DelayConfig) -> Result<(), EngineError> {
48 // Plan mode: record the pause without sleeping. The configured delay is
49 // the estimate when history has none -- a delay always lasts exactly
50 // as long as it was configured for.
51 if let Some(plan) = self.plan().cloned() {
52 self.position += 1;
53 let mut recorder = lock_plan(&plan);
54 recorder.seed_estimate(name, config.duration());
55 if recorder.record(
56 name,
57 StepKind::Custom("delay".to_string()),
58 &self.workflow_name,
59 None,
60 ) {
61 recorder.set_last(vec![name.to_string()]);
62 }
63 return Ok(());
64 }
65
66 let position = self.next_position();
67
68 if let Some(existing) = self.replay_steps().get(&position) {
69 check_replay_identity(
70 existing,
71 position,
72 name,
73 &StepKind::Custom("delay".to_string()),
74 )?;
75
76 if existing.status.state == StepStatus::Completed {
77 self.set_last_step_ids(vec![existing.id]);
78 info!(
79 run_id = %self.run_id(),
80 step = %name,
81 position,
82 "delay step replayed (already completed)"
83 );
84 return Ok(());
85 }
86 }
87
88 let trace_id = step_trace_id(self.run_id(), name, position);
89 let step = self
90 .store()
91 .create_step(NewStep {
92 run_id: self.run_id(),
93 trace_id,
94 name: name.to_string(),
95 kind: StepKind::Custom("delay".to_string()),
96 position,
97 input: Some(to_value(&config)?),
98 is_error_handler: false,
99 })
100 .await?;
101
102 let now = Utc::now();
103 self.start_step(step.id, now).await?;
104
105 if config.is_zero() {
106 self.store()
107 .update_step(
108 step.id,
109 StepUpdate {
110 status: Some(StepStatus::Completed),
111 completed_at: Some(now),
112 ..StepUpdate::default()
113 },
114 )
115 .await?;
116 self.set_last_step_ids(vec![step.id]);
117 info!(run_id = %self.run_id(), step = %name, "delay(0) completed immediately");
118 return Ok(());
119 }
120
121 let wake_at = now + Duration::seconds(config.duration_secs() as i64);
122
123 self.store()
124 .update_step(
125 step.id,
126 StepUpdate {
127 status: Some(StepStatus::Completed),
128 output: Some(json!({"wake_at": wake_at.to_rfc3339()})),
129 completed_at: Some(Utc::now()),
130 ..StepUpdate::default()
131 },
132 )
133 .await?;
134
135 self.set_last_step_ids(vec![step.id]);
136
137 Err(EngineError::DelaySleeping {
138 run_id: self.run_id(),
139 step_id: step.id,
140 wake_at,
141 })
142 }
143}