use axum::{
Router,
extract::DefaultBodyLimit,
http::{HeaderName, HeaderValue, Method, StatusCode, header},
routing::{any, get, post},
};
use tower_http::cors::CorsLayer;
use super::assistant::{assistant_descriptor, assistant_document};
use super::authoring::compile_source;
use super::awl::{
bind_run, check, create_document, deploy_authoring, edit, emit, format, get_document,
get_layout, get_revision, get_run_status, list_documents, put_document, put_layout, scaffold,
worker_availability,
};
use super::awl_deployed::{get_deployed_document, list_deployed};
use super::build::build_identity;
use super::changelog::changelog;
use super::cluster_command::cluster_command;
use super::deploy::{list_versions, route_version, unload_version, upload_package};
use super::describe_live::describe_live;
use super::dev_ui::{dev_register_mock, dev_replay_run, dev_trigger_run};
use super::events::subscribe_events_socket;
use super::history::{fetch_event, fetch_history};
use super::intervene::{intervene, list_attempts};
use super::managed_workers::{
list_managed_workers, restart_managed_worker, start_managed_worker, stop_managed_worker,
};
use super::outbox::list_dead_letters;
use super::queues::list_unserved_queues;
use super::schedules::{
create_schedule, delete_schedule, describe_schedule, list_schedules, pause_schedule,
resume_schedule, update_schedule,
};
use super::transcripts::{fetch_transcript, list_transcript_streams};
use super::unrecoverable::list_unrecoverable_runs;
use super::update_status::update_status;
use super::whoami::whoami;
use super::worker_deployments::{
delete_worker_deployment, get_worker_deployment, list_worker_deployments,
put_worker_deployment, set_worker_deployment_desired_state,
};
use super::workers::{drain_worker, stop_worker};
use super::workflows::{
cancel_workflow, count_workflows, describe_workflow, get_workflows, list_namespace_records,
list_namespaces, post_list_workflows, post_namespace, query_workflow, reopen_workflow,
set_namespace_placement, signal_workflow, start_workflow,
};
use crate::mcp::{McpRuntime, mcp_disabled_router, mcp_router};
use crate::{ServerError, ServerState, observability, ops_console::assets};
pub fn http_router(state: ServerState) -> Result<Router, ServerError> {
let ops_console = assets::ops_console_router(&state.runtime_config().ops_console)?;
let cors = cors_layer(&state.runtime_config().cors_allowed_origins)?;
let metrics = state.metrics().cloned();
let health = state.health().cloned();
let mut router = workflow_router(state.clone());
router = router.merge(mcp_family(&state)?.with_state(state));
if let Some(metrics) = metrics {
router = router.merge(Router::new().route(
"/metrics",
get(observability::metrics::metrics_handler).with_state(metrics),
));
}
if let Some(health) = health {
router = router.merge(
Router::new()
.route("/health/live", get(observability::health::live))
.route(
"/health/ready",
get(observability::health::ready).with_state(health),
),
);
}
let router = router.merge(ops_console);
Ok(match cors {
Some(cors) => router.layer(cors),
None => router,
})
}
fn cors_layer(allowed_origins: &[String]) -> Result<Option<CorsLayer>, ServerError> {
if allowed_origins.is_empty() {
return Ok(None);
}
let mut origins = Vec::with_capacity(allowed_origins.len());
for origin in allowed_origins {
let value = origin
.parse::<HeaderValue>()
.map_err(|source| ServerError::Config {
message: format!("invalid CORS origin `{origin}`: {source}"),
})?;
origins.push(value);
}
let layer = CorsLayer::new()
.allow_origin(origins)
.allow_methods([
Method::GET,
Method::POST,
Method::PUT,
Method::DELETE,
Method::OPTIONS,
])
.allow_headers([
header::CONTENT_TYPE,
header::AUTHORIZATION,
HeaderName::from_static("x-aion-namespaces"),
HeaderName::from_static("x-aion-subject"),
]);
Ok(Some(layer))
}
async fn deploy_disabled() -> StatusCode {
StatusCode::NOT_FOUND
}
async fn authoring_disabled() -> StatusCode {
StatusCode::NOT_FOUND
}
async fn dev_disabled() -> StatusCode {
StatusCode::NOT_FOUND
}
async fn deploy_surface_disabled() -> StatusCode {
StatusCode::NOT_FOUND
}
fn deployed_awl_router(deploy_enabled: bool) -> Router<ServerState> {
if deploy_enabled {
Router::new()
.route("/awl/deployed", get(list_deployed))
.route(
"/awl/deployed/{workflow_type}/{content_hash}",
get(get_deployed_document),
)
} else {
Router::new()
.route("/awl/deployed", any(deploy_surface_disabled))
.route("/awl/deployed/{*rest}", any(deploy_surface_disabled))
}
}
fn worker_deployment_router(deploy_enabled: bool) -> Router<ServerState> {
if deploy_enabled {
Router::new()
.route("/worker-deployments", get(list_worker_deployments))
.route(
"/worker-deployments/{name}",
get(get_worker_deployment)
.put(put_worker_deployment)
.delete(delete_worker_deployment),
)
.route(
"/worker-deployments/{name}/desired-state",
post(set_worker_deployment_desired_state),
)
} else {
Router::new()
.route("/worker-deployments", any(deploy_surface_disabled))
.route("/worker-deployments/{*rest}", any(deploy_surface_disabled))
}
}
fn managed_worker_router(deploy_enabled: bool) -> Router<ServerState> {
let read = Router::new().route("/workers/managed", get(list_managed_workers));
if deploy_enabled {
read.route("/workers/managed/{name}/start", post(start_managed_worker))
.route("/workers/managed/{name}/stop", post(stop_managed_worker))
.route(
"/workers/managed/{name}/restart",
post(restart_managed_worker),
)
} else {
read.route("/workers/managed/{*rest}", any(deploy_surface_disabled))
}
}
fn mcp_family(state: &ServerState) -> Result<Router<ServerState>, ServerError> {
if !state.runtime_config().mcp.enabled {
return Ok(mcp_disabled_router());
}
let runtime = McpRuntime::build(state.clone(), &state.runtime_config().mcp)?;
Ok(mcp_router(std::sync::Arc::new(runtime)))
}
pub fn workflow_router(state: ServerState) -> Router {
let deploy = if state.runtime_config().deploy.enabled {
Router::new()
.route(
"/deploy/packages",
post(upload_package).layer(DefaultBodyLimit::disable()),
)
.route("/deploy/versions", get(list_versions))
.route("/deploy/route", post(route_version))
.route("/deploy/unload", post(unload_version))
} else {
Router::new().route("/deploy/{*rest}", any(deploy_disabled))
};
let authoring = if state.runtime_config().authoring.gleam_path.is_some() {
Router::new().route("/authoring/compile", post(compile_source))
} else {
Router::new().route("/authoring/{*rest}", any(authoring_disabled))
};
let awl_documents = Router::new()
.route("/awl/documents", get(list_documents).post(create_document))
.route(
"/awl/documents/{*path}",
get(get_document).put(put_document),
)
.route("/awl/layout/{*path}", get(get_layout).put(put_layout));
let awl_deployed = deployed_awl_router(state.runtime_config().deploy.enabled);
let awl = Router::new()
.route("/awl/check", post(check))
.route("/awl/emit", post(emit))
.route("/awl/deploy", post(deploy_authoring))
.route("/awl/revisions/{hash}", get(get_revision))
.route("/awl/workers/availability", post(worker_availability))
.route("/awl/runs/{deployment_id}", get(get_run_status))
.route("/awl/runs/{deployment_id}/binding", post(bind_run))
.route("/awl/edit", post(edit))
.route("/awl/fmt", post(format))
.route("/awl/scaffold", post(scaffold))
.merge(awl_documents)
.merge(awl_deployed);
let dev = if state.runtime_config().dev.enabled {
Router::new()
.route("/dev/runs", post(dev_trigger_run))
.route("/dev/mocks", post(dev_register_mock))
.route("/dev/replay", post(dev_replay_run))
} else {
Router::new().route("/dev/{*rest}", any(dev_disabled))
};
deploy
.merge(authoring)
.merge(awl)
.merge(dev)
.merge(worker_deployment_router(
state.runtime_config().deploy.enabled,
))
.route("/whoami", get(whoami))
.route("/build", get(build_identity))
.route("/update-status", get(update_status))
.route("/changelog", get(changelog))
.route("/assistant", get(assistant_descriptor))
.route("/assistant/document", get(assistant_document))
.route("/namespaces", get(list_namespaces).post(post_namespace))
.route("/namespaces/records", get(list_namespace_records))
.route(
"/namespaces/{name}/placement",
axum::routing::put(set_namespace_placement),
)
.route("/workflows", get(get_workflows))
.route("/workflows/count", get(count_workflows))
.route("/workflows/start", post(start_workflow))
.route("/workflows/signal", post(signal_workflow))
.route("/workflows/query", post(query_workflow))
.route("/workflows/cancel", post(cancel_workflow))
.route("/workflows/reopen", post(reopen_workflow))
.route("/workflows/list", post(post_list_workflows))
.route("/workflows/describe", post(describe_workflow))
.route("/workflows/describe-live", post(describe_live))
.route("/workflows/history", post(fetch_history))
.route("/workflows/event", post(fetch_event))
.route("/workflows/intervene", post(intervene))
.route("/workflows/attempts", post(list_attempts))
.route("/workflows/transcript", post(fetch_transcript))
.route("/workflows/transcripts", post(list_transcript_streams))
.route("/workflows/unrecoverable", get(list_unrecoverable_runs))
.route("/events/stream", get(subscribe_events_socket))
.route("/cluster/command", post(cluster_command))
.route("/outbox/dead-letters", post(list_dead_letters))
.route("/queues/unserved", get(list_unserved_queues))
.merge(managed_worker_router(state.runtime_config().deploy.enabled))
.route("/workers/{worker_id}/drain", post(drain_worker))
.route("/workers/{worker_id}/stop", post(stop_worker))
.route("/schedules", post(create_schedule).get(list_schedules))
.route(
"/schedules/{id}",
get(describe_schedule)
.put(update_schedule)
.delete(delete_schedule),
)
.route("/schedules/{id}/pause", post(pause_schedule))
.route("/schedules/{id}/resume", post(resume_schedule))
.with_state(state)
}
#[cfg(test)]
mod tests {
use std::{fs, sync::Arc};
use aion::EngineBuilder;
use aion_store::{EventStore, InMemoryStore};
use axum::{body, http::Request, http::StatusCode};
use tower::ServiceExt;
use super::super::test_support::{
NAMESPACE, json_request, read_json, read_text, runtime_config, server_state,
};
use super::*;
use crate::{
NamespaceResolver, StaticScheduleNamespaces, StaticWorkflowNamespaces,
config::{NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig},
};
#[tokio::test]
async fn ops_console_assets_serve_index_asset_and_do_not_shadow_public_api()
-> Result<(), Box<dyn std::error::Error>> {
let bundle = crate::test_support::private_tempdir()?;
fs::write(
bundle.path().join("index.html"),
"<!doctype html><title>Aion</title><script src=\"/app.js\"></script>",
)?;
fs::write(bundle.path().join("app.js"), "window.AION = true;")?;
let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
let engine = Arc::new(
EngineBuilder::new()
.store_arc(Arc::clone(&store))
.in_memory_visibility()
.scheduler_threads(1)
.build()
.await?,
);
let resolver = NamespaceResolver::from_parts(
NamespaceMode::SharedEngine,
Some(engine),
Arc::new(StaticWorkflowNamespaces::default()),
Arc::new(StaticScheduleNamespaces::default()),
);
let mut config = runtime_config();
config.ops_console = OpsConsoleConfig {
source: OpsConsoleAssetSource::FileSystem {
asset_path: bundle.path().to_path_buf(),
},
};
let router = http_router(server_state(resolver, config).await?)?;
let root = router
.clone()
.oneshot(Request::builder().uri("/").body(body::Body::empty())?)
.await?;
assert_eq!(root.status(), StatusCode::OK);
assert!(read_text(root).await?.contains("<title>Aion</title>"));
let asset = router
.clone()
.oneshot(
Request::builder()
.uri("/app.js")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(asset.status(), StatusCode::OK);
assert_eq!(read_text(asset).await?, "window.AION = true;");
let spa = router
.clone()
.oneshot(
Request::builder()
.uri("/ops-console/workflows/demo")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(spa.status(), StatusCode::OK);
assert!(read_text(spa).await?.contains("<title>Aion</title>"));
let list = serde_json::json!({
"namespace": NAMESPACE,
"filter": { "workflow_type": "nonexistent" },
});
let list_response = router
.oneshot(json_request("/workflows/list", &list)?)
.await?;
assert_eq!(list_response.status(), StatusCode::OK);
let list_body: serde_json::Value = read_json(list_response).await?;
assert!(
list_body["summaries"]
.as_array()
.ok_or("summaries missing")?
.is_empty()
);
Ok(())
}
#[tokio::test]
async fn cors_preflight_and_actual_request_carry_allow_headers_for_configured_origin()
-> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
let engine = Arc::new(
EngineBuilder::new()
.store_arc(Arc::clone(&store))
.in_memory_visibility()
.scheduler_threads(1)
.build()
.await?,
);
let resolver = NamespaceResolver::from_parts(
NamespaceMode::SharedEngine,
Some(engine),
Arc::new(StaticWorkflowNamespaces::default()),
Arc::new(StaticScheduleNamespaces::default()),
);
let mut config = runtime_config();
config.cors_allowed_origins = vec!["http://localhost:5173".to_owned()];
let router = http_router(server_state(resolver, config).await?)?;
let preflight = router
.clone()
.oneshot(
Request::builder()
.method("OPTIONS")
.uri("/workflows/list")
.header("origin", "http://localhost:5173")
.header("access-control-request-method", "POST")
.header("access-control-request-headers", "x-aion-namespaces")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(
preflight
.headers()
.get("access-control-allow-origin")
.and_then(|value| value.to_str().ok()),
Some("http://localhost:5173")
);
let actual = router
.oneshot(
Request::builder()
.uri("/health/live")
.header("origin", "http://localhost:5173")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(actual.status(), StatusCode::OK);
assert_eq!(
actual
.headers()
.get("access-control-allow-origin")
.and_then(|value| value.to_str().ok()),
Some("http://localhost:5173")
);
Ok(())
}
#[tokio::test]
async fn cors_absent_origins_install_no_layer() -> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
let engine = Arc::new(
EngineBuilder::new()
.store_arc(Arc::clone(&store))
.in_memory_visibility()
.scheduler_threads(1)
.build()
.await?,
);
let resolver = NamespaceResolver::from_parts(
NamespaceMode::SharedEngine,
Some(engine),
Arc::new(StaticWorkflowNamespaces::default()),
Arc::new(StaticScheduleNamespaces::default()),
);
let router = http_router(server_state(resolver, runtime_config()).await?)?;
let response = router
.oneshot(
Request::builder()
.uri("/health/live")
.header("origin", "http://localhost:5173")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(response.status(), StatusCode::OK);
assert!(
response
.headers()
.get("access-control-allow-origin")
.is_none(),
"no CorsLayer must be installed when no origins are configured"
);
Ok(())
}
#[cfg(feature = "auth")]
#[tokio::test]
async fn worker_availability_allows_granted_jwt_namespace_and_denies_foreign_namespace()
-> Result<(), Box<dyn std::error::Error>> {
let (engine, _, _) = super::super::test_support::shared_engine().await?;
let resolver = NamespaceResolver::from_config(
NamespaceConfig {
mode: NamespaceMode::SharedEngine,
},
engine,
);
let router = http_router(server_state(resolver, runtime_config()).await?)?;
let allowed = router
.clone()
.oneshot(json_request(
"/awl/workers/availability",
&serde_json::json!({
"namespace": NAMESPACE,
"task_queue": "orders",
}),
)?)
.await?;
assert_eq!(allowed.status(), StatusCode::OK);
let foreign = router
.oneshot(json_request(
"/awl/workers/availability",
&serde_json::json!({
"namespace": "tenant-b",
"task_queue": "orders",
}),
)?)
.await?;
assert_eq!(foreign.status(), StatusCode::FORBIDDEN);
let body: serde_json::Value = read_json(foreign).await?;
assert_eq!(body["code"], "namespace_denied");
Ok(())
}
#[tokio::test]
async fn worker_availability_keeps_auth_off_operator_access()
-> Result<(), Box<dyn std::error::Error>> {
let (engine, _, _) = super::super::test_support::shared_engine().await?;
let resolver = NamespaceResolver::from_config(
NamespaceConfig {
mode: NamespaceMode::SharedEngine,
},
engine,
);
let mut config = runtime_config();
config.auth.enabled = false;
config.auth.jwks_url = None;
let router = http_router(crate::ServerState::from_parts(resolver, config))?;
let request = Request::builder()
.method("POST")
.uri("/awl/workers/availability")
.header("content-type", "application/json")
.body(body::Body::from(serde_json::to_vec(&serde_json::json!({
"namespace": "operator-selected-namespace",
"task_queue": "orders",
}))?))?;
let response = router.oneshot(request).await?;
assert_eq!(response.status(), StatusCode::OK);
let body: serde_json::Value = read_json(response).await?;
assert_eq!(body["connected_workers"], 0);
Ok(())
}
#[tokio::test]
async fn observability_routes_are_public_and_expose_expected_payloads()
-> Result<(), Box<dyn std::error::Error>> {
#[cfg(feature = "auth")]
let config = {
let mut config = runtime_config();
config.auth.jwks_url = Some(crate::auth::test_support::serve_jwks()?);
config
};
#[cfg(not(feature = "auth"))]
let config = runtime_config();
let router = http_router(
crate::ServerState::build_with_store(InMemoryStore::default(), config).await?,
)?;
let metrics_response = router
.clone()
.oneshot(
Request::builder()
.uri("/metrics")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(metrics_response.status(), StatusCode::OK);
assert_eq!(
metrics_response
.headers()
.get(axum::http::header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok()),
Some("text/plain; version=0.0.4; charset=utf-8")
);
let metrics_body = read_text(metrics_response).await?;
assert!(metrics_body.contains("# HELP aion_workflows_started_total"));
assert!(metrics_body.contains("# TYPE aion_workflows_started_total counter"));
assert!(metrics_body.contains("# HELP aion_activity_duration_seconds"));
assert!(metrics_body.contains("# TYPE aion_activity_duration_seconds histogram"));
assert!(metrics_body.contains("aion_activity_duration_seconds_bucket"));
assert!(metrics_body.contains("aion_store_operation_duration_seconds_bucket"));
let live_response = router
.clone()
.oneshot(
Request::builder()
.uri("/health/live")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(live_response.status(), StatusCode::OK);
let ready_response = router
.oneshot(
Request::builder()
.uri("/health/ready")
.body(body::Body::empty())?,
)
.await?;
assert_eq!(ready_response.status(), StatusCode::OK);
Ok(())
}
}