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::assistant::assistant_descriptor;
12use super::assistant_sessions::{
13 cancel_assistant_turn, create_assistant_session, current_assistant_session,
14 delete_assistant_session, list_assistant_sessions, push_assistant_context,
15 read_assistant_session, resume_assistant_session, set_assistant_config_option,
16 submit_assistant_turn,
17};
18use super::assistant_socket::assistant_session_socket;
19use super::authoring::compile_source;
20use super::awl::{
21 bind_run, check, create_document, deploy_authoring, doc as awl_doc, edit, emit, format,
22 get_document, get_layout, get_revision, get_run_status, list_documents, put_document,
23 put_layout, scaffold, worker_availability,
24};
25use super::awl_deployed::{get_deployed_doc, get_deployed_document, list_deployed};
26use super::awl_graph::{deployed_graph_svg, workspace_graph_svg};
27use super::build::build_identity;
28use super::changelog::changelog;
29use super::children::list_children;
30use super::cluster_command::cluster_command;
31use super::deploy::{list_versions, route_version, unload_version, upload_package};
32use super::describe_live::describe_live;
33use super::dev_ui::{dev_register_mock, dev_replay_run, dev_trigger_run};
34use super::events::subscribe_events_socket;
35use super::history::{fetch_event, fetch_history};
36use super::intervene::{intervene, list_attempts};
37use super::managed_workers::{
38 list_managed_workers, restart_managed_worker, start_managed_worker, stop_managed_worker,
39};
40use super::outbox::list_dead_letters;
41use super::queues::list_unserved_queues;
42use super::schedules::{
43 create_schedule, delete_schedule, describe_schedule, list_schedules, pause_schedule,
44 resume_schedule, update_schedule,
45};
46use super::transcripts::{fetch_transcript, list_transcript_streams};
47use super::unrecoverable::list_unrecoverable_runs;
48use super::update_status::update_status;
49use super::whoami::whoami;
50use super::worker_deployments::{
51 delete_worker_deployment, get_worker_deployment, list_worker_deployments,
52 put_worker_deployment, set_worker_deployment_desired_state,
53};
54use super::workers::{drain_worker, stop_worker};
55use super::workflow_document::get_run_document;
56use super::workflows::{
57 cancel_workflow, describe_workflow, list_namespace_records, list_namespaces,
58 post_list_workflows, post_namespace, query_workflow, rename_workflow, reopen_workflow,
59 retire_workloop, set_namespace_placement, signal_workflow, start_workflow,
60};
61use crate::assistant::mcp::{AssistantMcpRuntime, assistant_mcp_router};
62use crate::mcp::{McpRuntime, mcp_disabled_router, mcp_router};
63use crate::{ServerError, ServerState, observability, ops_console::assets};
64
65pub fn http_router(state: ServerState) -> Result<Router, ServerError> {
74 let ops_console = assets::ops_console_router(&state.runtime_config().ops_console)?;
75 let cors = cors_layer(&state.runtime_config().cors_allowed_origins)?;
76 let metrics = state.metrics().cloned();
77 let health = state.health().cloned();
78 let mut router = workflow_router(state.clone());
79 let assistant_mcp = assistant_mcp_family(&state)?;
88 router = router
89 .merge(assistant_mcp.with_state(state.clone()))
90 .merge(mcp_family(&state)?.with_state(state));
91 if let Some(metrics) = metrics {
92 router = router.merge(Router::new().route(
93 "/metrics",
94 get(observability::metrics::metrics_handler).with_state(metrics),
95 ));
96 }
97 if let Some(health) = health {
98 router = router.merge(
99 Router::new()
100 .route("/health/live", get(observability::health::live))
101 .route(
102 "/health/ready",
103 get(observability::health::ready).with_state(health),
104 ),
105 );
106 }
107 let router = router.merge(ops_console);
108 Ok(match cors {
114 Some(cors) => router.layer(cors),
115 None => router,
116 })
117}
118
119fn cors_layer(allowed_origins: &[String]) -> Result<Option<CorsLayer>, ServerError> {
135 if allowed_origins.is_empty() {
136 return Ok(None);
137 }
138 let mut origins = Vec::with_capacity(allowed_origins.len());
139 for origin in allowed_origins {
140 let value = origin
141 .parse::<HeaderValue>()
142 .map_err(|source| ServerError::Config {
143 message: format!("invalid CORS origin `{origin}`: {source}"),
144 })?;
145 origins.push(value);
146 }
147 let layer = CorsLayer::new()
148 .allow_origin(origins)
149 .allow_methods([
150 Method::GET,
151 Method::POST,
152 Method::PUT,
153 Method::DELETE,
154 Method::OPTIONS,
155 ])
156 .allow_headers([
157 header::CONTENT_TYPE,
158 header::AUTHORIZATION,
159 HeaderName::from_static("x-aion-namespaces"),
160 HeaderName::from_static("x-aion-subject"),
161 ]);
162 Ok(Some(layer))
163}
164
165async fn deploy_disabled() -> StatusCode {
168 StatusCode::NOT_FOUND
169}
170
171async fn authoring_disabled() -> StatusCode {
175 StatusCode::NOT_FOUND
176}
177
178async fn dev_disabled() -> StatusCode {
182 StatusCode::NOT_FOUND
183}
184
185async fn deploy_surface_disabled() -> StatusCode {
190 StatusCode::NOT_FOUND
191}
192
193fn deployed_awl_router(deploy_enabled: bool) -> Router<ServerState> {
201 if deploy_enabled {
202 Router::new()
203 .route("/awl/deployed", get(list_deployed))
204 .route(
205 "/awl/deployed/{workflow_type}/{content_hash}",
206 get(get_deployed_document),
207 )
208 .route(
209 "/awl/deployed/{workflow_type}/{content_hash}/doc",
210 get(get_deployed_doc),
211 )
212 .route(
213 "/awl/deployed/{workflow_type}/{content_hash}/doc/graph.svg",
214 get(deployed_graph_svg),
215 )
216 } else {
217 Router::new()
218 .route("/awl/deployed", any(deploy_surface_disabled))
219 .route("/awl/deployed/{*rest}", any(deploy_surface_disabled))
220 }
221}
222
223fn worker_deployment_router(deploy_enabled: bool) -> Router<ServerState> {
224 if deploy_enabled {
225 Router::new()
226 .route("/worker-deployments", get(list_worker_deployments))
227 .route(
228 "/worker-deployments/{name}",
229 get(get_worker_deployment)
230 .put(put_worker_deployment)
231 .delete(delete_worker_deployment),
232 )
233 .route(
234 "/worker-deployments/{name}/desired-state",
235 post(set_worker_deployment_desired_state),
236 )
237 } else {
238 Router::new()
239 .route("/worker-deployments", any(deploy_surface_disabled))
240 .route("/worker-deployments/{*rest}", any(deploy_surface_disabled))
241 }
242}
243
244fn managed_worker_router(deploy_enabled: bool) -> Router<ServerState> {
258 let read = Router::new().route("/workers/managed", get(list_managed_workers));
259 if deploy_enabled {
260 read.route("/workers/managed/{name}/start", post(start_managed_worker))
261 .route("/workers/managed/{name}/stop", post(stop_managed_worker))
262 .route(
263 "/workers/managed/{name}/restart",
264 post(restart_managed_worker),
265 )
266 } else {
267 read.route("/workers/managed/{*rest}", any(deploy_surface_disabled))
268 }
269}
270
271fn mcp_family(state: &ServerState) -> Result<Router<ServerState>, ServerError> {
287 if !state.runtime_config().mcp.enabled {
288 return Ok(mcp_disabled_router());
289 }
290 let runtime = McpRuntime::build(state.clone(), &state.runtime_config().mcp)?;
291 Ok(mcp_router(std::sync::Arc::new(runtime)))
292}
293
294fn assistant_mcp_family(state: &ServerState) -> Result<Router<ServerState>, ServerError> {
305 let runtime = AssistantMcpRuntime::build(state.clone(), &state.runtime_config().mcp)?;
306 Ok(assistant_mcp_router(std::sync::Arc::new(runtime)))
307}
308
309fn assistant_session_router() -> Router<ServerState> {
323 Router::new()
324 .route(
325 "/assistant/sessions",
326 get(list_assistant_sessions).post(create_assistant_session),
327 )
328 .route(
333 "/assistant/sessions/current",
334 get(current_assistant_session),
335 )
336 .route(
337 "/assistant/sessions/{id}",
338 get(read_assistant_session).delete(delete_assistant_session),
339 )
340 .route(
341 "/assistant/sessions/{id}/turns",
342 post(submit_assistant_turn),
343 )
344 .route(
345 "/assistant/sessions/{id}/context",
346 axum::routing::put(push_assistant_context),
347 )
348 .route(
349 "/assistant/sessions/{id}/cancel",
350 post(cancel_assistant_turn),
351 )
352 .route(
353 "/assistant/sessions/{id}/config",
354 axum::routing::put(set_assistant_config_option),
355 )
356 .route(
357 "/assistant/sessions/{id}/resume",
358 post(resume_assistant_session),
359 )
360 .route(
361 "/assistant/sessions/{id}/events",
362 get(assistant_session_socket),
363 )
364}
365
366fn dev_router(enabled: bool) -> Router<ServerState> {
368 if enabled {
369 Router::new()
370 .route("/dev/runs", post(dev_trigger_run))
371 .route("/dev/mocks", post(dev_register_mock))
372 .route("/dev/replay", post(dev_replay_run))
373 } else {
374 Router::new().route("/dev/{*rest}", any(dev_disabled))
375 }
376}
377
378fn deploy_routes(state: &ServerState) -> Router<ServerState> {
386 if state.runtime_config().deploy.enabled {
387 Router::new()
388 .route(
389 "/deploy/packages",
390 post(upload_package).layer(DefaultBodyLimit::disable()),
391 )
392 .route("/deploy/versions", get(list_versions))
393 .route("/deploy/route", post(route_version))
394 .route("/deploy/unload", post(unload_version))
395 } else {
396 Router::new().route("/deploy/{*rest}", any(deploy_disabled))
397 }
398}
399
400pub fn workflow_router(state: ServerState) -> Router {
403 let deploy = deploy_routes(&state);
411 let authoring = if state.runtime_config().authoring.gleam_path.is_some() {
418 Router::new().route("/authoring/compile", post(compile_source))
419 } else {
420 Router::new().route("/authoring/{*rest}", any(authoring_disabled))
421 };
422 let awl_documents = Router::new()
426 .route("/awl/documents", get(list_documents).post(create_document))
427 .route(
428 "/awl/documents/{*path}",
429 get(get_document).put(put_document),
430 )
431 .route("/awl/layout/{*path}", get(get_layout).put(put_layout));
432 let awl_deployed = deployed_awl_router(state.runtime_config().deploy.enabled);
433 let awl = Router::new()
434 .route("/awl/check", post(check))
435 .route("/awl/doc", post(awl_doc))
436 .route("/awl/doc/graph.svg", post(workspace_graph_svg))
437 .route("/awl/emit", post(emit))
438 .route("/awl/deploy", post(deploy_authoring))
439 .route("/awl/revisions/{hash}", get(get_revision))
440 .route("/awl/workers/availability", post(worker_availability))
441 .route("/awl/runs/{deployment_id}", get(get_run_status))
442 .route("/awl/runs/{deployment_id}/binding", post(bind_run))
443 .route("/awl/edit", post(edit))
444 .route("/awl/fmt", post(format))
445 .route("/awl/scaffold", post(scaffold))
446 .merge(awl_documents)
447 .merge(awl_deployed);
448 let dev = dev_router(state.runtime_config().dev.enabled);
453 deploy
454 .merge(authoring)
455 .merge(awl)
456 .merge(dev)
457 .merge(worker_deployment_router(
458 state.runtime_config().deploy.enabled,
459 ))
460 .route("/whoami", get(whoami))
461 .route("/build", get(build_identity))
462 .route("/update-status", get(update_status))
466 .route("/changelog", get(changelog))
469 .route("/assistant", get(assistant_descriptor))
473 .merge(assistant_session_router())
474 .route("/namespaces", get(list_namespaces).post(post_namespace))
475 .route("/namespaces/records", get(list_namespace_records))
476 .route(
477 "/namespaces/{name}/placement",
478 axum::routing::put(set_namespace_placement),
479 )
480 .route("/workflows/start", post(start_workflow))
481 .route("/workflows/signal", post(signal_workflow))
482 .route("/workflows/query", post(query_workflow))
483 .route("/workflows/cancel", post(cancel_workflow))
484 .route("/workflows/retire", post(retire_workloop))
485 .route("/workflows/rename", post(rename_workflow))
486 .route("/workflows/reopen", post(reopen_workflow))
487 .route("/workflows/list", post(post_list_workflows))
488 .route("/workflows/describe", post(describe_workflow))
489 .route("/workflows/describe-live", post(describe_live))
490 .route("/workflows/children", post(list_children))
491 .route("/workflows/history", post(fetch_history))
492 .route("/workflows/event", post(fetch_event))
493 .route("/workflows/intervene", post(intervene))
494 .route("/workflows/attempts", post(list_attempts))
495 .route("/workflows/transcript", post(fetch_transcript))
496 .route("/workflows/transcripts", post(list_transcript_streams))
497 .route("/workflows/unrecoverable", get(list_unrecoverable_runs))
498 .route(
501 "/workflows/{workflow_id}/document/{content_hash}",
502 get(get_run_document),
503 )
504 .route("/events/stream", get(subscribe_events_socket))
505 .route("/cluster/command", post(cluster_command))
506 .route("/outbox/dead-letters", post(list_dead_letters))
507 .route("/queues/unserved", get(list_unserved_queues))
508 .merge(managed_worker_router(state.runtime_config().deploy.enabled))
509 .route("/workers/{worker_id}/drain", post(drain_worker))
510 .route("/workers/{worker_id}/stop", post(stop_worker))
511 .route("/schedules", post(create_schedule).get(list_schedules))
512 .route(
513 "/schedules/{id}",
514 get(describe_schedule)
515 .put(update_schedule)
516 .delete(delete_schedule),
517 )
518 .route("/schedules/{id}/pause", post(pause_schedule))
519 .route("/schedules/{id}/resume", post(resume_schedule))
520 .with_state(state)
521}
522
523#[cfg(test)]
524#[path = "router_tests.rs"]
525mod tests;