Skip to main content

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}