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::error::EngineError;
15use crate::plan::lock_plan;
16
17impl WorkflowContext {
18    /// Execute a delay (timed pause) step.
19    ///
20    /// A zero-duration delay completes immediately. Otherwise, the
21    /// delay step is marked completed and the method returns
22    /// [`EngineError::DelaySleeping`] so the engine transitions the
23    /// run to [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping).
24    ///
25    /// On resume (after the worker picks up the re-queued run), the
26    /// delay step is replayed as completed via the replay mechanism.
27    ///
28    /// # Errors
29    ///
30    /// Returns [`EngineError::DelaySleeping`] to suspend the run.
31    ///
32    /// # Examples
33    ///
34    /// ```no_run
35    /// use ironflow_engine::context::WorkflowContext;
36    /// use ironflow_engine::config::delay::DelayConfig;
37    /// use ironflow_engine::error::EngineError;
38    ///
39    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
40    /// ctx.delay("cooldown", DelayConfig::from_secs(300)).await?;
41    /// # Ok(())
42    /// # }
43    /// ```
44    pub async fn delay(&mut self, name: &str, config: DelayConfig) -> Result<(), EngineError> {
45        // Plan mode: record the pause without sleeping. The configured delay is
46        // the estimate when history has none -- a delay always lasts exactly
47        // as long as it was configured for.
48        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}