use axum::Json;
use axum::extract::{Query, State};
use serde::{Deserialize, Serialize};
use crate::engine::api::state::ApiState;
use crate::engine::monitor::{SystemSnapshot, latest_system_snapshot};
use crate::sync::ClusterStatusView;
use crate::utils::logger::LogRecord;
pub async fn cluster_status(State(state): State<ApiState>) -> Json<Option<ClusterStatusView>> {
Json(
state
.state
.coordination
.as_ref()
.and_then(|c| c.cluster_status()),
)
}
#[derive(Debug, Clone, Serialize)]
pub struct PendingBreakdown {
pub task: usize,
pub download: usize,
pub response: usize,
pub parser: usize,
pub error: usize,
pub remote_task: usize,
pub total: usize,
}
#[derive(Debug, Clone, Serialize)]
pub struct EngineStats {
pub namespace: String,
pub single_node: bool,
pub clustered: bool,
pub pending: PendingBreakdown,
}
pub async fn engine_stats(State(state): State<ApiState>) -> Json<EngineStats> {
let (task, download, response, parser, error, remote_task) =
state.queue_manager.local_pending_breakdown().await;
let total = task + download + response + parser + error + remote_task;
let cfg = state.state.config.read().await;
Json(EngineStats {
namespace: cfg.name.clone(),
single_node: cfg.is_single_node_mode(),
clustered: state.state.coordination.is_some(),
pending: PendingBreakdown {
task,
download,
response,
parser,
error,
remote_task,
total,
},
})
}
pub async fn system_stats() -> Json<Option<SystemSnapshot>> {
Json(latest_system_snapshot())
}
#[derive(Debug, Deserialize)]
pub struct LogsQuery {
pub limit: Option<usize>,
}
pub async fn recent_logs(Query(q): Query<LogsQuery>) -> Json<Vec<LogRecord>> {
let limit = q.limit.unwrap_or(200).min(1000);
Json(crate::utils::logger::recent_logs(limit))
}