Skip to main content

ironflow_api/routes/
list_workflows.rs

1//! `GET /api/v1/workflows` — List registered workflows.
2
3use axum::extract::{Query, State};
4use axum::response::IntoResponse;
5use ironflow_auth::extractor::Authenticated;
6use serde::{Deserialize, Serialize};
7
8use crate::error::ApiError;
9use crate::response::ok;
10use crate::state::AppState;
11
12/// Query parameters for listing workflows.
13#[cfg_attr(feature = "openapi", derive(utoipa::IntoParams, utoipa::ToSchema))]
14#[derive(Debug, Deserialize)]
15pub struct ListWorkflowsQuery {
16    /// Optional case-insensitive partial match on workflow name.
17    pub name: Option<String>,
18    /// Optional case-insensitive partial match on the category path
19    /// (e.g. `etl` matches `Data/ETL` and `data/etl/nightly`).
20    ///
21    /// Pass `__uncategorized__` to list only workflows without any
22    /// category.
23    pub category: Option<String>,
24}
25
26/// Summary entry returned by `GET /api/v1/workflows`.
27#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
28#[derive(Debug, Serialize, Deserialize)]
29pub struct WorkflowSummary {
30    /// Workflow name (unique identifier).
31    pub name: String,
32    /// Optional `/`-separated category path.
33    pub category: Option<String>,
34    /// Current handler version.
35    pub version: Option<String>,
36    /// Optional 6-field cron expression for automatic execution.
37    #[serde(skip_serializing_if = "Option::is_none")]
38    pub schedule: Option<String>,
39}
40
41/// Sentinel value for the `category` query parameter that selects only
42/// uncategorized workflows.
43pub const UNCATEGORIZED_FILTER: &str = "__uncategorized__";
44
45/// List registered workflows, optionally filtered by name and category.
46///
47/// # Query Parameters
48///
49/// - `name` — Case-insensitive partial match on workflow name (optional)
50/// - `category` — Case-insensitive partial match on category path, or
51///   [`UNCATEGORIZED_FILTER`] to filter only uncategorized workflows
52///   (optional)
53#[cfg_attr(
54    feature = "openapi",
55    utoipa::path(
56        get,
57        path = "/api/v1/workflows",
58        tags = ["workflows"],
59        params(ListWorkflowsQuery),
60        responses(
61            (status = 200, description = "List of workflow summaries", body = [WorkflowSummary]),
62            (status = 401, description = "Unauthorized")
63        ),
64        security(("Bearer" = []))
65    )
66)]
67pub async fn list_workflows(
68    _auth: Authenticated,
69    State(state): State<AppState>,
70    Query(params): Query<ListWorkflowsQuery>,
71) -> Result<impl IntoResponse, ApiError> {
72    let mut summaries: Vec<WorkflowSummary> = state
73        .engine
74        .handler_names()
75        .into_iter()
76        .map(|name| {
77            let info = state.engine.handler_info(name);
78            let category = info.as_ref().and_then(|i| i.category.clone());
79            let version = info.as_ref().and_then(|i| i.version.clone());
80            let schedule = info.and_then(|i| i.schedule.map(|s| s.as_str().to_string()));
81            WorkflowSummary {
82                name: name.to_string(),
83                category,
84                version,
85                schedule,
86            }
87        })
88        .collect();
89
90    if let Some(ref filter) = params.name {
91        let lower = filter.to_lowercase();
92        summaries.retain(|s| s.name.to_lowercase().contains(&lower));
93    }
94
95    if let Some(ref cat_filter) = params.category {
96        if cat_filter == UNCATEGORIZED_FILTER {
97            summaries.retain(|s| s.category.is_none());
98        } else {
99            let needle = cat_filter.to_lowercase();
100            summaries.retain(|s| {
101                s.category
102                    .as_deref()
103                    .is_some_and(|c| c.to_lowercase().contains(&needle))
104            });
105        }
106    }
107
108    summaries.sort_by(|a, b| {
109        a.category
110            .as_deref()
111            .unwrap_or("")
112            .cmp(b.category.as_deref().unwrap_or(""))
113            .then_with(|| a.name.cmp(&b.name))
114    });
115
116    Ok(ok(summaries))
117}
118
119#[cfg(test)]
120mod tests {
121    use axum::Router;
122    use axum::body::Body;
123    use axum::http::{Request, StatusCode};
124    use axum::routing::get;
125    use http_body_util::BodyExt;
126    use ironflow_auth::jwt::AccessToken;
127    use ironflow_core::providers::claude::ClaudeCodeProvider;
128    use ironflow_engine::context::WorkflowContext;
129    use ironflow_engine::engine::Engine;
130    use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
131    use ironflow_engine::notify::Event;
132    use ironflow_store::memory::InMemoryStore;
133    use serde_json::{Value as JsonValue, from_slice, from_value};
134    use std::sync::Arc;
135    use tokio::sync::broadcast;
136    use tower::ServiceExt;
137    use uuid::Uuid;
138
139    use super::*;
140
141    fn make_auth_header(state: &AppState) -> String {
142        let user_id = Uuid::now_v7();
143        let token = AccessToken::for_user(user_id, "testuser", false, &state.jwt_config).unwrap();
144        format!("Bearer {}", token.0)
145    }
146
147    struct TestWorkflow;
148
149    impl WorkflowHandler for TestWorkflow {
150        fn name(&self) -> &str {
151            "test-workflow"
152        }
153
154        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
155            Box::pin(async move { Ok(()) })
156        }
157    }
158
159    struct AnotherWorkflow;
160
161    impl WorkflowHandler for AnotherWorkflow {
162        fn name(&self) -> &str {
163            "another-workflow"
164        }
165
166        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
167            Box::pin(async move { Ok(()) })
168        }
169    }
170
171    struct EtlNightlyWorkflow;
172
173    impl WorkflowHandler for EtlNightlyWorkflow {
174        fn name(&self) -> &str {
175            "etl-nightly"
176        }
177        fn category(&self) -> Option<&str> {
178            Some("data/etl")
179        }
180        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
181            Box::pin(async move { Ok(()) })
182        }
183    }
184
185    struct ReportsDailyWorkflow;
186
187    impl WorkflowHandler for ReportsDailyWorkflow {
188        fn name(&self) -> &str {
189            "reports-daily"
190        }
191        fn category(&self) -> Option<&str> {
192            Some("data/reports")
193        }
194        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
195            Box::pin(async move { Ok(()) })
196        }
197    }
198
199    fn base_state(engine: Engine) -> AppState {
200        let store = Arc::new(InMemoryStore::new());
201        Arc::new(InMemoryStore::new());
202        let jwt_config = Arc::new(ironflow_auth::jwt::JwtConfig {
203            secret: "test-secret".to_string(),
204            access_token_ttl_secs: 900,
205            refresh_token_ttl_secs: 604800,
206            cookie_domain: None,
207            cookie_secure: false,
208        });
209        let (event_sender, _) = broadcast::channel::<Event>(1);
210        AppState::new(
211            store,
212            Arc::new(engine),
213            jwt_config,
214            "test-worker-token".to_string(),
215            event_sender,
216        )
217    }
218
219    fn test_state() -> AppState {
220        let store = Arc::new(InMemoryStore::new());
221        let provider = Arc::new(ClaudeCodeProvider::new());
222        let mut engine = Engine::new(store.clone(), provider);
223        engine.register(TestWorkflow).unwrap();
224        engine.register(AnotherWorkflow).unwrap();
225        base_state(engine)
226    }
227
228    fn test_state_with_categories() -> AppState {
229        let store = Arc::new(InMemoryStore::new());
230        let provider = Arc::new(ClaudeCodeProvider::new());
231        let mut engine = Engine::new(store.clone(), provider);
232        engine.register(TestWorkflow).unwrap();
233        engine.register(EtlNightlyWorkflow).unwrap();
234        engine.register(ReportsDailyWorkflow).unwrap();
235        base_state(engine)
236    }
237
238    async fn run_request(state: AppState, uri: &str) -> (StatusCode, Vec<WorkflowSummary>) {
239        let auth_header = make_auth_header(&state);
240        let app = Router::new()
241            .route("/", get(list_workflows))
242            .with_state(state);
243
244        let req = Request::builder()
245            .uri(uri)
246            .header("authorization", auth_header)
247            .body(Body::empty())
248            .unwrap();
249
250        let resp = app.oneshot(req).await.unwrap();
251        let status = resp.status();
252        let body = resp.into_body().collect().await.unwrap().to_bytes();
253        let json_val: JsonValue = from_slice(&body).unwrap();
254        let summaries: Vec<WorkflowSummary> = from_value(json_val["data"].clone()).unwrap();
255        (status, summaries)
256    }
257
258    #[tokio::test]
259    async fn list_workflows_empty() {
260        let store = Arc::new(InMemoryStore::new());
261        let provider = Arc::new(ClaudeCodeProvider::new());
262        let engine = Engine::new(store.clone(), provider);
263        let state = base_state(engine);
264
265        let (status, summaries) = run_request(state, "/").await;
266        assert_eq!(status, StatusCode::OK);
267        assert!(summaries.is_empty());
268    }
269
270    #[tokio::test]
271    async fn list_workflows_multiple_returns_summaries() {
272        let state = test_state();
273        let (status, summaries) = run_request(state, "/").await;
274        assert_eq!(status, StatusCode::OK);
275        assert_eq!(summaries.len(), 2);
276        assert!(summaries.iter().any(|s| s.name == "test-workflow"));
277        assert!(summaries.iter().any(|s| s.name == "another-workflow"));
278        assert!(summaries.iter().all(|s| s.category.is_none()));
279    }
280
281    #[tokio::test]
282    async fn list_workflows_filtered_by_name() {
283        let state = test_state();
284        let (_, summaries) = run_request(state, "/?name=test").await;
285        assert_eq!(summaries.len(), 1);
286        assert_eq!(summaries[0].name, "test-workflow");
287    }
288
289    #[tokio::test]
290    async fn list_workflows_filter_name_case_insensitive() {
291        let state = test_state();
292        let (_, summaries) = run_request(state, "/?name=TEST").await;
293        assert_eq!(summaries.len(), 1);
294        assert_eq!(summaries[0].name, "test-workflow");
295    }
296
297    #[tokio::test]
298    async fn list_workflows_filter_no_match() {
299        let state = test_state();
300        let (_, summaries) = run_request(state, "/?name=nonexistent").await;
301        assert!(summaries.is_empty());
302    }
303
304    #[tokio::test]
305    async fn list_workflows_returns_category_when_present() {
306        let state = test_state_with_categories();
307        let (_, summaries) = run_request(state, "/").await;
308        let etl = summaries.iter().find(|s| s.name == "etl-nightly").unwrap();
309        assert_eq!(etl.category.as_deref(), Some("data/etl"));
310        let test = summaries
311            .iter()
312            .find(|s| s.name == "test-workflow")
313            .unwrap();
314        assert!(test.category.is_none());
315    }
316
317    #[tokio::test]
318    async fn list_workflows_filter_by_category_partial() {
319        let state = test_state_with_categories();
320        let (_, summaries) = run_request(state, "/?category=data").await;
321        assert_eq!(summaries.len(), 2);
322        assert!(
323            summaries
324                .iter()
325                .all(|s| { s.category.as_deref().is_some_and(|c| c.contains("data")) })
326        );
327    }
328
329    #[tokio::test]
330    async fn list_workflows_filter_by_category_nested_substring() {
331        let state = test_state_with_categories();
332        let (_, summaries) = run_request(state, "/?category=etl").await;
333        assert_eq!(summaries.len(), 1);
334        assert_eq!(summaries[0].name, "etl-nightly");
335    }
336
337    #[tokio::test]
338    async fn list_workflows_filter_by_category_case_insensitive() {
339        let state = test_state_with_categories();
340        let (_, summaries) = run_request(state, "/?category=DATA").await;
341        assert_eq!(summaries.len(), 2);
342        let (_, summaries) = run_request(test_state_with_categories(), "/?category=Etl").await;
343        assert_eq!(summaries.len(), 1);
344        assert_eq!(summaries[0].name, "etl-nightly");
345    }
346
347    #[tokio::test]
348    async fn list_workflows_filter_by_category_no_match_excludes_uncategorized() {
349        let state = test_state_with_categories();
350        let (_, summaries) = run_request(state, "/?category=nonexistent").await;
351        assert!(summaries.is_empty());
352    }
353
354    #[tokio::test]
355    async fn list_workflows_filter_uncategorized_sentinel() {
356        let state = test_state_with_categories();
357        let (_, summaries) = run_request(state, "/?category=__uncategorized__").await;
358        assert_eq!(summaries.len(), 1);
359        assert_eq!(summaries[0].name, "test-workflow");
360    }
361
362    struct ScheduledWorkflow {
363        schedule: ironflow_engine::prelude::CronSchedule,
364    }
365    impl ScheduledWorkflow {
366        fn new() -> Self {
367            Self {
368                schedule: ironflow_engine::prelude::CronSchedule::new("0 30 9 * * MON-FRI")
369                    .unwrap(),
370            }
371        }
372    }
373
374    impl WorkflowHandler for ScheduledWorkflow {
375        fn name(&self) -> &str {
376            "scheduled-workflow"
377        }
378        fn schedule(&self) -> Option<&ironflow_engine::prelude::CronSchedule> {
379            Some(&self.schedule)
380        }
381        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
382            Box::pin(async move { Ok(()) })
383        }
384    }
385
386    #[tokio::test]
387    async fn list_workflows_returns_schedule_when_present() {
388        let store = Arc::new(InMemoryStore::new());
389        let provider = Arc::new(ClaudeCodeProvider::new());
390        let mut engine = Engine::new(store.clone(), provider);
391        engine.register(TestWorkflow).unwrap();
392        engine.register(ScheduledWorkflow::new()).unwrap();
393        let state = base_state(engine);
394
395        let (_, summaries) = run_request(state, "/").await;
396        let scheduled = summaries
397            .iter()
398            .find(|s| s.name == "scheduled-workflow")
399            .unwrap();
400        assert_eq!(scheduled.schedule.as_deref(), Some("0 30 9 * * MON-FRI"));
401
402        let unscheduled = summaries
403            .iter()
404            .find(|s| s.name == "test-workflow")
405            .unwrap();
406        assert!(unscheduled.schedule.is_none());
407    }
408}