pub mod api_docs;
mod chat_handlers;
mod hitl_handlers;
mod registration_handlers;
mod registry_handlers;
mod status_handlers;
use super::SharedAgentStatus;
use super::agent_events::AgentEventStore;
use crate::agents::{AgentConfig, ChatCapable};
use crate::orchestrator_registry::OrchestratorRegistry;
use crate::workers::buffer::ResponseBuffer;
use axum::{
Router,
extract::{FromRef, Path, Request, State},
http::{HeaderMap, StatusCode, header},
middleware::{self, Next},
response::{Html, IntoResponse, Json, Response},
routing::{get, post, put},
};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::net::SocketAddr;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64};
use tokio::sync::RwLock;
use tracing::{error, info};
use utoipa::ToSchema;
#[derive(Clone)]
pub(crate) struct MultiAppState {
statuses: HashMap<String, SharedAgentStatus>,
chat_agents: HashMap<String, Arc<dyn ChatCapable>>,
configs: HashMap<String, Arc<RwLock<AgentConfig>>>,
buffers: HashMap<String, Arc<ResponseBuffer>>,
pause_handles: HashMap<String, Arc<AtomicBool>>,
event_stores: HashMap<String, AgentEventStore>,
orchestrator_registry: Option<OrchestratorRegistry>,
base_hold_secs: Arc<AtomicU64>,
response_sla_secs: Arc<AtomicU64>,
buffer_floor_pct: Arc<AtomicU64>,
before_release_middleware: Option<Arc<crate::middleware::pipeline::MiddlewarePipeline>>,
auth_token: Option<Arc<str>>,
}
#[derive(Serialize, ToSchema)]
pub(super) struct GlobalConfig {
base_hold_secs: u64,
response_sla_secs: u64,
buffer_floor_pct: u64,
}
#[derive(Deserialize, ToSchema)]
pub(super) struct GlobalConfigUpdate {
base_hold_secs: Option<u64>,
response_sla_secs: Option<u64>,
buffer_floor_pct: Option<u64>,
}
pub struct MultiAgentStatusServer;
fn resolve_dashboard_bind(raw: Option<&str>) -> std::net::IpAddr {
raw.and_then(|s| s.parse().ok())
.unwrap_or(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST))
}
fn resolve_dashboard_token(raw: Option<String>) -> Option<Arc<str>> {
raw.filter(|s| !s.is_empty()).map(Arc::from)
}
fn refuse_unauthenticated_exposure(ip: &std::net::IpAddr, auth_enabled: bool) -> bool {
!ip.is_loopback() && !auth_enabled
}
fn ct_eq(a: &[u8], b: &[u8]) -> bool {
if a.len() != b.len() {
return false;
}
let mut diff = 0u8;
for (x, y) in a.iter().zip(b.iter()) {
diff |= x ^ y;
}
diff == 0
}
#[derive(Clone)]
struct DashAuth {
token: Option<Arc<str>>,
}
impl FromRef<MultiAppState> for DashAuth {
fn from_ref(state: &MultiAppState) -> Self {
DashAuth {
token: state.auth_token.clone(),
}
}
}
impl DashAuth {
fn is_authorized(&self, headers: &HeaderMap) -> bool {
let Some(expected) = &self.token else {
return true;
};
headers
.get(header::AUTHORIZATION)
.and_then(|value| value.to_str().ok())
.and_then(|value| value.strip_prefix("Bearer "))
.map(|token| ct_eq(token.as_bytes(), expected.as_bytes()))
.unwrap_or(false)
}
}
async fn require_bearer(State(auth): State<DashAuth>, req: Request, next: Next) -> Response {
if auth.is_authorized(req.headers()) {
next.run(req).await
} else {
(
StatusCode::UNAUTHORIZED,
[(header::WWW_AUTHENTICATE, "Bearer")],
"unauthorized",
)
.into_response()
}
}
#[derive(Serialize)]
struct AuthStatus {
auth_required: bool,
authenticated: bool,
}
async fn auth_status(State(auth): State<DashAuth>, headers: HeaderMap) -> Json<AuthStatus> {
Json(AuthStatus {
auth_required: auth.token.is_some(),
authenticated: auth.is_authorized(&headers),
})
}
impl MultiAgentStatusServer {
pub async fn run(
port: u16,
statuses: HashMap<String, SharedAgentStatus>,
chat_agents: HashMap<String, Arc<dyn ChatCapable>>,
configs: HashMap<String, AgentConfig>,
) {
Self::run_with_registry(port, statuses, chat_agents, configs, None).await;
}
pub async fn run_with_registry(
port: u16,
statuses: HashMap<String, SharedAgentStatus>,
chat_agents: HashMap<String, Arc<dyn ChatCapable>>,
configs: HashMap<String, AgentConfig>,
registry: Option<OrchestratorRegistry>,
) {
let rw_configs = configs
.into_iter()
.map(|(k, v)| (k, Arc::new(RwLock::new(v))))
.collect();
Self::run_control_plane(
port,
statuses,
chat_agents,
rw_configs,
HashMap::new(),
HashMap::new(),
HashMap::new(),
registry,
None, )
.await;
}
#[allow(clippy::too_many_arguments)]
pub async fn run_control_plane(
port: u16,
statuses: HashMap<String, SharedAgentStatus>,
chat_agents: HashMap<String, Arc<dyn ChatCapable>>,
configs: HashMap<String, Arc<RwLock<AgentConfig>>>,
buffers: HashMap<String, Arc<ResponseBuffer>>,
pause_handles: HashMap<String, Arc<AtomicBool>>,
event_stores: HashMap<String, AgentEventStore>,
registry: Option<OrchestratorRegistry>,
middleware: Option<Arc<crate::middleware::pipeline::MiddlewarePipeline>>,
) {
let base_hold = buffers
.values()
.map(|b| b.base_hold_duration())
.max()
.unwrap_or(std::time::Duration::ZERO);
let base_secs = base_hold.as_secs();
let sla_secs = base_secs;
let agent_count = configs.len();
let state = MultiAppState {
statuses,
chat_agents,
configs,
buffers,
pause_handles,
event_stores,
orchestrator_registry: registry,
base_hold_secs: Arc::new(AtomicU64::new(base_secs)),
response_sla_secs: Arc::new(AtomicU64::new(sla_secs)),
buffer_floor_pct: Arc::new(AtomicU64::new(0)), before_release_middleware: middleware,
auth_token: resolve_dashboard_token(std::env::var("QUORUM_DASHBOARD_TOKEN").ok()),
};
for buf in state.buffers.values() {
buf.set_response_sla(base_hold);
}
let auth_enabled = state.auth_token.is_some();
let app = build_router(state);
let ip = resolve_dashboard_bind(std::env::var("QUORUM_DASHBOARD_BIND").ok().as_deref());
let addr = SocketAddr::from((ip, port));
if refuse_unauthenticated_exposure(&ip, auth_enabled) {
error!(
bind = %ip,
"refusing to start the dashboard: bound to a non-loopback address with no \
QUORUM_DASHBOARD_TOKEN — that would expose the control plane unauthenticated. \
Set QUORUM_DASHBOARD_TOKEN, or bind to loopback (QUORUM_DASHBOARD_BIND)."
);
return;
}
info!(
"Multi-agent dashboard → http://{}/ ({} agents) Swagger UI → http://{}/swagger-ui/",
addr, agent_count, addr
);
let listener = match tokio::net::TcpListener::bind(addr).await {
Ok(l) => l,
Err(e) => {
error!("Failed to bind multi-agent server on port {}: {}", port, e);
return;
}
};
if let Err(e) = axum::serve(listener, app).await {
error!("Multi-agent server error: {}", e);
}
}
}
#[derive(Serialize, ToSchema)]
pub(super) struct AgentSummary {
name: String,
model_name: String,
provider_id: String,
nats_connected: bool,
current_job: Option<String>,
current_phase: Option<String>,
has_chat: bool,
is_paused: bool,
buffered_count: u32,
error_rate: f32,
mean_score: Option<f32>,
score_std_dev: Option<f32>,
avg_response_ms: Option<u64>,
is_flagged: bool,
flag_reason: Option<String>,
auto_approve: bool,
auto_approve_threshold: f32,
}
#[utoipa::path(
get,
path = "/api/agents",
responses(
(status = 200, description = "List of all agents with summary status", body = Vec<AgentSummary>)
),
tag = "Agents"
)]
pub(super) async fn list_agents(State(state): State<MultiAppState>) -> Json<Vec<AgentSummary>> {
let mut agents = Vec::new();
for (name, config) in &state.configs {
let config = config.read().await;
let (
nats_connected,
current_job,
current_phase,
is_paused,
buffered_count,
error_rate,
mean_score,
score_std_dev,
avg_response_ms,
is_flagged,
flag_reason,
) = if let Some(status) = state.statuses.get(name) {
let snap = status.read().await;
let avg_ms = if snap.recent_tasks.is_empty() {
None
} else {
let total: u64 = snap.recent_tasks.iter().map(|t| t.duration_ms).sum();
Some(total / snap.recent_tasks.len() as u64)
};
(
snap.nats_connected,
snap.current_job.clone(),
snap.current_phase.clone(),
snap.is_paused,
snap.buffered_count,
snap.error_rate,
snap.mean_score,
snap.score_std_dev,
avg_ms,
snap.is_flagged,
snap.flag_reason.clone(),
)
} else {
(
false, None, None, false, 0, 0.0, None, None, None, false, None,
)
};
let buffered_count = if let Some(buf) = state.buffers.get(name) {
buf.len().await as u32
} else {
buffered_count
};
let (auto_approve, auto_approve_threshold) = state
.buffers
.get(name)
.map(|b| (b.is_auto_approve(), b.auto_approve_threshold()))
.unwrap_or((true, 1.0));
agents.push(AgentSummary {
name: name.clone(),
model_name: config.model_name.clone(),
provider_id: config.provider_id.clone(),
nats_connected,
current_job,
current_phase,
has_chat: state.chat_agents.contains_key(name),
is_paused,
buffered_count,
error_rate,
mean_score,
score_std_dev,
avg_response_ms,
is_flagged,
flag_reason,
auto_approve,
auto_approve_threshold,
});
}
agents.sort_by(|a, b| a.name.cmp(&b.name));
Json(agents)
}
#[derive(Serialize, ToSchema)]
pub(super) struct AgentDiagnostics {
name: String,
model_name: String,
uptime_secs: u64,
tasks_completed: u64,
tasks_failed: u64,
error_rate: f32,
is_paused: bool,
is_flagged: bool,
flag_reason: Option<String>,
#[schema(value_type = Vec<crate::status::EventLogEntry>)]
recent_errors: Vec<crate::status::EventLogEntry>,
#[schema(value_type = Vec<crate::status::TaskLogEntry>)]
recent_failed_tasks: Vec<crate::status::TaskLogEntry>,
}
fn diagnostics_from_snapshot(
name: &str,
snap: &crate::status::AgentStatusSnapshot,
) -> AgentDiagnostics {
const MAX: usize = 20;
let recent_errors: Vec<_> = snap
.event_log
.iter()
.rev()
.filter(|e| e.event_type == "agent_error")
.take(MAX)
.cloned()
.collect();
let recent_failed_tasks: Vec<_> = snap
.recent_tasks
.iter()
.rev()
.filter(|t| t.status == "error")
.take(MAX)
.cloned()
.collect();
AgentDiagnostics {
name: name.to_string(),
model_name: snap.model_name.clone(),
uptime_secs: snap.uptime_secs,
tasks_completed: snap.tasks_completed,
tasks_failed: snap.tasks_failed,
error_rate: snap.error_rate,
is_paused: snap.is_paused,
is_flagged: snap.is_flagged,
flag_reason: snap.flag_reason.clone(),
recent_errors,
recent_failed_tasks,
}
}
#[utoipa::path(
get,
path = "/api/agents/{name}/diagnostics",
params(("name" = String, Path, description = "Agent name")),
responses(
(status = 200, description = "Agent metrics + latest errors", body = AgentDiagnostics),
(status = 404, description = "Unknown agent")
),
tag = "Agents"
)]
pub(super) async fn agent_diagnostics(
State(state): State<MultiAppState>,
Path(name): Path<String>,
) -> Result<Json<AgentDiagnostics>, StatusCode> {
let status = state.statuses.get(&name).ok_or(StatusCode::NOT_FOUND)?;
let snap = status.read().await;
Ok(Json(diagnostics_from_snapshot(&name, &snap)))
}
use super::agent_events::{AgentEvent, TasksView, ToolCallsView};
const ERROR_WINDOW_HOURS: i64 = 24;
#[derive(Serialize, ToSchema)]
pub(super) struct AgentErrorEntry {
agent: String,
model_name: String,
timestamp: String,
job_id: Option<String>,
detail: String,
}
#[derive(Serialize, ToSchema)]
pub(super) struct AgentErrorsReport {
window_hours: i64,
stream_cap: i64,
total: usize,
errors: Vec<AgentErrorEntry>,
}
struct AgentErrorSource {
agent: String,
model_name: String,
events: Vec<AgentEvent>,
}
fn flatten_error_feed(sources: &[AgentErrorSource]) -> Vec<AgentErrorEntry> {
let mut dated: Vec<(chrono::DateTime<chrono::Utc>, AgentErrorEntry)> = sources
.iter()
.flat_map(|source| {
source.events.iter().filter_map(move |event| {
let at = event.parsed_time()?;
Some((
at,
AgentErrorEntry {
agent: source.agent.clone(),
model_name: source.model_name.clone(),
timestamp: event.timestamp.clone(),
job_id: event.job_id.clone(),
detail: event.detail.clone(),
},
))
})
})
.collect();
dated.sort_by_key(|(at, _)| std::cmp::Reverse(*at));
dated.into_iter().map(|(_, entry)| entry).collect()
}
async fn model_name_for(state: &MultiAppState, agent: &str) -> String {
match state.configs.get(agent) {
Some(config) => config.read().await.model_name.clone(),
None => String::new(),
}
}
#[utoipa::path(
get,
path = "/api/agents/errors",
responses(
(status = 200, description = "Fleet-wide API errors over the last 24h", body = AgentErrorsReport)
),
tag = "Agents"
)]
pub(super) async fn agents_errors(State(state): State<MultiAppState>) -> Json<AgentErrorsReport> {
let cutoff = chrono::Utc::now() - chrono::Duration::hours(ERROR_WINDOW_HOURS);
let mut sources = Vec::with_capacity(state.event_stores.len());
for (agent, store) in &state.event_stores {
let events = match store.read_since(cutoff).await {
Ok(events) => super::agent_events::collect_errors(&events),
Err(e) => {
error!(agent = %agent, error = %e, "failed to read agent error log");
continue;
}
};
sources.push(AgentErrorSource {
agent: agent.clone(),
model_name: model_name_for(&state, agent).await,
events,
});
}
let errors = flatten_error_feed(&sources);
Json(AgentErrorsReport {
window_hours: ERROR_WINDOW_HOURS,
stream_cap: super::agent_events::STREAM_MAX_MESSAGES,
total: errors.len(),
errors,
})
}
async fn read_agent_window(
state: &MultiAppState,
name: &str,
) -> Option<Result<Vec<AgentEvent>, StatusCode>> {
let store = state.event_stores.get(name)?;
let cutoff = chrono::Utc::now() - chrono::Duration::hours(ERROR_WINDOW_HOURS);
Some(match store.read_since(cutoff).await {
Ok(events) => Ok(events),
Err(e) => {
error!(agent = %name, error = %e, "failed to read agent event log");
Err(StatusCode::BAD_GATEWAY)
}
})
}
#[utoipa::path(
get,
path = "/api/agents/{name}/tasks",
params(("name" = String, Path, description = "Agent name")),
responses(
(status = 200, description = "In-flight and finished tasks", body = TasksView),
(status = 404, description = "Unknown agent")
),
tag = "Agents"
)]
pub(super) async fn agent_tasks(
State(state): State<MultiAppState>,
Path(name): Path<String>,
) -> Result<Json<TasksView>, StatusCode> {
if !state.configs.contains_key(&name) && !state.statuses.contains_key(&name) {
return Err(StatusCode::NOT_FOUND);
}
let events = match read_agent_window(&state, &name).await {
Some(result) => result?,
None => return Ok(Json(TasksView::default())),
};
Ok(Json(super::agent_events::reconcile_tasks(&events)))
}
#[utoipa::path(
get,
path = "/api/agents/{name}/tool-calls",
params(("name" = String, Path, description = "Agent name")),
responses(
(status = 200, description = "Pending and finished tool calls", body = ToolCallsView),
(status = 404, description = "Unknown agent")
),
tag = "Agents"
)]
pub(super) async fn agent_tool_calls(
State(state): State<MultiAppState>,
Path(name): Path<String>,
) -> Result<Json<ToolCallsView>, StatusCode> {
if !state.configs.contains_key(&name) && !state.statuses.contains_key(&name) {
return Err(StatusCode::NOT_FOUND);
}
let events = match read_agent_window(&state, &name).await {
Some(result) => result?,
None => return Ok(Json(ToolCallsView::default())),
};
Ok(Json(super::agent_events::reconcile_tool_calls(&events)))
}
fn build_router(state: MultiAppState) -> Router {
use utoipa::OpenApi;
let swagger_ui = utoipa_swagger_ui::SwaggerUi::new("/swagger-ui")
.url("/api-docs/openapi.json", api_docs::ApiDoc::openapi());
let protected = Router::new()
.route(
"/api/config",
get(registry_handlers::get_global_config).put(registry_handlers::update_global_config),
)
.route("/api/agents", get(list_agents))
.route("/api/agents/errors", get(agents_errors))
.route("/api/agents/{name}/diagnostics", get(agent_diagnostics))
.route("/api/agents/{name}/tasks", get(agent_tasks))
.route("/api/agents/{name}/tool-calls", get(agent_tool_calls))
.route(
"/api/agents/register",
post(registration_handlers::register_agent),
)
.route(
"/api/agents/bulk",
post(registration_handlers::bulk_register),
)
.route(
"/api/agents/pause-all",
put(hitl_handlers::pause_all_agents),
)
.route("/api/agents/auto-all", put(hitl_handlers::auto_all_agents))
.route(
"/api/agents/{name}/status",
get(status_handlers::agent_status),
)
.route(
"/api/agents/{name}/config",
get(status_handlers::agent_config).put(hitl_handlers::agent_config_update),
)
.route("/api/agents/{name}/chat", post(chat_handlers::agent_chat))
.route("/api/agents/{name}/pause", put(hitl_handlers::agent_pause))
.route(
"/api/agents/{name}/auto",
put(hitl_handlers::agent_auto_approve),
)
.route(
"/api/agents/{name}/buffer",
get(hitl_handlers::agent_buffer_list),
)
.route(
"/api/agents/{name}/buffer/{id}",
get(hitl_handlers::agent_buffer_detail).put(hitl_handlers::agent_buffer_edit),
)
.route(
"/api/agents/{name}/buffer/{id}/release",
post(hitl_handlers::agent_buffer_release),
)
.route(
"/api/agents/{name}/buffer/{id}/reject",
post(hitl_handlers::agent_buffer_reject),
)
.route(
"/api/agents/{name}/buffer/{id}/stop",
post(hitl_handlers::agent_buffer_stop),
)
.route(
"/api/agents/{name}/buffer/{id}/unstop",
post(hitl_handlers::agent_buffer_unstop),
)
.route(
"/api/agents/{id}/manage",
put(registration_handlers::replace_agent)
.patch(registration_handlers::patch_agent)
.delete(registration_handlers::delete_agent),
)
.route(
"/api/orchestrators",
get(registry_handlers::list_orchestrators).post(registry_handlers::add_orchestrator),
)
.route(
"/api/orchestrators/budgets",
get(registry_handlers::get_orchestrator_budgets),
)
.route(
"/api/orchestrators/{orch_id}/proxy/{*path}",
get(registry_handlers::proxy_orchestrator_get)
.post(registry_handlers::proxy_orchestrator_post),
)
.route(
"/api/orchestrators/{orch_id}/stream/{job_id}",
get(registry_handlers::proxy_orchestrator_sse),
)
.route_layer(middleware::from_fn_with_state(
DashAuth::from_ref(&state),
require_bearer,
));
Router::new()
.merge(swagger_ui)
.route("/", get(dashboard_page))
.route("/auth/status", get(auth_status))
.merge(protected)
.with_state(state)
}
#[utoipa::path(
get,
path = "/",
responses(
(status = 200, description = "Multi-agent dashboard HTML page", content_type = "text/html")
),
tag = "Dashboard"
)]
async fn dashboard_page() -> impl IntoResponse {
let html = include_str!("../multi_status.html");
([(header::CACHE_CONTROL, "no-store")], Html(html))
}
#[cfg(test)]
mod tests;