1use axum::{
4 Router,
5 extract::DefaultBodyLimit,
6 http::{HeaderName, HeaderValue, Method, StatusCode, header},
7 routing::{any, get, post},
8};
9use tower_http::cors::CorsLayer;
10
11use super::authoring::compile_source;
12use super::awl::{
13 bind_run, check, create_document, deploy_authoring, edit, emit, format, get_document,
14 get_layout, get_revision, get_run_status, list_documents, put_document, put_layout, scaffold,
15 worker_availability,
16};
17use super::cluster_command::cluster_command;
18use super::deploy::{list_versions, route_version, unload_version, upload_package};
19use super::dev_ui::{dev_register_mock, dev_replay_run, dev_trigger_run};
20use super::events::subscribe_events_socket;
21use super::intervene::{intervene, list_attempts};
22use super::schedules::{
23 create_schedule, delete_schedule, describe_schedule, list_schedules, pause_schedule,
24 resume_schedule, update_schedule,
25};
26use super::transcripts::{fetch_transcript, list_transcript_streams};
27use super::whoami::whoami;
28use super::workflows::{
29 cancel_workflow, count_workflows, describe_workflow, get_workflows, list_namespace_records,
30 list_namespaces, post_list_workflows, post_namespace, query_workflow, reopen_workflow,
31 set_namespace_placement, signal_workflow, start_workflow,
32};
33use crate::{ServerError, ServerState, observability, ops_console::assets};
34
35pub fn http_router(state: ServerState) -> Result<Router, ServerError> {
42 let ops_console = assets::ops_console_router(&state.runtime_config().ops_console)?;
43 let cors = cors_layer(&state.runtime_config().cors_allowed_origins)?;
44 let metrics = state.metrics().cloned();
45 let health = state.health().cloned();
46 let mut router = workflow_router(state);
47 if let Some(metrics) = metrics {
48 router = router.merge(Router::new().route(
49 "/metrics",
50 get(observability::metrics::metrics_handler).with_state(metrics),
51 ));
52 }
53 if let Some(health) = health {
54 router = router.merge(
55 Router::new()
56 .route("/health/live", get(observability::health::live))
57 .route(
58 "/health/ready",
59 get(observability::health::ready).with_state(health),
60 ),
61 );
62 }
63 let router = router.merge(ops_console);
64 Ok(match cors {
70 Some(cors) => router.layer(cors),
71 None => router,
72 })
73}
74
75fn cors_layer(allowed_origins: &[String]) -> Result<Option<CorsLayer>, ServerError> {
91 if allowed_origins.is_empty() {
92 return Ok(None);
93 }
94 let mut origins = Vec::with_capacity(allowed_origins.len());
95 for origin in allowed_origins {
96 let value = origin
97 .parse::<HeaderValue>()
98 .map_err(|source| ServerError::Config {
99 message: format!("invalid CORS origin `{origin}`: {source}"),
100 })?;
101 origins.push(value);
102 }
103 let layer = CorsLayer::new()
104 .allow_origin(origins)
105 .allow_methods([Method::GET, Method::POST, Method::PUT, Method::OPTIONS])
106 .allow_headers([
107 header::CONTENT_TYPE,
108 header::AUTHORIZATION,
109 HeaderName::from_static("x-aion-namespaces"),
110 HeaderName::from_static("x-aion-subject"),
111 ]);
112 Ok(Some(layer))
113}
114
115async fn deploy_disabled() -> StatusCode {
118 StatusCode::NOT_FOUND
119}
120
121async fn authoring_disabled() -> StatusCode {
125 StatusCode::NOT_FOUND
126}
127
128async fn dev_disabled() -> StatusCode {
132 StatusCode::NOT_FOUND
133}
134
135pub fn workflow_router(state: ServerState) -> Router {
137 let deploy = if state.runtime_config().deploy.enabled {
145 Router::new()
146 .route(
147 "/deploy/packages",
148 post(upload_package).layer(DefaultBodyLimit::disable()),
149 )
150 .route("/deploy/versions", get(list_versions))
151 .route("/deploy/route", post(route_version))
152 .route("/deploy/unload", post(unload_version))
153 } else {
154 Router::new().route("/deploy/{*rest}", any(deploy_disabled))
155 };
156 let authoring = if state.runtime_config().authoring.gleam_path.is_some() {
163 Router::new().route("/authoring/compile", post(compile_source))
164 } else {
165 Router::new().route("/authoring/{*rest}", any(authoring_disabled))
166 };
167 let awl_documents = Router::new()
171 .route("/awl/documents", get(list_documents).post(create_document))
172 .route(
173 "/awl/documents/{*path}",
174 get(get_document).put(put_document),
175 )
176 .route("/awl/layout/{*path}", get(get_layout).put(put_layout));
177 let awl = Router::new()
178 .route("/awl/check", post(check))
179 .route("/awl/emit", post(emit))
180 .route("/awl/deploy", post(deploy_authoring))
181 .route("/awl/revisions/{hash}", get(get_revision))
182 .route("/awl/workers/availability", post(worker_availability))
183 .route("/awl/runs/{deployment_id}", get(get_run_status))
184 .route("/awl/runs/{deployment_id}/binding", post(bind_run))
185 .route("/awl/edit", post(edit))
186 .route("/awl/fmt", post(format))
187 .route("/awl/scaffold", post(scaffold))
188 .merge(awl_documents);
189 let dev = if state.runtime_config().dev.enabled {
194 Router::new()
195 .route("/dev/runs", post(dev_trigger_run))
196 .route("/dev/mocks", post(dev_register_mock))
197 .route("/dev/replay", post(dev_replay_run))
198 } else {
199 Router::new().route("/dev/{*rest}", any(dev_disabled))
200 };
201 deploy
202 .merge(authoring)
203 .merge(awl)
204 .merge(dev)
205 .route("/whoami", get(whoami))
206 .route("/namespaces", get(list_namespaces).post(post_namespace))
207 .route("/namespaces/records", get(list_namespace_records))
208 .route(
209 "/namespaces/{name}/placement",
210 axum::routing::put(set_namespace_placement),
211 )
212 .route("/workflows", get(get_workflows))
213 .route("/workflows/count", get(count_workflows))
214 .route("/workflows/start", post(start_workflow))
215 .route("/workflows/signal", post(signal_workflow))
216 .route("/workflows/query", post(query_workflow))
217 .route("/workflows/cancel", post(cancel_workflow))
218 .route("/workflows/reopen", post(reopen_workflow))
219 .route("/workflows/list", post(post_list_workflows))
220 .route("/workflows/describe", post(describe_workflow))
221 .route("/workflows/intervene", post(intervene))
222 .route("/workflows/attempts", post(list_attempts))
223 .route("/workflows/transcript", post(fetch_transcript))
224 .route("/workflows/transcripts", post(list_transcript_streams))
225 .route("/events/stream", get(subscribe_events_socket))
226 .route("/cluster/command", post(cluster_command))
227 .route("/schedules", post(create_schedule).get(list_schedules))
228 .route(
229 "/schedules/{id}",
230 get(describe_schedule)
231 .put(update_schedule)
232 .delete(delete_schedule),
233 )
234 .route("/schedules/{id}/pause", post(pause_schedule))
235 .route("/schedules/{id}/resume", post(resume_schedule))
236 .with_state(state)
237}
238
239#[cfg(test)]
240mod tests {
241 use std::{fs, sync::Arc};
242
243 use aion::EngineBuilder;
244 use aion_store::{EventStore, InMemoryStore};
245 use axum::{body, http::Request, http::StatusCode};
246 use tower::ServiceExt;
247
248 use super::super::test_support::{
249 NAMESPACE, json_request, read_json, read_text, runtime_config, server_state,
250 };
251 use super::*;
252 use crate::{
253 NamespaceResolver, StaticScheduleNamespaces, StaticWorkflowNamespaces,
254 config::{NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig},
255 };
256
257 #[tokio::test]
258 async fn ops_console_assets_serve_index_asset_and_do_not_shadow_public_api()
259 -> Result<(), Box<dyn std::error::Error>> {
260 let bundle = crate::test_support::private_tempdir()?;
261 fs::write(
262 bundle.path().join("index.html"),
263 "<!doctype html><title>Aion</title><script src=\"/app.js\"></script>",
264 )?;
265 fs::write(bundle.path().join("app.js"), "window.AION = true;")?;
266
267 let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
268 let engine = Arc::new(
269 EngineBuilder::new()
270 .store_arc(Arc::clone(&store))
271 .in_memory_visibility()
272 .scheduler_threads(1)
273 .build()
274 .await?,
275 );
276 let resolver = NamespaceResolver::from_parts(
277 NamespaceMode::SharedEngine,
278 Some(engine),
279 Arc::new(StaticWorkflowNamespaces::default()),
280 Arc::new(StaticScheduleNamespaces::default()),
281 );
282 let mut config = runtime_config();
283 config.ops_console = OpsConsoleConfig {
284 source: OpsConsoleAssetSource::FileSystem {
285 asset_path: bundle.path().to_path_buf(),
286 },
287 };
288 let router = http_router(server_state(resolver, config).await?)?;
289
290 let root = router
291 .clone()
292 .oneshot(Request::builder().uri("/").body(body::Body::empty())?)
293 .await?;
294 assert_eq!(root.status(), StatusCode::OK);
295 assert!(read_text(root).await?.contains("<title>Aion</title>"));
296
297 let asset = router
298 .clone()
299 .oneshot(
300 Request::builder()
301 .uri("/app.js")
302 .body(body::Body::empty())?,
303 )
304 .await?;
305 assert_eq!(asset.status(), StatusCode::OK);
306 assert_eq!(read_text(asset).await?, "window.AION = true;");
307
308 let spa = router
309 .clone()
310 .oneshot(
311 Request::builder()
312 .uri("/ops-console/workflows/demo")
313 .body(body::Body::empty())?,
314 )
315 .await?;
316 assert_eq!(spa.status(), StatusCode::OK);
317 assert!(read_text(spa).await?.contains("<title>Aion</title>"));
318
319 let list = serde_json::json!({
320 "namespace": NAMESPACE,
321 "filter": { "workflow_type": "nonexistent" },
322 });
323 let list_response = router
324 .oneshot(json_request("/workflows/list", &list)?)
325 .await?;
326 assert_eq!(list_response.status(), StatusCode::OK);
327 let list_body: serde_json::Value = read_json(list_response).await?;
328 assert!(
329 list_body["summaries"]
330 .as_array()
331 .ok_or("summaries missing")?
332 .is_empty()
333 );
334 Ok(())
335 }
336
337 #[tokio::test]
338 async fn cors_preflight_and_actual_request_carry_allow_headers_for_configured_origin()
339 -> Result<(), Box<dyn std::error::Error>> {
340 let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
341 let engine = Arc::new(
342 EngineBuilder::new()
343 .store_arc(Arc::clone(&store))
344 .in_memory_visibility()
345 .scheduler_threads(1)
346 .build()
347 .await?,
348 );
349 let resolver = NamespaceResolver::from_parts(
350 NamespaceMode::SharedEngine,
351 Some(engine),
352 Arc::new(StaticWorkflowNamespaces::default()),
353 Arc::new(StaticScheduleNamespaces::default()),
354 );
355 let mut config = runtime_config();
356 config.cors_allowed_origins = vec!["http://localhost:5173".to_owned()];
357 let router = http_router(server_state(resolver, config).await?)?;
358
359 let preflight = router
362 .clone()
363 .oneshot(
364 Request::builder()
365 .method("OPTIONS")
366 .uri("/workflows/list")
367 .header("origin", "http://localhost:5173")
368 .header("access-control-request-method", "POST")
369 .header("access-control-request-headers", "x-aion-namespaces")
370 .body(body::Body::empty())?,
371 )
372 .await?;
373 assert_eq!(
374 preflight
375 .headers()
376 .get("access-control-allow-origin")
377 .and_then(|value| value.to_str().ok()),
378 Some("http://localhost:5173")
379 );
380
381 let actual = router
383 .oneshot(
384 Request::builder()
385 .uri("/health/live")
386 .header("origin", "http://localhost:5173")
387 .body(body::Body::empty())?,
388 )
389 .await?;
390 assert_eq!(actual.status(), StatusCode::OK);
391 assert_eq!(
392 actual
393 .headers()
394 .get("access-control-allow-origin")
395 .and_then(|value| value.to_str().ok()),
396 Some("http://localhost:5173")
397 );
398 Ok(())
399 }
400
401 #[tokio::test]
402 async fn cors_absent_origins_install_no_layer() -> Result<(), Box<dyn std::error::Error>> {
403 let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
404 let engine = Arc::new(
405 EngineBuilder::new()
406 .store_arc(Arc::clone(&store))
407 .in_memory_visibility()
408 .scheduler_threads(1)
409 .build()
410 .await?,
411 );
412 let resolver = NamespaceResolver::from_parts(
413 NamespaceMode::SharedEngine,
414 Some(engine),
415 Arc::new(StaticWorkflowNamespaces::default()),
416 Arc::new(StaticScheduleNamespaces::default()),
417 );
418 let router = http_router(server_state(resolver, runtime_config()).await?)?;
420
421 let response = router
422 .oneshot(
423 Request::builder()
424 .uri("/health/live")
425 .header("origin", "http://localhost:5173")
426 .body(body::Body::empty())?,
427 )
428 .await?;
429 assert_eq!(response.status(), StatusCode::OK);
430 assert!(
431 response
432 .headers()
433 .get("access-control-allow-origin")
434 .is_none(),
435 "no CorsLayer must be installed when no origins are configured"
436 );
437 Ok(())
438 }
439
440 #[cfg(feature = "auth")]
441 #[tokio::test]
442 async fn worker_availability_allows_granted_jwt_namespace_and_denies_foreign_namespace()
443 -> Result<(), Box<dyn std::error::Error>> {
444 let (engine, _, _) = super::super::test_support::shared_engine().await?;
445 let resolver = NamespaceResolver::from_config(
446 NamespaceConfig {
447 mode: NamespaceMode::SharedEngine,
448 },
449 engine,
450 );
451 let router = http_router(server_state(resolver, runtime_config()).await?)?;
452
453 let allowed = router
454 .clone()
455 .oneshot(json_request(
456 "/awl/workers/availability",
457 &serde_json::json!({
458 "namespace": NAMESPACE,
459 "task_queue": "orders",
460 }),
461 )?)
462 .await?;
463 assert_eq!(allowed.status(), StatusCode::OK);
464
465 let foreign = router
466 .oneshot(json_request(
467 "/awl/workers/availability",
468 &serde_json::json!({
469 "namespace": "tenant-b",
470 "task_queue": "orders",
471 }),
472 )?)
473 .await?;
474 assert_eq!(foreign.status(), StatusCode::FORBIDDEN);
475 let body: serde_json::Value = read_json(foreign).await?;
476 assert_eq!(body["code"], "namespace_denied");
477 Ok(())
478 }
479
480 #[tokio::test]
481 async fn worker_availability_keeps_auth_off_operator_access()
482 -> Result<(), Box<dyn std::error::Error>> {
483 let (engine, _, _) = super::super::test_support::shared_engine().await?;
484 let resolver = NamespaceResolver::from_config(
485 NamespaceConfig {
486 mode: NamespaceMode::SharedEngine,
487 },
488 engine,
489 );
490 let mut config = runtime_config();
491 config.auth.enabled = false;
492 config.auth.jwks_url = None;
493 let router = http_router(crate::ServerState::from_parts(resolver, config))?;
494 let request = Request::builder()
495 .method("POST")
496 .uri("/awl/workers/availability")
497 .header("content-type", "application/json")
498 .body(body::Body::from(serde_json::to_vec(&serde_json::json!({
499 "namespace": "operator-selected-namespace",
500 "task_queue": "orders",
501 }))?))?;
502
503 let response = router.oneshot(request).await?;
504 assert_eq!(response.status(), StatusCode::OK);
505 let body: serde_json::Value = read_json(response).await?;
506 assert_eq!(body["connected_workers"], 0);
507 Ok(())
508 }
509
510 #[tokio::test]
511 async fn observability_routes_are_public_and_expose_expected_payloads()
512 -> Result<(), Box<dyn std::error::Error>> {
513 #[cfg(feature = "auth")]
517 let config = {
518 let mut config = runtime_config();
519 config.auth.jwks_url = Some(crate::auth::test_support::serve_jwks()?);
520 config
521 };
522 #[cfg(not(feature = "auth"))]
523 let config = runtime_config();
524 let router = http_router(
525 crate::ServerState::build_with_store(InMemoryStore::default(), config).await?,
526 )?;
527
528 let metrics_response = router
529 .clone()
530 .oneshot(
531 Request::builder()
532 .uri("/metrics")
533 .body(body::Body::empty())?,
534 )
535 .await?;
536 assert_eq!(metrics_response.status(), StatusCode::OK);
537 assert_eq!(
538 metrics_response
539 .headers()
540 .get(axum::http::header::CONTENT_TYPE)
541 .and_then(|value| value.to_str().ok()),
542 Some("text/plain; version=0.0.4; charset=utf-8")
543 );
544 let metrics_body = read_text(metrics_response).await?;
545 assert!(metrics_body.contains("# HELP aion_workflows_started_total"));
546 assert!(metrics_body.contains("# TYPE aion_workflows_started_total counter"));
547 assert!(metrics_body.contains("# HELP aion_activity_duration_seconds"));
548 assert!(metrics_body.contains("# TYPE aion_activity_duration_seconds histogram"));
549 assert!(metrics_body.contains("aion_activity_duration_seconds_bucket"));
550 assert!(metrics_body.contains("aion_store_operation_duration_seconds_bucket"));
551
552 let live_response = router
553 .clone()
554 .oneshot(
555 Request::builder()
556 .uri("/health/live")
557 .body(body::Body::empty())?,
558 )
559 .await?;
560 assert_eq!(live_response.status(), StatusCode::OK);
561
562 let ready_response = router
563 .oneshot(
564 Request::builder()
565 .uri("/health/ready")
566 .body(body::Body::empty())?,
567 )
568 .await?;
569 assert_eq!(ready_response.status(), StatusCode::OK);
570 Ok(())
571 }
572}