1use 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#[cfg_attr(feature = "openapi", derive(utoipa::IntoParams, utoipa::ToSchema))]
14#[derive(Debug, Deserialize)]
15pub struct ListWorkflowsQuery {
16 pub name: Option<String>,
18 pub category: Option<String>,
24}
25
26#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
28#[derive(Debug, Serialize, Deserialize)]
29pub struct WorkflowSummary {
30 pub name: String,
32 pub category: Option<String>,
34 pub version: Option<String>,
36 #[serde(skip_serializing_if = "Option::is_none")]
38 pub schedule: Option<String>,
39}
40
41pub const UNCATEGORIZED_FILTER: &str = "__uncategorized__";
44
45#[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}