#![allow(clippy::result_large_err)]
#[cfg(feature = "http-api")]
use axum::{
extract::{Extension, Path, State},
http::StatusCode,
response::Json,
};
#[cfg(feature = "http-api")]
use std::sync::Arc;
#[cfg(feature = "http-api")]
use utoipa;
#[cfg(feature = "http-api")]
use super::traits::RuntimeApiProvider;
#[cfg(feature = "http-api")]
use super::types::{
AddIdentityMappingRequest, AgentStatusResponse, ChannelActionResponse, ChannelAuditResponse,
ChannelDetail, ChannelHealthResponse, ChannelSummary, CreateAgentRequest, CreateAgentResponse,
CreateScheduleRequest, CreateScheduleResponse, DeleteAgentResponse, DeleteChannelResponse,
DeleteScheduleResponse, ErrorResponse, ExecuteAgentRequest, GetAgentHistoryResponse,
HeartbeatRequest, IdentityMappingEntry, MessageStatusResponse, NextRunsResponse,
PushEventRequest, ReceiveMessagesResponse, RegisterChannelRequest, RegisterChannelResponse,
ScheduleActionResponse, ScheduleDetail, ScheduleHistoryResponse, ScheduleSummary,
SchedulerHealthResponse, SendMessageRequest, SendMessageResponse, StatusResponse,
UpdateAgentRequest, UpdateAgentResponse, UpdateChannelRequest, UpdateScheduleRequest,
WorkflowExecutionRequest,
};
#[cfg(feature = "http-api")]
use super::api_keys::ValidatedKey;
#[cfg(feature = "http-api")]
use crate::types::AgentId;
#[cfg(feature = "http-api")]
fn check_agent_access(
validated: Option<&ValidatedKey>,
agent_id: &AgentId,
) -> Result<(), (StatusCode, Json<ErrorResponse>)> {
match validated {
Some(key) => match &key.agent_scope {
None => Ok(()),
Some(scope) => {
let target = agent_id.0.to_string();
if scope.iter().any(|s| s == &target) {
Ok(())
} else {
Err((
StatusCode::FORBIDDEN,
Json(ErrorResponse {
error: "API key is not scoped to this agent".to_string(),
code: "AGENT_SCOPE_DENIED".to_string(),
details: None,
}),
))
}
}
},
None => Ok(()),
}
}
#[cfg(feature = "http-api")]
pub(super) fn require_admin(
validated: Option<&ValidatedKey>,
) -> Result<(), (StatusCode, Json<ErrorResponse>)> {
match validated {
Some(key) if key.agent_scope.is_some() => Err((
StatusCode::FORBIDDEN,
Json(ErrorResponse {
error: "Scoped API keys may not access this endpoint".to_string(),
code: "ADMIN_REQUIRED".to_string(),
details: None,
}),
)),
_ => Ok(()),
}
}
#[cfg(feature = "http-api")]
fn scope_filter(validated: Option<&ValidatedKey>) -> Option<&Vec<String>> {
validated.and_then(|k| k.agent_scope.as_ref())
}
#[derive(serde::Deserialize)]
pub struct RunInspectionQuery {
public_key: String,
}
async fn capacity_response(
provider: Arc<dyn RuntimeApiProvider>,
validated: Option<Extension<ValidatedKey>>,
lease: Option<uuid::Uuid>,
) -> axum::response::Response {
use axum::response::IntoResponse;
let mut response = match require_admin(validated.as_deref()) {
Err(error) => error.into_response(),
Ok(()) => match provider.inspect_capacity(lease).await {
Ok(snapshot) => Json(snapshot).into_response(),
Err(error) => (
StatusCode::SERVICE_UNAVAILABLE,
Json(ErrorResponse {
code: "CAPACITY_UNAVAILABLE".into(),
error,
details: None,
}),
)
.into_response(),
},
};
response.headers_mut().insert(
axum::http::header::CACHE_CONTROL,
axum::http::HeaderValue::from_static("no-store"),
);
response
}
#[utoipa::path(get, path="/api/v1/sandbox/capacity",
responses((status=200,description="Current retained worker reservations"),(status=403,description="Administrative key required"),(status=503,description="Supervisor accounting unavailable")),
security(("bearer_auth"=[])))]
pub async fn inspect_capacity(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> axum::response::Response {
capacity_response(provider, validated, None).await
}
#[utoipa::path(get, path="/api/v1/sandbox/workers/{lease}/usage",
params(("lease"=String,Path,description="Retained supervisor lease UUID")),
responses((status=200,description="One separately sampled worker measurement"),(status=403,description="Administrative key required"),(status=503,description="Measurement unavailable")),
security(("bearer_auth"=[])))]
pub async fn inspect_worker_usage(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
Path(lease): Path<uuid::Uuid>,
) -> axum::response::Response {
capacity_response(provider, validated, Some(lease)).await
}
#[utoipa::path(
get,
path = "/api/v1/audit/runs/{agent_id}/{run_id}",
params(
("agent_id" = String, Path, description = "Agent UUID"),
("run_id" = String, Path, description = "Run UUID"),
("public_key" = String, Query, description = "Independently retained Ed25519 public key, 64 hex characters")
),
responses((status = 200, description = "Verified bounded snapshot"),
(status = 403, description = "Administrative key required"),
(status = 422, description = "Evidence unavailable or invalid")),
security(("bearer_auth" = []))
)]
pub async fn inspect_run(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
Path((agent_id, run_id)): Path<(AgentId, uuid::Uuid)>,
axum::extract::Query(query): axum::extract::Query<RunInspectionQuery>,
) -> axum::response::Response {
use axum::response::IntoResponse;
let mut response = match require_admin(validated.as_deref()) {
Err(error) => error.into_response(),
Ok(()) => {
let key = if query.public_key.len() == 64 {
hex::decode(&query.public_key)
.ok()
.and_then(|bytes| bytes.try_into().ok())
} else {
None
};
match key {
None => (
StatusCode::BAD_REQUEST,
Json(ErrorResponse {
code: "INVALID_AUDIT_KEY".into(),
error: "public_key must be 64 hexadecimal characters".into(),
details: None,
}),
)
.into_response(),
Some(key) => match provider.inspect_run(agent_id, run_id, key).await {
Ok(view) => Json(view).into_response(),
Err(error) => (
StatusCode::UNPROCESSABLE_ENTITY,
Json(ErrorResponse {
code: "AUDIT_UNAVAILABLE".into(),
error,
details: None,
}),
)
.into_response(),
},
}
}
};
response.headers_mut().insert(
axum::http::header::CACHE_CONTROL,
axum::http::HeaderValue::from_static("no-store"),
);
response
}
#[cfg(feature = "http-api")]
pub async fn serve_agents_md() -> Result<
(
StatusCode,
[(axum::http::header::HeaderName, &'static str); 1],
String,
),
StatusCode,
> {
let content = tokio::fs::read_to_string("AGENTS.md")
.await
.map_err(|_| StatusCode::NOT_FOUND)?;
let filtered = strip_sensitive_sections(&content);
Ok((
StatusCode::OK,
[(
axum::http::header::CONTENT_TYPE,
"text/markdown; charset=utf-8",
)],
filtered,
))
}
#[cfg(feature = "http-api")]
fn strip_sensitive_sections(content: &str) -> String {
const SENSITIVE_START: &str = "<!-- agents-md:sensitive-start -->";
const SENSITIVE_END: &str = "<!-- agents-md:sensitive-end -->";
let mut result = content.to_string();
while let (Some(start), Some(end)) = (result.find(SENSITIVE_START), result.find(SENSITIVE_END))
{
if end <= start {
break;
}
let end_pos = end + SENSITIVE_END.len();
let end_pos = if result[end_pos..].starts_with('\n') {
end_pos + 1
} else {
end_pos
};
let start_pos = if start > 0 && result.as_bytes()[start - 1] == b'\n' {
start - 1
} else {
start
};
result = format!("{}{}", &result[..start_pos], &result[end_pos..]);
}
result
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/workflows/execute",
params(("Idempotency-Key" = String, Header, description = "UUID retained for retries")),
request_body = WorkflowExecutionRequest,
responses(
(status = 200, description = "Workflow invocation queued; execution_id identifies its history", body = serde_json::Value),
(status = 401, description = "Authentication required"),
(status = 403, description = "Administrative authority required", body = ErrorResponse),
(status = 400, description = "Invalid UUID header or workflow source", body = serde_json::Value),
(status = 409, description = "Active, unresolved or conflicting invocation", body = serde_json::Value),
(status = 422, description = "Saved known terminal failure", body = serde_json::Value),
(status = 503, description = "Persistent admission unavailable", body = serde_json::Value)
),
tag = "workflows"
)]
pub async fn execute_workflow(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
caller: Option<Extension<super::invocations::AuthenticatedCaller>>,
headers: axum::http::HeaderMap,
Json(request): Json<WorkflowExecutionRequest>,
) -> Result<axum::response::Response, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
Ok(super::invocations::submit(
provider.as_ref(),
caller.map(|Extension(c)| c),
headers,
super::invocations::ExecutionTarget::Workflow(request),
)
.await)
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/agents/{id}/status",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
responses(
(status = 200, description = "Agent status retrieved successfully", body = AgentStatusResponse),
(status = 404, description = "Agent not found", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn get_agent_status(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<AgentStatusResponse>, (StatusCode, Json<ErrorResponse>)> {
check_agent_access(validated.as_ref().map(|Extension(k)| k), &agent_id)?;
match provider.get_agent_status(agent_id).await {
Ok(status) => Ok(Json(status)),
Err(e) => Err((
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "AGENT_NOT_FOUND".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/agents",
responses(
(status = 200, description = "Agents listed successfully", body = Vec<AgentSummary>),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn list_agents(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<Vec<super::types::AgentSummary>>, (StatusCode, Json<ErrorResponse>)> {
let scope = scope_filter(validated.as_ref().map(|Extension(k)| k)).cloned();
match provider.list_agents_detailed().await {
Ok(mut agents) => {
if let Some(scope) = scope {
agents.retain(|a| scope.iter().any(|s| s == &a.id.0.to_string()));
}
Ok(Json(agents))
}
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "LIST_AGENTS_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/metrics",
responses(
(status = 200, description = "Metrics retrieved successfully", body = serde_json::Value),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "system"
)]
pub async fn get_metrics(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<serde_json::Value>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
match provider.get_metrics().await {
Ok(metrics) => Ok(Json(metrics)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "METRICS_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/agents",
request_body = CreateAgentRequest,
responses(
(status = 200, description = "Agent created successfully", body = CreateAgentResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn create_agent(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<CreateAgentRequest>,
) -> Result<Json<CreateAgentResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
match provider.create_agent(request).await {
Ok(response) => Ok(Json(response)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "AGENT_CREATION_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
put,
path = "/api/v1/agents/{id}",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
request_body = UpdateAgentRequest,
responses(
(status = 200, description = "Agent updated successfully", body = UpdateAgentResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn update_agent(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<UpdateAgentRequest>,
) -> Result<Json<UpdateAgentResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
match provider.update_agent(agent_id, request).await {
Ok(response) => Ok(Json(response)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "AGENT_UPDATE_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
delete,
path = "/api/v1/agents/{id}",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
responses(
(status = 200, description = "Agent deleted successfully", body = DeleteAgentResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn delete_agent(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<DeleteAgentResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
match provider.delete_agent(agent_id).await {
Ok(response) => Ok(Json(response)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "AGENT_DELETION_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/agents/{id}/execute",
params(
("id" = AgentId, Path, description = "Agent identifier"),
("Idempotency-Key" = String, Header, description = "UUID retained for retries")
),
request_body = ExecuteAgentRequest,
responses(
(status = 200, description = "Queued invocation or previously recorded completion", body = serde_json::Value),
(status = 400, description = "Invalid UUID header", body = serde_json::Value),
(status = 404, description = "Registered agent unavailable", body = serde_json::Value),
(status = 409, description = "Active, unresolved or conflicting invocation", body = serde_json::Value),
(status = 422, description = "Saved known terminal failure", body = serde_json::Value),
(status = 503, description = "Persistent admission unavailable", body = serde_json::Value)
),
tag = "agents"
)]
pub async fn execute_agent(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
caller: Option<Extension<super::invocations::AuthenticatedCaller>>,
headers: axum::http::HeaderMap,
Json(request): Json<ExecuteAgentRequest>,
) -> Result<axum::response::Response, (StatusCode, Json<ErrorResponse>)> {
check_agent_access(validated.as_ref().map(|Extension(k)| k), &agent_id)?;
Ok(super::invocations::submit(
provider.as_ref(),
caller.map(|Extension(c)| c),
headers,
super::invocations::ExecutionTarget::Agent {
id: agent_id,
input: request.input,
},
)
.await)
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/agents/{id}/history",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
responses(
(status = 200, description = "Agent history retrieved successfully", body = GetAgentHistoryResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn get_agent_history(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<GetAgentHistoryResponse>, (StatusCode, Json<ErrorResponse>)> {
check_agent_access(validated.as_ref().map(|Extension(k)| k), &agent_id)?;
match provider.get_agent_history(agent_id).await {
Ok(response) => Ok(Json(response)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "AGENT_HISTORY_FAILED".to_string(),
details: None,
}),
)),
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/schedules",
responses(
(status = 200, description = "Schedules listed", body = Vec<ScheduleSummary>),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn list_schedules(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<Vec<ScheduleSummary>>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.list_schedules().await.map(Json).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "LIST_SCHEDULES_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/schedules",
request_body = CreateScheduleRequest,
responses(
(status = 201, description = "Schedule created", body = CreateScheduleResponse),
(status = 400, description = "Bad request", body = ErrorResponse),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn create_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<CreateScheduleRequest>,
) -> Result<(StatusCode, Json<CreateScheduleResponse>), (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.create_schedule(request)
.await
.map(|r| (StatusCode::CREATED, Json(r)))
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "CREATE_SCHEDULE_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/schedules/{id}",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Schedule details", body = ScheduleDetail),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn get_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ScheduleDetail>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.get_schedule(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_NOT_FOUND".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
put,
path = "/api/v1/schedules/{id}",
params(("id" = String, Path, description = "Job UUID")),
request_body = UpdateScheduleRequest,
responses(
(status = 200, description = "Schedule updated", body = ScheduleDetail),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn update_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<UpdateScheduleRequest>,
) -> Result<Json<ScheduleDetail>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.update_schedule(&id, request)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "UPDATE_SCHEDULE_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
delete,
path = "/api/v1/schedules/{id}",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Schedule deleted", body = DeleteScheduleResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn delete_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<DeleteScheduleResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.delete_schedule(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_NOT_FOUND".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/schedules/{id}/pause",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Schedule paused", body = ScheduleActionResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn pause_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ScheduleActionResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.pause_schedule(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_PAUSE_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/schedules/{id}/resume",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Schedule resumed", body = ScheduleActionResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn resume_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ScheduleActionResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.resume_schedule(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_RESUME_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/schedules/{id}/trigger",
params(
("id" = String, Path, description = "Job UUID"),
("Idempotency-Key" = String, Header, description = "UUID retained for retries")
),
responses(
(status = 200, description = "Queued or verified saved completion", body = serde_json::Value),
(status = 400, description = "Invalid invocation identity", body = serde_json::Value),
(status = 403, description = "Administrative or schedule authority refused", body = serde_json::Value),
(status = 404, description = "Schedule not found", body = serde_json::Value),
(status = 409, description = "In progress, unresolved or conflicting identity", body = serde_json::Value),
(status = 422, description = "Verified saved failure", body = serde_json::Value),
(status = 503, description = "Persistent admission unavailable", body = serde_json::Value)
),
tag = "schedules"
)]
pub async fn trigger_schedule(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
caller: Option<Extension<super::invocations::AuthenticatedCaller>>,
headers: axum::http::HeaderMap,
) -> Result<axum::response::Response, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
Ok(super::invocations::submit(
provider.as_ref(),
caller.map(|Extension(c)| c),
headers,
super::invocations::ExecutionTarget::Schedule { job_id: id },
)
.await)
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/schedules/{id}/history",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Schedule history", body = ScheduleHistoryResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn get_schedule_history(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ScheduleHistoryResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.get_schedule_history(&id, 50)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_HISTORY_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/schedules/{id}/next-runs",
params(("id" = String, Path, description = "Job UUID")),
responses(
(status = 200, description = "Next runs", body = NextRunsResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "schedules"
)]
pub async fn get_schedule_next_runs(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<NextRunsResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.get_schedule_next_runs(&id, 10)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULE_NEXT_RUNS_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/health/scheduler",
responses(
(status = 200, description = "Scheduler health", body = SchedulerHealthResponse),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "system"
)]
pub async fn get_scheduler_health(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<SchedulerHealthResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.get_scheduler_health()
.await
.map(Json)
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "SCHEDULER_HEALTH_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/channels",
responses(
(status = 200, description = "Channels listed", body = Vec<ChannelSummary>),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn list_channels(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<Vec<ChannelSummary>>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.list_channels().await.map(Json).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "LIST_CHANNELS_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/channels",
request_body = RegisterChannelRequest,
responses(
(status = 201, description = "Channel registered", body = RegisterChannelResponse),
(status = 400, description = "Bad request", body = ErrorResponse),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn register_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<RegisterChannelRequest>,
) -> Result<(StatusCode, Json<RegisterChannelResponse>), (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.register_channel(request)
.await
.map(|r| (StatusCode::CREATED, Json(r)))
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "REGISTER_CHANNEL_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/channels/{id}",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Channel details", body = ChannelDetail),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn get_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ChannelDetail>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.get_channel(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_NOT_FOUND".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
put,
path = "/api/v1/channels/{id}",
params(("id" = String, Path, description = "Channel ID")),
request_body = UpdateChannelRequest,
responses(
(status = 200, description = "Channel updated", body = ChannelDetail),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn update_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<UpdateChannelRequest>,
) -> Result<Json<ChannelDetail>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.update_channel(&id, request)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "UPDATE_CHANNEL_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
delete,
path = "/api/v1/channels/{id}",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Channel deleted", body = DeleteChannelResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn delete_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<DeleteChannelResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.delete_channel(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_NOT_FOUND".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/channels/{id}/start",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Channel started", body = ChannelActionResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn start_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ChannelActionResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.start_channel(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_START_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/channels/{id}/stop",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Channel stopped", body = ChannelActionResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn stop_channel(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ChannelActionResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider.stop_channel(&id).await.map(Json).map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_STOP_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/channels/{id}/health",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Channel health", body = ChannelHealthResponse),
(status = 404, description = "Not found", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn get_channel_health(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ChannelHealthResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.get_channel_health(&id)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_HEALTH_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/channels/{id}/mappings",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Identity mappings listed", body = Vec<IdentityMappingEntry>),
(status = 404, description = "Not found", body = ErrorResponse),
(status = 501, description = "Not implemented (community edition)", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn list_channel_mappings(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<Vec<IdentityMappingEntry>>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.list_channel_mappings(&id)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_MAPPINGS_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/channels/{id}/mappings",
params(("id" = String, Path, description = "Channel ID")),
request_body = AddIdentityMappingRequest,
responses(
(status = 201, description = "Mapping added", body = IdentityMappingEntry),
(status = 404, description = "Not found", body = ErrorResponse),
(status = 501, description = "Not implemented (community edition)", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn add_channel_mapping(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<AddIdentityMappingRequest>,
) -> Result<(StatusCode, Json<IdentityMappingEntry>), (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.add_channel_mapping(&id, request)
.await
.map(|r| (StatusCode::CREATED, Json(r)))
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "ADD_CHANNEL_MAPPING_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
delete,
path = "/api/v1/channels/{id}/mappings/{user_id}",
params(
("id" = String, Path, description = "Channel ID"),
("user_id" = String, Path, description = "Platform user ID to remove")
),
responses(
(status = 204, description = "Mapping removed"),
(status = 404, description = "Not found", body = ErrorResponse),
(status = 501, description = "Not implemented (community edition)", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn remove_channel_mapping(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path((id, user_id)): Path<(String, String)>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.remove_channel_mapping(&id, &user_id)
.await
.map(|_| StatusCode::NO_CONTENT)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "REMOVE_CHANNEL_MAPPING_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/channels/{id}/audit",
params(("id" = String, Path, description = "Channel ID")),
responses(
(status = 200, description = "Audit log", body = ChannelAuditResponse),
(status = 404, description = "Not found", body = ErrorResponse),
(status = 501, description = "Not implemented (community edition)", body = ErrorResponse)
),
tag = "channels"
)]
pub async fn get_channel_audit(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ChannelAuditResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
provider
.get_channel_audit(&id, 50)
.await
.map(Json)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: e.to_string(),
code: "CHANNEL_AUDIT_FAILED".to_string(),
details: None,
}),
)
})
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/agents/{id}/heartbeat",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
request_body = HeartbeatRequest,
responses(
(status = 200, description = "Heartbeat accepted"),
(status = 404, description = "Agent not found or not external", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn agent_heartbeat(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<HeartbeatRequest>,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
check_agent_access(validated.as_ref().map(|Extension(k)| k), &agent_id)?;
match provider.update_agent_heartbeat(agent_id, request).await {
Ok(()) => Ok(StatusCode::OK),
Err(e) => {
use crate::types::RuntimeError;
let (status, code) = match &e {
RuntimeError::Authentication(_) => {
(StatusCode::UNAUTHORIZED, "AGENTPIN_VERIFICATION_FAILED")
}
_ => (StatusCode::NOT_FOUND, "HEARTBEAT_FAILED"),
};
Err((
status,
Json(ErrorResponse {
error: e.to_string(),
code: code.to_string(),
details: None,
}),
))
}
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/agents/{id}/events",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
request_body = PushEventRequest,
responses(
(status = 200, description = "Event accepted"),
(status = 404, description = "Agent not found or not external", body = ErrorResponse)
),
tag = "agents"
)]
pub async fn agent_push_event(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<PushEventRequest>,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
check_agent_access(validated.as_ref().map(|Extension(k)| k), &agent_id)?;
match provider.push_agent_event(agent_id, request).await {
Ok(()) => Ok(StatusCode::OK),
Err(e) => {
use crate::types::RuntimeError;
let (status, code) = match &e {
RuntimeError::Authentication(_) => {
(StatusCode::UNAUTHORIZED, "AGENTPIN_VERIFICATION_FAILED")
}
_ => (StatusCode::NOT_FOUND, "PUSH_EVENT_FAILED"),
};
Err((
status,
Json(ErrorResponse {
error: e.to_string(),
code: code.to_string(),
details: None,
}),
))
}
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
post,
path = "/api/v1/agents/{id}/messages",
params(
("id" = AgentId, Path, description = "Recipient agent identifier")
),
request_body = SendMessageRequest,
responses(
(status = 200, description = "Message queued", body = SendMessageResponse),
(status = 400, description = "Invalid request", body = ErrorResponse),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "messages"
)]
pub async fn send_agent_message(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
Json(request): Json<SendMessageRequest>,
) -> Result<Json<SendMessageResponse>, (StatusCode, Json<ErrorResponse>)> {
let key = validated.as_ref().map(|Extension(k)| k);
check_agent_access(key, &request.sender)?;
match provider.send_agent_message(agent_id, request).await {
Ok(response) => Ok(Json(response)),
Err(e) => {
use crate::types::RuntimeError;
let (status, code) = match &e {
RuntimeError::Authentication(_) => {
(StatusCode::UNAUTHORIZED, "AGENTPIN_VERIFICATION_FAILED")
}
_ => (StatusCode::INTERNAL_SERVER_ERROR, "SEND_MESSAGE_FAILED"),
};
Err((
status,
Json(ErrorResponse {
error: e.to_string(),
code: code.to_string(),
details: None,
}),
))
}
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/agents/{id}/messages",
params(
("id" = AgentId, Path, description = "Agent identifier")
),
responses(
(status = 200, description = "Pending messages", body = ReceiveMessagesResponse),
(status = 500, description = "Internal error", body = ErrorResponse)
),
tag = "messages"
)]
pub async fn receive_agent_messages(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(agent_id): Path<AgentId>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<ReceiveMessagesResponse>, (StatusCode, Json<ErrorResponse>)> {
let key = validated.as_ref().map(|Extension(k)| k);
check_agent_access(key, &agent_id)?;
match provider.receive_agent_messages(agent_id).await {
Ok(response) => Ok(Json(response)),
Err(e) => {
use crate::types::{CommunicationError, RuntimeError};
let (status, code) = match &e {
RuntimeError::Communication(CommunicationError::AgentNotRegistered { .. }) => {
(StatusCode::NOT_FOUND, "AGENT_NOT_FOUND")
}
_ => (StatusCode::INTERNAL_SERVER_ERROR, "RECEIVE_MESSAGES_FAILED"),
};
Err((
status,
Json(ErrorResponse {
error: e.to_string(),
code: code.to_string(),
details: None,
}),
))
}
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/messages/{id}/status",
params(
("id" = String, Path, description = "Message identifier (UUID)")
),
responses(
(status = 200, description = "Delivery status", body = MessageStatusResponse),
(status = 400, description = "Invalid message ID", body = ErrorResponse),
(status = 404, description = "Message not found", body = ErrorResponse)
),
tag = "messages"
)]
pub async fn get_message_status(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
Path(message_id): Path<String>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<MessageStatusResponse>, (StatusCode, Json<ErrorResponse>)> {
if let Some(Extension(key)) = validated.as_ref() {
if key.agent_scope.is_some() {
return Err((
StatusCode::FORBIDDEN,
Json(ErrorResponse {
error: "Scoped API keys may not query message status".to_string(),
code: "AGENT_SCOPE_DENIED".to_string(),
details: None,
}),
));
}
}
match provider.get_message_status(&message_id).await {
Ok(response) => Ok(Json(response)),
Err(e) => {
use crate::types::{CommunicationError, RuntimeError};
let (status, code) = match &e {
RuntimeError::Communication(CommunicationError::InvalidFormat(_)) => {
(StatusCode::BAD_REQUEST, "INVALID_MESSAGE_ID")
}
_ => (StatusCode::NOT_FOUND, "MESSAGE_STATUS_FAILED"),
};
Err((
status,
Json(ErrorResponse {
error: e.to_string(),
code: code.to_string(),
details: None,
}),
))
}
}
}
#[cfg(feature = "http-api")]
#[utoipa::path(
get,
path = "/api/v1/status",
responses(
(status = 200, description = "Operational status retrieved successfully", body = StatusResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "system"
)]
pub async fn get_status(
State(provider): State<Arc<dyn RuntimeApiProvider>>,
validated: Option<Extension<ValidatedKey>>,
) -> Result<Json<StatusResponse>, (StatusCode, Json<ErrorResponse>)> {
require_admin(validated.as_ref().map(|Extension(k)| k))?;
let healthy = provider
.get_system_health()
.await
.map(|v| {
v.get("status")
.and_then(|s| s.as_str())
.map(|s| s == "healthy")
.unwrap_or(true)
})
.unwrap_or(false);
let agents = provider.list_agents().await.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "STATUS_AGENTS_FAILED".to_string(),
details: None,
}),
)
})?;
let schedules = provider.list_schedules().await.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "STATUS_SCHEDULES_FAILED".to_string(),
details: None,
}),
)
})?;
let channels = provider.list_channels().await.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: e.to_string(),
code: "STATUS_CHANNELS_FAILED".to_string(),
details: None,
}),
)
})?;
Ok(Json(StatusResponse {
version: env!("CARGO_PKG_VERSION").to_string(),
healthy,
agent_count: agents.len() as u64,
schedule_count: schedules.len() as u64,
channel_count: channels.len() as u64,
}))
}
#[cfg(all(test, feature = "http-api"))]
mod scope_tests {
use super::*;
fn key_with_scope(scope: Option<Vec<String>>) -> ValidatedKey {
ValidatedKey {
key_id: "test".to_string(),
agent_scope: scope,
}
}
#[test]
fn require_admin_rejects_scoped_keys() {
let k = key_with_scope(Some(vec!["a".into()]));
assert!(require_admin(Some(&k)).is_err());
}
#[test]
fn require_admin_allows_unscoped_key() {
let k = key_with_scope(None);
assert!(require_admin(Some(&k)).is_ok());
}
#[test]
fn require_admin_allows_legacy_env_token() {
assert!(require_admin(None).is_ok());
}
#[test]
fn check_agent_access_blocks_out_of_scope() {
let target = AgentId::new();
let k = key_with_scope(Some(vec![uuid::Uuid::new_v4().to_string()]));
assert!(check_agent_access(Some(&k), &target).is_err());
}
#[test]
fn check_agent_access_allows_in_scope() {
let target = AgentId::new();
let k = key_with_scope(Some(vec![target.0.to_string()]));
assert!(check_agent_access(Some(&k), &target).is_ok());
}
#[test]
fn scope_filter_returns_none_for_admin() {
let k = key_with_scope(None);
assert!(scope_filter(Some(&k)).is_none());
}
#[test]
fn scope_filter_returns_scope_list() {
let k = key_with_scope(Some(vec!["a".into(), "b".into()]));
assert_eq!(scope_filter(Some(&k)).unwrap().len(), 2);
}
}
#[cfg(all(test, feature = "http-api"))]
mod status_tests {
use super::*;
use crate::api::types::ExecuteAgentResponse;
use crate::types::RuntimeError;
use async_trait::async_trait;
use axum::extract::State;
struct MockStatusProvider {
agents: usize,
schedules: usize,
channels: usize,
healthy: bool,
}
#[async_trait]
impl RuntimeApiProvider for MockStatusProvider {
async fn execute_workflow(
&self,
_: WorkflowExecutionRequest,
) -> Result<serde_json::Value, RuntimeError> {
unimplemented!()
}
async fn get_agent_status(&self, _: AgentId) -> Result<AgentStatusResponse, RuntimeError> {
unimplemented!()
}
async fn get_system_health(&self) -> Result<serde_json::Value, RuntimeError> {
if self.healthy {
Ok(serde_json::json!({"status": "healthy"}))
} else {
Err(RuntimeError::Internal("unhealthy".to_string()))
}
}
async fn list_agents(&self) -> Result<Vec<AgentId>, RuntimeError> {
Ok((0..self.agents).map(|_| AgentId::new()).collect())
}
async fn shutdown_agent(&self, _: AgentId) -> Result<(), RuntimeError> {
unimplemented!()
}
async fn get_metrics(&self) -> Result<serde_json::Value, RuntimeError> {
unimplemented!()
}
async fn create_agent(
&self,
_: CreateAgentRequest,
) -> Result<CreateAgentResponse, RuntimeError> {
unimplemented!()
}
async fn update_agent(
&self,
_: AgentId,
_: UpdateAgentRequest,
) -> Result<UpdateAgentResponse, RuntimeError> {
unimplemented!()
}
async fn delete_agent(&self, _: AgentId) -> Result<DeleteAgentResponse, RuntimeError> {
unimplemented!()
}
async fn execute_agent(
&self,
_: AgentId,
_: ExecuteAgentRequest,
) -> Result<ExecuteAgentResponse, RuntimeError> {
unimplemented!()
}
async fn get_agent_history(
&self,
_: AgentId,
) -> Result<GetAgentHistoryResponse, RuntimeError> {
unimplemented!()
}
async fn list_schedules(&self) -> Result<Vec<ScheduleSummary>, RuntimeError> {
Ok((0..self.schedules)
.map(|i| ScheduleSummary {
job_id: i.to_string(),
name: format!("job-{i}"),
cron_expression: "0 * * * * *".to_string(),
timezone: "UTC".to_string(),
status: "active".to_string(),
enabled: true,
next_run: None,
run_count: 0,
})
.collect())
}
async fn create_schedule(
&self,
_: CreateScheduleRequest,
) -> Result<CreateScheduleResponse, RuntimeError> {
unimplemented!()
}
async fn get_schedule(&self, _: &str) -> Result<ScheduleDetail, RuntimeError> {
unimplemented!()
}
async fn update_schedule(
&self,
_: &str,
_: UpdateScheduleRequest,
) -> Result<ScheduleDetail, RuntimeError> {
unimplemented!()
}
async fn delete_schedule(&self, _: &str) -> Result<DeleteScheduleResponse, RuntimeError> {
unimplemented!()
}
async fn pause_schedule(&self, _: &str) -> Result<ScheduleActionResponse, RuntimeError> {
unimplemented!()
}
async fn resume_schedule(&self, _: &str) -> Result<ScheduleActionResponse, RuntimeError> {
unimplemented!()
}
async fn trigger_schedule(&self, _: &str) -> Result<ScheduleActionResponse, RuntimeError> {
unimplemented!()
}
async fn get_schedule_history(
&self,
_: &str,
_: usize,
) -> Result<ScheduleHistoryResponse, RuntimeError> {
unimplemented!()
}
async fn get_schedule_next_runs(
&self,
_: &str,
_: usize,
) -> Result<NextRunsResponse, RuntimeError> {
unimplemented!()
}
async fn get_scheduler_health(&self) -> Result<SchedulerHealthResponse, RuntimeError> {
unimplemented!()
}
async fn list_channels(&self) -> Result<Vec<ChannelSummary>, RuntimeError> {
Ok((0..self.channels)
.map(|i| ChannelSummary {
id: i.to_string(),
name: format!("chan-{i}"),
platform: "slack".to_string(),
status: "running".to_string(),
})
.collect())
}
async fn register_channel(
&self,
_: RegisterChannelRequest,
) -> Result<RegisterChannelResponse, RuntimeError> {
unimplemented!()
}
async fn get_channel(&self, _: &str) -> Result<ChannelDetail, RuntimeError> {
unimplemented!()
}
async fn update_channel(
&self,
_: &str,
_: UpdateChannelRequest,
) -> Result<ChannelDetail, RuntimeError> {
unimplemented!()
}
async fn delete_channel(&self, _: &str) -> Result<DeleteChannelResponse, RuntimeError> {
unimplemented!()
}
async fn start_channel(&self, _: &str) -> Result<ChannelActionResponse, RuntimeError> {
unimplemented!()
}
async fn stop_channel(&self, _: &str) -> Result<ChannelActionResponse, RuntimeError> {
unimplemented!()
}
async fn get_channel_health(&self, _: &str) -> Result<ChannelHealthResponse, RuntimeError> {
unimplemented!()
}
async fn list_channel_mappings(
&self,
_: &str,
) -> Result<Vec<IdentityMappingEntry>, RuntimeError> {
unimplemented!()
}
async fn add_channel_mapping(
&self,
_: &str,
_: AddIdentityMappingRequest,
) -> Result<IdentityMappingEntry, RuntimeError> {
unimplemented!()
}
async fn remove_channel_mapping(&self, _: &str, _: &str) -> Result<(), RuntimeError> {
unimplemented!()
}
async fn get_channel_audit(
&self,
_: &str,
_: usize,
) -> Result<ChannelAuditResponse, RuntimeError> {
unimplemented!()
}
async fn update_agent_heartbeat(
&self,
_: AgentId,
_: HeartbeatRequest,
) -> Result<(), RuntimeError> {
unimplemented!()
}
async fn push_agent_event(
&self,
_: AgentId,
_: PushEventRequest,
) -> Result<(), RuntimeError> {
unimplemented!()
}
async fn send_agent_message(
&self,
_: AgentId,
_: SendMessageRequest,
) -> Result<SendMessageResponse, RuntimeError> {
unimplemented!()
}
async fn receive_agent_messages(
&self,
_: AgentId,
) -> Result<ReceiveMessagesResponse, RuntimeError> {
unimplemented!()
}
async fn get_message_status(&self, _: &str) -> Result<MessageStatusResponse, RuntimeError> {
unimplemented!()
}
}
#[tokio::test]
async fn get_status_returns_correct_counts_and_version() {
let provider: Arc<dyn RuntimeApiProvider> = Arc::new(MockStatusProvider {
agents: 3,
schedules: 2,
channels: 1,
healthy: true,
});
let result = get_status(State(provider), None).await;
assert!(result.is_ok(), "handler should succeed");
let Json(status) = result.unwrap();
assert_eq!(status.agent_count, 3);
assert_eq!(status.schedule_count, 2);
assert_eq!(status.channel_count, 1);
assert!(status.healthy);
assert_eq!(status.version, env!("CARGO_PKG_VERSION"));
}
#[tokio::test]
async fn get_status_healthy_false_when_health_check_errors() {
let provider: Arc<dyn RuntimeApiProvider> = Arc::new(MockStatusProvider {
agents: 0,
schedules: 0,
channels: 0,
healthy: false,
});
let result = get_status(State(provider), None).await;
assert!(
result.is_ok(),
"handler should still succeed even when the health check returns an error"
);
let Json(status) = result.unwrap();
assert!(!status.healthy);
}
#[tokio::test]
async fn get_status_rejects_scoped_key() {
let provider: Arc<dyn RuntimeApiProvider> = Arc::new(MockStatusProvider {
agents: 0,
schedules: 0,
channels: 0,
healthy: true,
});
let k = ValidatedKey {
key_id: "scoped".to_string(),
agent_scope: Some(vec!["some-agent".to_string()]),
};
let result = get_status(State(provider), Some(Extension(k))).await;
assert!(result.is_err(), "scoped key should be rejected");
let (status_code, _) = result.unwrap_err();
assert_eq!(status_code, StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn capacity_requires_admin_and_unavailable_samples_are_never_zero_or_cached() {
let provider: Arc<dyn RuntimeApiProvider> = Arc::new(MockStatusProvider {
agents: 0,
schedules: 0,
channels: 0,
healthy: true,
});
for lease in [None, Some(uuid::Uuid::new_v4())] {
let denied = capacity_response(
provider.clone(),
Some(Extension(ValidatedKey {
key_id: "scoped".into(),
agent_scope: Some(vec![AgentId::new().to_string()]),
})),
lease,
)
.await;
assert_eq!(denied.status(), StatusCode::FORBIDDEN);
assert_eq!(
denied.headers()[axum::http::header::CACHE_CONTROL],
"no-store"
);
let missing = capacity_response(provider.clone(), None, lease).await;
assert_eq!(missing.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
missing.headers()[axum::http::header::CACHE_CONTROL],
"no-store"
);
let body = axum::body::to_bytes(missing.into_body(), 4096)
.await
.unwrap();
let body: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(body["code"], "CAPACITY_UNAVAILABLE");
assert!(body.get("reserved").is_none());
assert!(body.get("memory_bytes").is_none());
}
}
#[tokio::test]
async fn run_inspection_rejects_scoped_keys_before_validation_and_never_caches() {
let provider: Arc<dyn RuntimeApiProvider> = Arc::new(MockStatusProvider {
agents: 0,
schedules: 0,
channels: 0,
healthy: true,
});
let agent = AgentId::new();
let response = inspect_run(
State(provider.clone()),
Some(Extension(ValidatedKey {
key_id: "scoped".into(),
agent_scope: Some(vec![agent.to_string()]),
})),
Path((agent, uuid::Uuid::new_v4())),
axum::extract::Query(RunInspectionQuery {
public_key: "invalid".into(),
}),
)
.await;
assert_eq!(response.status(), StatusCode::FORBIDDEN);
assert_eq!(
response.headers()[axum::http::header::CACHE_CONTROL],
"no-store"
);
let response = inspect_run(
State(provider),
None,
Path((agent, uuid::Uuid::new_v4())),
axum::extract::Query(RunInspectionQuery {
public_key: "invalid".into(),
}),
)
.await;
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
}