Skip to main content

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}