use crate::authz::KbAccess;
use crate::error::{ApiError, ApiErrorResponse};
use crate::middleware::extract_request_id;
use crate::state::AppState;
use axum::Json;
use axum::extract::{Path, Request, State};
use axum::http::header;
use axum::response::{IntoResponse, Response};
use notedthat_core::{Verb, unix_to_rfc3339};
use notedthat_indexer::{IndexFailure, IndexState, KbHealthSnapshot, ReconcileSummary};
use serde::Serialize;
#[derive(Serialize)]
struct IndexHealthResponse {
kb_slug: String,
state: &'static str,
pending: usize,
queue: QueueView,
worker: &'static str,
last_indexed_at: Option<String>,
last_failure: Option<FailureView>,
last_reconcile: Option<ReconcileView>,
}
#[derive(Serialize)]
struct QueueView {
depth: usize,
capacity: usize,
}
#[derive(Serialize)]
struct FailureView {
at: String,
#[serde(skip_serializing_if = "Option::is_none")]
object_key: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
summary: Option<String>,
}
#[derive(Serialize)]
struct ReconcileView {
at: String,
#[serde(skip_serializing_if = "Option::is_none")]
scope: Option<String>,
#[serde(flatten, skip_serializing_if = "Option::is_none")]
counts: Option<ReconcileCounts>,
}
#[derive(Serialize)]
struct ReconcileCounts {
objects_on_disk: usize,
unchanged: usize,
changed: usize,
orphaned: usize,
}
pub(super) async fn get_index_health(
State(state): State<AppState>,
Path(kb_slug): Path<String>,
req: Request,
) -> Result<Response, ApiErrorResponse> {
let request_id = extract_request_id(&req);
let err = |error: ApiError| ApiErrorResponse {
error,
request_id: request_id.clone(),
};
let access = KbAccess::resolve(&state, &kb_slug, &req).map_err(&err)?;
access.require_visible().map_err(&err)?;
let snapshot = state.index_health.snapshot(access.kb().as_str());
let queue = QueueView {
capacity: state.indexer_tx.max_capacity(),
depth: state
.indexer_tx
.max_capacity()
.saturating_sub(state.indexer_tx.capacity()),
};
let body = render(&access, &kb_slug, snapshot, queue);
Ok(([(header::CACHE_CONTROL, "no-store")], Json(body)).into_response())
}
fn render(
access: &KbAccess,
kb_slug: &str,
snapshot: KbHealthSnapshot,
queue: QueueView,
) -> IndexHealthResponse {
let state = if queue.depth >= queue.capacity
&& rank(snapshot.state) < rank(IndexState::Backpressured)
{
IndexState::Backpressured
} else {
snapshot.state
};
IndexHealthResponse {
kb_slug: kb_slug.to_string(),
state: state.as_str(),
pending: snapshot.pending,
queue,
worker: if snapshot.worker_alive {
"running"
} else {
"stopped"
},
last_indexed_at: snapshot.last_indexed_at.map(unix_to_rfc3339),
last_failure: snapshot
.last_failure
.map(|failure| failure_view(access, failure)),
last_reconcile: snapshot
.last_reconcile
.map(|summary| reconcile_view(access, summary)),
}
}
fn rank(state: IndexState) -> u8 {
match state {
IndexState::Healthy => 0,
IndexState::Indexing => 1,
IndexState::Backpressured => 2,
IndexState::Stale => 3,
IndexState::Failed => 4,
}
}
fn failure_view(access: &KbAccess, failure: IndexFailure) -> FailureView {
let object_key = access
.allows(Verb::List, &failure.object_key)
.then_some(failure.object_key);
let summary = access
.filter(Verb::List)
.covers_whole_kb()
.then_some(failure.summary);
FailureView {
at: unix_to_rfc3339(failure.at),
object_key,
summary,
}
}
fn reconcile_view(access: &KbAccess, summary: ReconcileSummary) -> ReconcileView {
let visible = match &summary.scope {
Some(prefix) => access.allows(Verb::List, prefix),
None => access.filter(Verb::List).covers_whole_kb(),
};
let (scope, counts) = if visible {
(
summary.scope,
Some(ReconcileCounts {
objects_on_disk: summary.objects_on_disk,
unchanged: summary.unchanged,
changed: summary.changed,
orphaned: summary.orphaned,
}),
)
} else {
(None, None)
};
ReconcileView {
at: unix_to_rfc3339(summary.at),
scope,
counts,
}
}