Skip to main content

ironflow_engine/
wake.rs

1//! [`RunWaker`] -- wakes `Sleeping` runs whose `scheduled_at` has passed.
2//!
3//! A run goes `Sleeping` when a [`delay`](crate::context::WorkflowContext::delay)
4//! step pauses it, or when
5//! [`wait_for_signal`](crate::context::WorkflowContext::wait_for_signal) suspends
6//! it until a signal or its deadline. Either way the wake-up time lives on the
7//! run itself (`scheduled_at`), so the timer survives an API or worker restart:
8//! nothing is held in memory, and a brand-new process picks the run up on its
9//! next tick.
10//!
11//! [`RunWaker::tick`] claims every due run and moves it back to `Pending` in
12//! the same transaction that clears `scheduled_at`, so a run wakes **exactly
13//! once** even with several API instances running a waker. Under
14//! [`ExecutionMode::Workers`] a worker then picks the run up; under
15//! [`ExecutionMode::Local`] the waker resumes it in-process.
16
17use std::sync::Arc;
18
19use tracing::info;
20
21use ironflow_store::models::Run;
22
23use crate::engine::{Engine, ExecutionMode};
24use crate::error::EngineError;
25
26/// How many runs a single [`RunWaker::tick`] wakes.
27pub const DEFAULT_WAKE_BATCH_SIZE: u32 = 50;
28
29/// Wakes `Sleeping` runs whose wake-up time has passed.
30///
31/// # Examples
32///
33/// ```no_run
34/// use std::sync::Arc;
35/// use ironflow_core::providers::claude::ClaudeCodeProvider;
36/// use ironflow_engine::engine::Engine;
37/// use ironflow_engine::wake::RunWaker;
38/// use ironflow_store::memory::InMemoryStore;
39/// use ironflow_store::store::Store;
40///
41/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
42/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
43/// let engine = Arc::new(Engine::new(store, Arc::new(ClaudeCodeProvider::new())));
44///
45/// let waker = RunWaker::new(engine).batch_size(10);
46/// let woken = waker.tick().await?;
47/// println!("{} runs woken", woken.len());
48/// # Ok(())
49/// # }
50/// ```
51#[derive(Debug, Clone)]
52pub struct RunWaker {
53    engine: Arc<Engine>,
54    batch_size: u32,
55}
56
57impl RunWaker {
58    /// Create a waker with the default batch size.
59    ///
60    /// # Examples
61    ///
62    /// ```no_run
63    /// use std::sync::Arc;
64    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
65    /// use ironflow_engine::engine::Engine;
66    /// use ironflow_engine::wake::RunWaker;
67    /// use ironflow_store::memory::InMemoryStore;
68    /// use ironflow_store::store::Store;
69    ///
70    /// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
71    /// let engine = Arc::new(Engine::new(store, Arc::new(ClaudeCodeProvider::new())));
72    /// let waker = RunWaker::new(engine);
73    /// ```
74    pub fn new(engine: Arc<Engine>) -> Self {
75        Self {
76            engine,
77            batch_size: DEFAULT_WAKE_BATCH_SIZE,
78        }
79    }
80
81    /// Set how many runs a single [`tick`](Self::tick) wakes.
82    ///
83    /// # Examples
84    ///
85    /// ```no_run
86    /// use std::sync::Arc;
87    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
88    /// use ironflow_engine::engine::Engine;
89    /// use ironflow_engine::wake::RunWaker;
90    /// use ironflow_store::memory::InMemoryStore;
91    /// use ironflow_store::store::Store;
92    ///
93    /// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
94    /// let engine = Arc::new(Engine::new(store, Arc::new(ClaudeCodeProvider::new())));
95    /// let waker = RunWaker::new(engine).batch_size(10);
96    /// ```
97    pub fn batch_size(mut self, batch_size: u32) -> Self {
98        self.batch_size = batch_size;
99        self
100    }
101
102    /// Claim and wake one batch of due `Sleeping` runs.
103    ///
104    /// Every claimed run is `Pending` when this returns. Under
105    /// [`ExecutionMode::Local`] each one is then resumed in a background task,
106    /// so a long workflow never holds the waker loop; a failed resume is
107    /// logged, not rolled back.
108    ///
109    /// Returns the woken runs as they were right after the transition.
110    ///
111    /// # Errors
112    ///
113    /// Returns [`EngineError::Store`] if the batch cannot be claimed.
114    ///
115    /// # Examples
116    ///
117    /// ```no_run
118    /// use ironflow_engine::wake::RunWaker;
119    ///
120    /// # use ironflow_engine::error::EngineError;
121    /// # async fn example(waker: &RunWaker) -> Result<(), EngineError> {
122    /// for run in waker.tick().await? {
123    ///     println!("woke {}", run.id);
124    /// }
125    /// # Ok(())
126    /// # }
127    /// ```
128    pub async fn tick(&self) -> Result<Vec<Run>, EngineError> {
129        let runs = self
130            .engine
131            .store()
132            .claim_due_sleeping_runs(self.batch_size)
133            .await?;
134
135        for run in &runs {
136            info!(
137                run_id = %run.id,
138                workflow = %run.workflow_name,
139                "sleeping run woken"
140            );
141            if self.engine.execution_mode() == ExecutionMode::Local {
142                self.engine.spawn_local_resume(run.id);
143            }
144        }
145
146        Ok(runs)
147    }
148}