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 QueueBreakdown {
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 StageOutcome {
pub success: u64,
pub failure: u64,
pub total: u64,
pub success_rate: Option<f64>,
}
impl StageOutcome {
fn new(counter: &crate::engine::runner::StageCounter) -> Self {
let (success, failure) = counter.snapshot();
let total = success + failure;
Self {
success,
failure,
total,
success_rate: (total > 0).then(|| success as f64 / total as f64),
}
}
}
#[derive(Debug, Clone, Serialize)]
pub struct ThroughputStats {
pub task: StageOutcome,
pub request: StageOutcome,
pub parse: StageOutcome,
pub store: StageOutcome,
}
#[derive(Debug, Clone, Serialize)]
pub struct EngineStats {
pub namespace: String,
pub single_node: bool,
pub clustered: bool,
pub pending: QueueBreakdown,
pub inflight: QueueBreakdown,
pub outstanding: QueueBreakdown,
pub throughput: ThroughputStats,
}
pub async fn engine_stats(State(state): State<ApiState>) -> Json<EngineStats> {
let (p_task, p_download, p_response, p_parser, p_error, p_remote) =
state.queue_manager.local_pending_breakdown().await;
let (i_task, i_download, i_response, i_parser, i_error) = state.inflight.breakdown();
let pending = QueueBreakdown {
task: p_task,
download: p_download,
response: p_response,
parser: p_parser,
error: p_error,
remote_task: p_remote,
total: p_task + p_download + p_response + p_parser + p_error + p_remote,
};
let inflight = QueueBreakdown {
task: i_task,
download: i_download,
response: i_response,
parser: i_parser,
error: i_error,
remote_task: 0,
total: i_task + i_download + i_response + i_parser + i_error,
};
let outstanding = QueueBreakdown {
task: p_task + i_task,
download: p_download + i_download,
response: p_response + i_response,
parser: p_parser + i_parser,
error: p_error + i_error,
remote_task: p_remote,
total: pending.total + inflight.total,
};
let throughput = ThroughputStats {
task: StageOutcome::new(&state.outcomes.task),
request: StageOutcome::new(&state.outcomes.request),
parse: StageOutcome::new(&state.outcomes.parse),
store: StageOutcome::new(&state.outcomes.store),
};
let cfg = state.state.config.read().await;
Json(EngineStats {
namespace: cfg.name.clone(),
single_node: state.state.coordination.is_none(),
clustered: state.state.coordination.is_some(),
pending,
inflight,
outstanding,
throughput,
})
}
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))
}