ironflow_api/state.rs
1//! Application state and dependency injection.
2//!
3//! [`AppState`] holds the shared [`Store`] and [`Engine`] used by all handlers.
4
5use std::sync::Arc;
6#[cfg(feature = "prometheus")]
7use std::sync::OnceLock;
8
9use axum::extract::FromRef;
10#[cfg(feature = "prometheus")]
11use metrics_exporter_prometheus::{PrometheusBuilder, PrometheusHandle};
12use tokio::sync::broadcast;
13use uuid::Uuid;
14
15use tokio_util::sync::CancellationToken;
16use tracing::warn;
17
18use ironflow_artifacts::blob_store::BlobStore;
19use ironflow_auth::jwt::JwtConfig;
20use ironflow_engine::engine::Engine;
21use ironflow_engine::notify::{Event, WorkflowEventBus};
22use ironflow_store::entities::Run;
23use ironflow_store::store::Store;
24
25use crate::error::ApiError;
26use crate::escalator::Escalator;
27use crate::reaper::Reaper;
28use crate::schedule_sync::sync_handler_schedules;
29use crate::schedule_ticker::ScheduleTicker;
30
31/// Global application state.
32///
33/// Holds the shared store (runs, users, API keys, secrets) and engine,
34/// extracted by handlers using Axum's state extraction mechanism.
35///
36/// # Examples
37///
38/// ```no_run
39/// use ironflow_api::state::AppState;
40/// use ironflow_auth::jwt::JwtConfig;
41/// use ironflow_store::prelude::*;
42/// use ironflow_store::store::Store;
43/// use ironflow_engine::engine::Engine;
44/// use ironflow_core::providers::claude::ClaudeCodeProvider;
45/// use std::sync::Arc;
46///
47/// # async fn example() {
48/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
49/// let provider = Arc::new(ClaudeCodeProvider::new());
50/// let engine = Arc::new(Engine::new(store.clone(), provider));
51/// let jwt_config = Arc::new(JwtConfig {
52/// secret: "secret".to_string(),
53/// access_token_ttl_secs: 900,
54/// refresh_token_ttl_secs: 604800,
55/// cookie_domain: None,
56/// cookie_secure: false,
57/// });
58/// let broadcaster = ironflow_api::sse::SseBroadcaster::new();
59/// let state = AppState::new(store, engine, jwt_config, "token".to_string(), broadcaster.sender());
60/// # }
61/// ```
62#[derive(Clone)]
63pub struct AppState {
64 /// The unified backing store for runs, steps, users, API keys, and secrets.
65 pub store: Arc<dyn Store>,
66 /// The workflow orchestration engine.
67 pub engine: Arc<Engine>,
68 /// JWT configuration for auth tokens.
69 pub jwt_config: Arc<JwtConfig>,
70 /// Static token for worker-to-API authentication.
71 pub worker_token: String,
72 /// Broadcast sender for SSE event streaming.
73 pub event_sender: broadcast::Sender<Event>,
74 /// Per-run event bus for real-time workflow monitoring.
75 ///
76 /// When set, the `GET /api/v1/runs/{id}/events` route subscribes to this
77 /// bus and streams [`WorkflowEvent`](ironflow_engine::notify::WorkflowEvent)s
78 /// via SSE. `None` when the engine was not configured with a bus.
79 pub event_bus: Option<WorkflowEventBus>,
80 /// Where artifact bytes live, when artifacts are enabled.
81 ///
82 /// `None` on a deployment that has not configured artifact storage: the
83 /// artifact routes answer `501` and every other endpoint is unaffected.
84 pub blob_store: Option<Arc<dyn BlobStore>>,
85 /// Prometheus metrics handle (only when `prometheus` feature is enabled).
86 #[cfg(feature = "prometheus")]
87 pub prometheus_handle: PrometheusHandle,
88}
89
90impl FromRef<AppState> for Arc<dyn Store> {
91 fn from_ref(state: &AppState) -> Self {
92 Arc::clone(&state.store)
93 }
94}
95
96impl FromRef<AppState> for Arc<JwtConfig> {
97 fn from_ref(state: &AppState) -> Self {
98 Arc::clone(&state.jwt_config)
99 }
100}
101
102#[cfg(feature = "prometheus")]
103impl FromRef<AppState> for PrometheusHandle {
104 fn from_ref(state: &AppState) -> Self {
105 state.prometheus_handle.clone()
106 }
107}
108
109impl AppState {
110 /// Create a new `AppState`.
111 ///
112 /// When the `prometheus` feature is enabled, a global Prometheus recorder
113 /// is installed (once) and its handle is stored in the state.
114 ///
115 /// # Panics
116 ///
117 /// Panics if a Prometheus recorder cannot be installed (should only
118 /// happen if another incompatible recorder was set elsewhere).
119 pub fn new(
120 store: Arc<dyn Store>,
121 engine: Arc<Engine>,
122 jwt_config: Arc<JwtConfig>,
123 worker_token: String,
124 event_sender: broadcast::Sender<Event>,
125 ) -> Self {
126 Self {
127 store,
128 engine,
129 jwt_config,
130 worker_token,
131 event_sender,
132 event_bus: None,
133 blob_store: None,
134 #[cfg(feature = "prometheus")]
135 prometheus_handle: Self::global_prometheus_handle(),
136 }
137 }
138
139 /// Enable artifacts by attaching the backend that holds their bytes.
140 ///
141 /// # Examples
142 ///
143 /// ```no_run
144 /// use std::sync::Arc;
145 ///
146 /// use ironflow_api::state::AppState;
147 /// use ironflow_artifacts::blob_store::BlobStore;
148 /// use ironflow_artifacts::local::LocalBlobStore;
149 ///
150 /// # fn example(state: AppState) -> AppState {
151 /// let blob: Arc<dyn BlobStore> = Arc::new(LocalBlobStore::new("/var/lib/ironflow/artifacts"));
152 /// state.with_blob_store(blob)
153 /// # }
154 /// ```
155 pub fn with_blob_store(mut self, blob_store: Arc<dyn BlobStore>) -> Self {
156 self.blob_store = Some(blob_store);
157 self
158 }
159
160 /// Attach a [`WorkflowEventBus`] for per-run SSE streaming.
161 ///
162 /// When set, `GET /api/v1/runs/{id}/events` streams step-level events
163 /// for a specific workflow run. When absent the route returns an empty
164 /// SSE stream (with keep-alive).
165 ///
166 /// # Examples
167 ///
168 /// ```no_run
169 /// use ironflow_api::state::AppState;
170 /// use ironflow_engine::notify::WorkflowEventBus;
171 ///
172 /// # fn example(state: AppState) -> AppState {
173 /// state.with_event_bus(WorkflowEventBus::new())
174 /// # }
175 /// ```
176 pub fn with_event_bus(mut self, bus: WorkflowEventBus) -> Self {
177 self.event_bus = Some(bus);
178 self
179 }
180
181 /// The artifact backend, or a `501` error when artifacts are disabled.
182 ///
183 /// # Errors
184 ///
185 /// Returns [`ApiError::ArtifactStorageUnavailable`] when no backend is attached.
186 pub fn blob_store_or_501(&self) -> Result<&Arc<dyn BlobStore>, ApiError> {
187 self.blob_store
188 .as_ref()
189 .ok_or(ApiError::ArtifactStorageUnavailable)
190 }
191
192 /// Install (or reuse) a global Prometheus recorder and return its handle.
193 #[cfg(feature = "prometheus")]
194 fn global_prometheus_handle() -> PrometheusHandle {
195 static HANDLE: OnceLock<PrometheusHandle> = OnceLock::new();
196 HANDLE
197 .get_or_init(|| {
198 PrometheusBuilder::new()
199 .install_recorder()
200 .expect("failed to install Prometheus recorder")
201 })
202 .clone()
203 }
204
205 /// Fetch a run by ID or return 404.
206 ///
207 /// # Errors
208 ///
209 /// Returns `ApiError::RunNotFound` if the run does not exist.
210 /// Returns `ApiError::Store` if there is a store error.
211 pub async fn get_run_or_404(&self, id: Uuid) -> Result<Run, ApiError> {
212 self.store
213 .get_run(id)
214 .await
215 .map_err(ApiError::from)?
216 .ok_or(ApiError::RunNotFound(id))
217 }
218
219 /// Spawn the built-in background tasks and return their shared shutdown token.
220 ///
221 /// This starts:
222 /// - **Schedule sync**: seeds DB rows for handler-declared schedules.
223 /// - **Schedule ticker**: polls due schedules and creates runs.
224 /// - **Reaper**: recovers runs abandoned by dead workers.
225 /// - **Escalator**: resolves approval gates that missed their SLA deadline.
226 ///
227 /// Call this once after building the `AppState`, before serving requests.
228 /// Drop the returned [`CancellationToken`] (or call `.cancel()`) to stop
229 /// all tasks gracefully.
230 ///
231 /// # Examples
232 ///
233 /// ```no_run
234 /// use ironflow_api::state::AppState;
235 ///
236 /// # async fn example(state: AppState) {
237 /// let shutdown = state.spawn_background_tasks().await;
238 /// // ... serve requests ...
239 /// shutdown.cancel();
240 /// # }
241 /// ```
242 pub async fn spawn_background_tasks(&self) -> CancellationToken {
243 if let Err(err) = sync_handler_schedules(&self.engine, self.store.as_ref()).await {
244 warn!(error = %err, "failed to sync handler-declared schedules");
245 }
246
247 let shutdown = CancellationToken::new();
248 tokio::spawn(ScheduleTicker::new(self.store.clone()).run(shutdown.clone()));
249 tokio::spawn(Reaper::new(self.store.clone(), self.engine.clone()).run(shutdown.clone()));
250 tokio::spawn(Escalator::new(self.engine.clone()).run(shutdown.clone()));
251 shutdown
252 }
253}
254
255#[cfg(test)]
256mod tests {
257 use super::*;
258 use ironflow_core::providers::claude::ClaudeCodeProvider;
259 use ironflow_store::memory::InMemoryStore;
260 use ironflow_store::store::Store;
261
262 fn test_state() -> AppState {
263 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
264 let provider = Arc::new(ClaudeCodeProvider::new());
265 let engine = Arc::new(Engine::new(store.clone(), provider));
266 let jwt_config = Arc::new(JwtConfig {
267 secret: "test-secret".to_string(),
268 access_token_ttl_secs: 900,
269 refresh_token_ttl_secs: 604800,
270 cookie_domain: None,
271 cookie_secure: false,
272 });
273 let (event_sender, _) = broadcast::channel::<Event>(1);
274 AppState::new(
275 store,
276 engine,
277 jwt_config,
278 "test-worker-token".to_string(),
279 event_sender,
280 )
281 }
282
283 #[test]
284 fn app_state_cloneable() {
285 let state = test_state();
286 let _cloned = state.clone();
287 }
288
289 #[test]
290 fn app_state_from_ref() {
291 let state = test_state();
292 let extracted: Arc<dyn Store> = Arc::from_ref(&state);
293 assert!(Arc::ptr_eq(&extracted, &state.store));
294 }
295}