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}