use axum::Json;
use axum::extract::{Path, Query, State};
use axum::http::{HeaderMap, StatusCode};
use kanade_shared::wire::MetaEntry;
use serde::{Deserialize, Serialize};
use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool};
use tracing::{info, warn};
use crate::api::AppState;
use crate::audit::{self, Caller};
#[derive(Serialize)]
pub struct AgentRow {
pub pc_id: String,
pub hostname: Option<String>,
pub os_family: Option<String>,
pub agent_version: Option<String>,
pub last_heartbeat: Option<chrono::DateTime<chrono::Utc>>,
pub updated_at: Option<chrono::DateTime<chrono::Utc>>,
pub agent_cpu_pct: Option<f64>,
pub agent_rss_bytes: Option<i64>,
pub agent_disk_read_bytes: Option<i64>,
pub agent_disk_written_bytes: Option<i64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub quarantined_versions: Vec<String>,
pub last_logon_user: Option<String>,
pub last_logon_display_name: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub meta: Vec<MetaEntry>,
}
#[derive(Debug, Default, Deserialize)]
pub struct ListParams {
pub q: Option<String>,
pub user: Option<String>,
pub version: Option<String>,
pub quarantined: Option<String>,
pub limit: Option<u32>,
pub offset: Option<u32>,
pub status: Option<String>,
pub meta_key: Option<String>,
pub meta_value: Option<String>,
pub meta_any: Option<String>,
pub sort: Option<String>,
pub dir: Option<String>,
}
pub const ALIVE_THRESHOLD: chrono::Duration = chrono::Duration::minutes(2);
const MAX_FETCH: i64 = 10_000;
fn quarantined_like(version: Option<&str>) -> Option<String> {
version
.map(str::trim)
.filter(|s| !s.is_empty())
.map(|s| format!("%\"{}\"%", escape_like(s)))
}
fn escape_like(s: &str) -> String {
s.replace('\\', "\\\\")
.replace('%', "\\%")
.replace('_', "\\_")
}
fn contains_like(value: &str) -> String {
format!("%{}%", escape_like(value))
}
fn starts_like(value: &str) -> String {
format!("{}%", escape_like(value))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum MetaOp {
Contains,
Eq,
Neq,
Starts,
Empty,
Set,
Absent,
}
impl MetaOp {
fn parse(s: &str) -> Option<Self> {
Some(match s {
"contains" => Self::Contains,
"eq" => Self::Eq,
"neq" => Self::Neq,
"starts" => Self::Starts,
"empty" => Self::Empty,
"set" => Self::Set,
"absent" => Self::Absent,
_ => return None,
})
}
fn wants_value(self) -> bool {
matches!(self, Self::Contains | Self::Eq | Self::Neq | Self::Starts)
}
}
#[derive(Debug, Clone)]
struct MetaCond {
key: String,
op: MetaOp,
value: String,
}
fn parse_meta_conditions(raw: &[(String, String)]) -> Result<Vec<MetaCond>, (StatusCode, String)> {
let mut out = Vec::new();
for (k, v) in raw {
let Some(rest) = k.strip_prefix("meta.") else {
continue;
};
let Some((key, op_str)) = rest.rsplit_once('.') else {
return Err((
StatusCode::BAD_REQUEST,
format!("malformed meta filter `{k}` — expected meta.<key>.<op>"),
));
};
let key = key.trim();
if key.is_empty() {
return Err((
StatusCode::BAD_REQUEST,
format!("meta filter `{k}` has an empty key"),
));
}
let op = MetaOp::parse(op_str).ok_or_else(|| {
(
StatusCode::BAD_REQUEST,
format!("unknown meta filter operator `{op_str}` in `{k}`"),
)
})?;
let value = v.trim().to_string();
if op.wants_value() && value.is_empty() {
continue; }
out.push(MetaCond {
key: key.to_string(),
op,
value,
});
}
Ok(out)
}
fn push_meta_cond(qb: &mut QueryBuilder<Sqlite>, started: &mut bool, cond: &MetaCond) {
qb.push(sep(started));
match cond.op {
MetaOp::Absent => {
qb.push("pc_id NOT IN (SELECT pc_id FROM agent_meta WHERE key = ")
.push_bind(cond.key.clone())
.push(")");
}
MetaOp::Set => {
qb.push("pc_id IN (SELECT pc_id FROM agent_meta WHERE key = ")
.push_bind(cond.key.clone())
.push(")");
}
MetaOp::Empty => {
qb.push("pc_id IN (SELECT pc_id FROM agent_meta WHERE key = ")
.push_bind(cond.key.clone())
.push(" AND value = '')");
}
MetaOp::Eq | MetaOp::Neq | MetaOp::Contains | MetaOp::Starts => {
qb.push("pc_id IN (SELECT pc_id FROM agent_meta WHERE key = ")
.push_bind(cond.key.clone());
match cond.op {
MetaOp::Eq => {
qb.push(" AND value = ").push_bind(cond.value.clone());
}
MetaOp::Neq => {
qb.push(" AND value <> ").push_bind(cond.value.clone());
}
MetaOp::Contains => {
qb.push(" AND value LIKE ")
.push_bind(contains_like(&cond.value))
.push(" ESCAPE '\\'");
}
MetaOp::Starts => {
qb.push(" AND value LIKE ")
.push_bind(starts_like(&cond.value))
.push(" ESCAPE '\\'");
}
_ => unreachable!(),
}
qb.push(")");
}
}
}
fn push_filters(
qb: &mut QueryBuilder<Sqlite>,
started: &mut bool,
quar_like: &Option<String>,
meta_conds: &[MetaCond],
meta_any: &Option<String>,
) {
if let Some(p) = quar_like {
qb.push(sep(started))
.push("quarantined_versions LIKE ")
.push_bind(p.clone())
.push(" ESCAPE '\\'");
}
for cond in meta_conds {
push_meta_cond(qb, started, cond);
}
if let Some(text) = meta_any {
qb.push(sep(started))
.push("pc_id IN (SELECT pc_id FROM agent_meta WHERE value LIKE ")
.push_bind(contains_like(text))
.push(" ESCAPE '\\')");
}
}
#[derive(Debug, Clone)]
enum SortField {
Column(&'static str),
Meta(String),
}
struct SortSpec {
field: SortField,
desc: bool,
}
fn parse_sort(sort: Option<&str>, dir: Option<&str>) -> Result<SortSpec, (StatusCode, String)> {
let sort = sort.map(str::trim).filter(|s| !s.is_empty());
let field = match sort {
None => SortField::Column("updated_at"),
Some(s) => {
if let Some(key) = s.strip_prefix("meta:") {
let key = key.trim();
if key.is_empty() {
return Err((
StatusCode::BAD_REQUEST,
"sort=meta: requires a key (meta:<key>)".to_string(),
));
}
SortField::Meta(key.to_string())
} else {
let col = match s {
"pc_id" => "pc_id",
"hostname" => "hostname",
"os_family" | "os" => "os_family",
"agent_version" | "agent" => "agent_version",
"last_heartbeat" => "last_heartbeat",
"last_logon_user" | "last_logon" => "last_logon_user",
"updated_at" => "updated_at",
_ => {
return Err((StatusCode::BAD_REQUEST, format!("unknown sort field `{s}`")));
}
};
SortField::Column(col)
}
}
};
let desc = match dir.map(str::trim).filter(|s| !s.is_empty()) {
None => matches!(field, SortField::Column("updated_at")),
Some("asc") => false,
Some("desc") => true,
Some(other) => {
return Err((
StatusCode::BAD_REQUEST,
format!("sort direction must be asc/desc, got `{other}`"),
));
}
};
Ok(SortSpec { field, desc })
}
fn push_order_by(qb: &mut QueryBuilder<Sqlite>, spec: &SortSpec) {
let dir = if spec.desc { " DESC" } else { " ASC" };
qb.push(" ORDER BY ");
match &spec.field {
SortField::Column(col) => {
qb.push(col).push(" IS NULL, ").push(col).push(dir);
}
SortField::Meta(key) => {
qb.push("(SELECT value FROM agent_meta WHERE pc_id = agents.pc_id AND key = ")
.push_bind(key.clone())
.push(") IS NULL, (SELECT value FROM agent_meta WHERE pc_id = agents.pc_id AND key = ")
.push_bind(key.clone())
.push(") COLLATE NOCASE")
.push(dir);
}
}
}
fn sep(started: &mut bool) -> &'static str {
if *started {
" AND "
} else {
*started = true;
" WHERE "
}
}
const META_IN_CHUNK: usize = 500;
async fn attach_meta(pool: &SqlitePool, page: &mut [AgentRow]) -> Result<(), (StatusCode, String)> {
if page.is_empty() {
return Ok(());
}
let mut by_pc: std::collections::HashMap<String, Vec<MetaEntry>> =
std::collections::HashMap::new();
for chunk in page.chunks(META_IN_CHUNK) {
let mut qb: QueryBuilder<Sqlite> =
QueryBuilder::new("SELECT pc_id, key, value FROM agent_meta WHERE pc_id IN (");
{
let mut list = qb.separated(", ");
for a in chunk {
list.push_bind(a.pc_id.clone());
}
}
qb.push(") ORDER BY pc_id, key");
let rows = qb.build().fetch_all(pool).await.map_err(|e| {
warn!(error = %e, "attach agent_meta");
(
StatusCode::INTERNAL_SERVER_ERROR,
"load agent metadata failed".to_string(),
)
})?;
for r in rows {
let pc: String = r.try_get("pc_id").unwrap_or_default();
let key: String = r.try_get("key").unwrap_or_default();
let value: String = r.try_get("value").unwrap_or_default();
by_pc.entry(pc).or_default().push(MetaEntry { key, value });
}
}
for a in page.iter_mut() {
if let Some(m) = by_pc.remove(&a.pc_id) {
a.meta = m;
}
}
Ok(())
}
fn is_online(a: &AgentRow, cutoff: chrono::DateTime<chrono::Utc>) -> bool {
a.last_heartbeat.is_some_and(|hb| hb >= cutoff)
}
fn total_count(
needs_count: bool,
status: Option<&str>,
matched: i64,
online: i64,
fallback: i64,
) -> i64 {
if !needs_count {
return fallback;
}
match status {
Some("online") => online,
Some("offline") => matched - online,
_ => matched,
}
}
fn build_headers(needs_count: bool, total: i64, matched: i64, online: i64) -> HeaderMap {
let mut headers = HeaderMap::new();
if let Ok(v) = total.to_string().parse() {
headers.insert("X-Total-Count", v);
}
if needs_count {
if let Ok(v) = online.to_string().parse() {
headers.insert("X-Online-Count", v);
}
if let Ok(v) = (matched - online).to_string().parse() {
headers.insert("X-Offline-Count", v);
}
}
headers
}
pub async fn list(
State(pool): State<SqlitePool>,
Query(params): Query<ListParams>,
Query(raw): Query<Vec<(String, String)>>,
) -> Result<(HeaderMap, Json<Vec<AgentRow>>), (StatusCode, String)> {
let status = match params.status.as_deref().map(str::trim) {
None | Some("") => None,
Some(s @ ("online" | "offline")) => Some(s.to_string()),
Some(_) => {
return Err((
StatusCode::BAD_REQUEST,
"status must be 'online' or 'offline'".to_string(),
));
}
};
let q_re = super::compile(params.q.as_deref())?;
let user_re = super::compile(params.user.as_deref())?;
let version_re = super::compile(params.version.as_deref())?;
let has_regex = q_re.is_some() || user_re.is_some() || version_re.is_some();
let quar_like = quarantined_like(params.quarantined.as_deref());
let mut meta_conds = parse_meta_conditions(&raw)?;
if let Some(key) = params
.meta_key
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
{
let value = params
.meta_value
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty());
meta_conds.push(MetaCond {
key: key.to_string(),
op: if value.is_some() {
MetaOp::Contains
} else {
MetaOp::Set
},
value: value.unwrap_or_default().to_string(),
});
}
let meta_any = params
.meta_any
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
.map(|s| s.to_string());
let sort_spec = parse_sort(params.sort.as_deref(), params.dir.as_deref())?;
let cutoff = chrono::Utc::now() - ALIVE_THRESHOLD;
let needs_count = params.limit.is_some();
if !has_regex {
let limit = params.limit.map(i64::from).unwrap_or(-1);
let offset = params.offset.map(i64::from).unwrap_or(0);
let (matched, online): (i64, i64) = if !needs_count {
(0, 0)
} else {
let mut qb: QueryBuilder<Sqlite> = QueryBuilder::new(
"SELECT COUNT(*) AS matched, CAST(COALESCE(SUM(CASE WHEN \
last_heartbeat IS NOT NULL AND last_heartbeat >= ",
);
qb.push_bind(cutoff)
.push(" THEN 1 ELSE 0 END), 0) AS INTEGER) AS online FROM agents");
let mut started = false;
push_filters(&mut qb, &mut started, &quar_like, &meta_conds, &meta_any);
let row = qb.build().fetch_one(&pool).await.map_err(|e| {
warn!(error = %e, "count agents");
(
StatusCode::INTERNAL_SERVER_ERROR,
"count agents failed".to_string(),
)
})?;
(
row.try_get("matched").unwrap_or(0),
row.try_get("online").unwrap_or(0),
)
};
let mut qb: QueryBuilder<Sqlite> = QueryBuilder::new("SELECT * FROM agents");
let mut started = false;
push_filters(&mut qb, &mut started, &quar_like, &meta_conds, &meta_any);
match status.as_deref() {
Some("online") => {
qb.push(sep(&mut started))
.push("last_heartbeat IS NOT NULL AND last_heartbeat >= ")
.push_bind(cutoff);
}
Some("offline") => {
qb.push(sep(&mut started))
.push("(last_heartbeat IS NULL OR last_heartbeat < ")
.push_bind(cutoff)
.push(")");
}
_ => {}
}
push_order_by(&mut qb, &sort_spec);
qb.push(" LIMIT ")
.push_bind(limit)
.push(" OFFSET ")
.push_bind(offset);
let rows = qb.build().fetch_all(&pool).await.map_err(|e| {
warn!(error = %e, "list agents");
(
StatusCode::INTERNAL_SERVER_ERROR,
"list agents failed".to_string(),
)
})?;
let mut page: Vec<AgentRow> = rows.into_iter().map(row_to_agent).collect();
attach_meta(&pool, &mut page).await?;
let total = total_count(
needs_count,
status.as_deref(),
matched,
online,
offset + page.len() as i64,
);
return Ok((
build_headers(needs_count, total, matched, online),
Json(page),
));
}
let mut qb: QueryBuilder<Sqlite> = QueryBuilder::new("SELECT * FROM agents");
let mut started = false;
push_filters(&mut qb, &mut started, &quar_like, &meta_conds, &meta_any);
push_order_by(&mut qb, &sort_spec);
qb.push(" LIMIT ").push_bind(MAX_FETCH);
let rows = qb.build().fetch_all(&pool).await.map_err(|e| {
warn!(error = %e, "list agents");
(
StatusCode::INTERNAL_SERVER_ERROR,
"list agents failed".to_string(),
)
})?;
if rows.len() as i64 >= MAX_FETCH {
warn!(
cap = MAX_FETCH,
"agents regex prefilter hit the fetch cap; results may be truncated"
);
}
let matched_rows: Vec<AgentRow> = rows
.into_iter()
.filter(|r| {
if let Some(re) = &q_re {
let pc: &str = r.try_get("pc_id").unwrap_or("");
let host: &str = r.try_get("hostname").unwrap_or("");
if !(re.is_match(pc) || re.is_match(host)) {
return false;
}
}
if let Some(re) = &user_re {
let user: &str = r.try_get("last_logon_user").unwrap_or("");
let display: &str = r.try_get("last_logon_display_name").unwrap_or("");
if !(re.is_match(user) || re.is_match(display)) {
return false;
}
}
if let Some(re) = &version_re {
let version: &str = r.try_get("agent_version").unwrap_or("");
if !re.is_match(version) {
return false;
}
}
true
})
.map(row_to_agent)
.collect();
let matched = matched_rows.len() as i64;
let online = matched_rows.iter().filter(|a| is_online(a, cutoff)).count() as i64;
let offset = params.offset.unwrap_or(0) as usize;
let take = params.limit.map(|n| n as usize).unwrap_or(usize::MAX);
let mut page: Vec<AgentRow> = matched_rows
.into_iter()
.filter(|a| match status.as_deref() {
Some("online") => is_online(a, cutoff),
Some("offline") => !is_online(a, cutoff),
_ => true,
})
.skip(offset)
.take(take)
.collect();
attach_meta(&pool, &mut page).await?;
let total = total_count(
needs_count,
status.as_deref(),
matched,
online,
offset as i64 + page.len() as i64,
);
Ok((
build_headers(needs_count, total, matched, online),
Json(page),
))
}
pub async fn detail(
State(pool): State<SqlitePool>,
Path(pc_id): Path<String>,
) -> Result<Json<AgentRow>, StatusCode> {
let row = sqlx::query("SELECT * FROM agents WHERE pc_id = ?")
.bind(&pc_id)
.fetch_optional(&pool)
.await
.map_err(|e| {
warn!(error = %e, "detail agent");
StatusCode::INTERNAL_SERVER_ERROR
})?;
match row {
Some(r) => Ok(Json(row_to_agent(r))),
None => Err(StatusCode::NOT_FOUND),
}
}
pub async fn delete(
State(s): State<AppState>,
caller: Caller,
Path(pc_id): Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
let res = sqlx::query("DELETE FROM agents WHERE pc_id = ?")
.bind(&pc_id)
.execute(&s.pool)
.await
.map_err(|e| {
warn!(error = %e, pc_id = %pc_id, "delete agent");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("delete agent {pc_id}: {e}"),
)
})?;
if res.rows_affected() > 0 {
info!(pc_id = %pc_id, "agent deleted from registry");
audit::record(
&s.nats,
"operator",
"agent_delete",
Some(&pc_id),
Some(&caller),
serde_json::json!({ "pc_id": pc_id }),
)
.await;
}
Ok(StatusCode::NO_CONTENT)
}
fn row_to_agent(r: sqlx::sqlite::SqliteRow) -> AgentRow {
AgentRow {
pc_id: r.try_get("pc_id").unwrap_or_default(),
hostname: r.try_get("hostname").ok(),
os_family: r.try_get("os_family").ok(),
agent_version: r.try_get("agent_version").ok(),
last_heartbeat: r.try_get("last_heartbeat").ok(),
updated_at: r.try_get("updated_at").ok(),
agent_cpu_pct: r.try_get("agent_cpu_pct").ok(),
agent_rss_bytes: r.try_get("agent_rss_bytes").ok(),
agent_disk_read_bytes: r.try_get("agent_disk_read_bytes").ok(),
agent_disk_written_bytes: r.try_get("agent_disk_written_bytes").ok(),
quarantined_versions: r
.try_get::<Option<String>, _>("quarantined_versions")
.ok()
.flatten()
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default(),
last_logon_user: r
.try_get::<Option<String>, _>("last_logon_user")
.ok()
.flatten()
.filter(|s| !s.is_empty()),
last_logon_display_name: r
.try_get::<Option<String>, _>("last_logon_display_name")
.ok()
.flatten()
.filter(|s| !s.is_empty()),
meta: Vec::new(),
}
}
pub async fn meta_keys(State(pool): State<SqlitePool>) -> Result<Json<Vec<String>>, StatusCode> {
let rows = sqlx::query("SELECT DISTINCT key FROM agent_meta ORDER BY key")
.fetch_all(&pool)
.await
.map_err(|e| {
warn!(error = %e, "agent meta keys");
StatusCode::INTERNAL_SERVER_ERROR
})?;
Ok(Json(
rows.into_iter()
.map(|r| r.try_get::<String, _>("key").unwrap_or_default())
.collect(),
))
}
#[derive(Serialize)]
pub struct VersionCount {
pub version: Option<String>,
pub total: i64,
pub active: i64,
}
pub async fn versions(
State(pool): State<SqlitePool>,
) -> Result<Json<Vec<VersionCount>>, StatusCode> {
let rows = sqlx::query("SELECT agent_version, last_heartbeat FROM agents")
.fetch_all(&pool)
.await
.map_err(|e| {
warn!(error = %e, "agent versions");
StatusCode::INTERNAL_SERVER_ERROR
})?;
let cutoff = chrono::Utc::now() - ALIVE_THRESHOLD;
let mut buckets: std::collections::HashMap<String, (i64, i64)> =
std::collections::HashMap::new();
for r in rows {
let version: String = r
.try_get::<Option<String>, _>("agent_version")
.ok()
.flatten()
.unwrap_or_default();
let alive = r
.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("last_heartbeat")
.ok()
.flatten()
.is_some_and(|hb| hb >= cutoff);
let entry = buckets.entry(version).or_insert((0, 0));
entry.0 += 1;
if alive {
entry.1 += 1;
}
}
let mut out: Vec<VersionCount> = buckets
.into_iter()
.map(|(version, (total, active))| VersionCount {
version: (!version.is_empty()).then_some(version),
total,
active,
})
.collect();
out.sort_by(|a, b| {
b.total
.cmp(&a.total)
.then_with(|| a.version.cmp(&b.version))
});
Ok(Json(out))
}
#[cfg(test)]
mod tests {
use super::*;
use sqlx::sqlite::SqlitePoolOptions;
async fn seeded_pool() -> SqlitePool {
let pool = SqlitePoolOptions::new()
.max_connections(1)
.connect("sqlite::memory:")
.await
.unwrap();
sqlx::migrate!("./migrations").run(&pool).await.unwrap();
for (pc, host) in [
("PC001", "alpha"),
("PC002", "beta"),
("WS-9", "gamma"),
("web%01", "delta"),
] {
sqlx::query("INSERT INTO agents (pc_id, hostname) VALUES (?, ?)")
.bind(pc)
.bind(host)
.execute(&pool)
.await
.unwrap();
}
pool
}
async fn call(
pool: SqlitePool,
params: ListParams,
) -> Result<(HeaderMap, Json<Vec<AgentRow>>), (StatusCode, String)> {
list(State(pool), Query(params), Query(Vec::new())).await
}
async fn call_raw(
pool: SqlitePool,
params: ListParams,
raw: &[(&str, &str)],
) -> Result<(HeaderMap, Json<Vec<AgentRow>>), (StatusCode, String)> {
let raw: Vec<(String, String)> = raw
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect();
list(State(pool), Query(params), Query(raw)).await
}
async fn ids_of(pool: SqlitePool, params: ListParams) -> Vec<String> {
let (_headers, Json(rows)) = call(pool, params).await.unwrap();
rows.into_iter().map(|r| r.pc_id).collect()
}
async fn ids(pool: SqlitePool, q: Option<&str>, limit: Option<u32>) -> Vec<String> {
ids_of(
pool,
ListParams {
q: q.map(Into::into),
limit,
..Default::default()
},
)
.await
}
#[tokio::test]
async fn quarantined_versions_decode_through_the_api() {
let pool = seeded_pool().await;
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = 'PC001'")
.bind(r#"["0.43.51","0.43.52"]"#)
.execute(&pool)
.await
.unwrap();
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = 'PC002'")
.bind("not json")
.execute(&pool)
.await
.unwrap();
let (_h, Json(rows)) = call(pool, ListParams::default()).await.unwrap();
let by_id = |id: &str| {
rows.iter()
.find(|r| r.pc_id == id)
.unwrap()
.quarantined_versions
.clone()
};
assert_eq!(by_id("PC001"), vec!["0.43.51", "0.43.52"]);
assert!(
by_id("PC002").is_empty(),
"malformed JSON → empty, not error"
);
assert!(by_id("WS-9").is_empty(), "NULL column → empty");
}
async fn set_heartbeat(pool: &SqlitePool, pc_id: &str, online: bool) {
let hb = if online {
chrono::Utc::now()
} else {
chrono::Utc::now() - chrono::Duration::hours(1)
};
sqlx::query("UPDATE agents SET last_heartbeat = ? WHERE pc_id = ?")
.bind(hb)
.bind(pc_id)
.execute(pool)
.await
.unwrap();
}
fn get_header(h: &HeaderMap, k: &str) -> i64 {
h.get(k)
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse().ok())
.unwrap_or_else(|| panic!("{k} header missing or unparseable"))
}
#[tokio::test]
async fn status_filter_is_server_side_and_counts_are_fleet_wide() {
let pool = seeded_pool().await;
set_heartbeat(&pool, "PC001", true).await;
set_heartbeat(&pool, "PC002", false).await;
let (headers, Json(rows)) = call(
pool,
ListParams {
limit: Some(2),
status: Some("offline".into()),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(rows.len(), 2);
assert!(rows.iter().all(|r| r.pc_id != "PC001"));
assert_eq!(get_header(&headers, "X-Total-Count"), 3);
assert_eq!(get_header(&headers, "X-Online-Count"), 1);
assert_eq!(get_header(&headers, "X-Offline-Count"), 3);
}
#[tokio::test]
async fn online_filter_returns_only_live_agents() {
let pool = seeded_pool().await;
set_heartbeat(&pool, "PC001", true).await;
let (headers, Json(rows)) = call(
pool,
ListParams {
limit: Some(10),
status: Some("online".into()),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(
rows.iter().map(|r| r.pc_id.as_str()).collect::<Vec<_>>(),
vec!["PC001"]
);
assert_eq!(get_header(&headers, "X-Total-Count"), 1);
}
#[tokio::test]
async fn invalid_status_is_a_bad_request() {
let pool = seeded_pool().await;
match call(
pool,
ListParams {
limit: Some(10),
status: Some("onlin".into()),
..Default::default()
},
)
.await
{
Err((code, _)) => assert_eq!(code, StatusCode::BAD_REQUEST),
Ok(_) => panic!("a typo'd status must be a 400, not silently 'all'"),
}
}
#[tokio::test]
async fn invalid_regex_is_a_bad_request() {
let pool = seeded_pool().await;
match call(
pool,
ListParams {
q: Some("[unterminated".into()),
..Default::default()
},
)
.await
{
Err((code, _)) => assert_eq!(code, StatusCode::BAD_REQUEST),
Ok(_) => panic!("an invalid regex must be a 400"),
}
}
#[tokio::test]
async fn offset_pages_and_total_header_reports_match_count() {
let pool = seeded_pool().await;
let (headers, Json(page2)) = call(
pool,
ListParams {
limit: Some(1),
offset: Some(1),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(page2.len(), 1);
assert_eq!(
get_header(&headers, "X-Total-Count"),
4,
"seeded fleet has exactly four agents"
);
}
#[tokio::test]
async fn no_query_returns_whole_fleet() {
let got = ids(seeded_pool().await, None, None).await;
assert_eq!(got.len(), 4);
}
#[tokio::test]
async fn blank_query_is_treated_as_no_filter() {
let got = ids(seeded_pool().await, Some(" "), None).await;
assert_eq!(got.len(), 4);
}
#[tokio::test]
async fn q_is_a_regex_over_pc_id() {
let mut got = ids(seeded_pool().await, Some("^PC00"), None).await;
got.sort();
assert_eq!(got, vec!["PC001".to_string(), "PC002".to_string()]);
}
#[tokio::test]
async fn q_alternation_matches_pc_id_or_hostname() {
let mut got = ids(seeded_pool().await, Some("PC002|gamma"), None).await;
got.sort();
assert_eq!(got, vec!["PC002".to_string(), "WS-9".to_string()]);
}
#[tokio::test]
async fn q_matches_hostname_too() {
let got = ids(seeded_pool().await, Some("^alpha$"), None).await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn user_regex_matches_either_logon_field() {
let pool = seeded_pool().await;
sqlx::query(
"UPDATE agents SET last_logon_user = ?, last_logon_display_name = ? WHERE pc_id = 'PC001'",
)
.bind(r"CORP\taro")
.bind("Yamada Taro")
.execute(&pool)
.await
.unwrap();
let got = ids_of(
pool.clone(),
ListParams {
user: Some("Yamada".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
let got = ids_of(
pool,
ListParams {
user: Some(r"taro".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn version_regex_filters_agent_version() {
let pool = seeded_pool().await;
sqlx::query("UPDATE agents SET agent_version = ? WHERE pc_id = 'PC001'")
.bind("0.43.62")
.execute(&pool)
.await
.unwrap();
sqlx::query("UPDATE agents SET agent_version = ? WHERE pc_id = 'PC002'")
.bind("0.43.61")
.execute(&pool)
.await
.unwrap();
let got = ids_of(
pool,
ListParams {
version: Some(r"^0\.43\.62$".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn quarantined_filter_pre_filters_by_version_token() {
let pool = seeded_pool().await;
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = 'PC001'")
.bind(r#"["0.43.62"]"#)
.execute(&pool)
.await
.unwrap();
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = 'PC002'")
.bind(r#"["0.43.61"]"#)
.execute(&pool)
.await
.unwrap();
let got = ids_of(
pool,
ListParams {
quarantined: Some("0.43.62".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn quarantined_token_match_is_not_a_substring_match() {
let pool = seeded_pool().await;
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = 'PC001'")
.bind(r#"["0.43.62"]"#)
.execute(&pool)
.await
.unwrap();
let got = ids_of(
pool,
ListParams {
quarantined: Some("0.43.6".into()),
..Default::default()
},
)
.await;
assert!(got.is_empty(), "0.43.6 must not match the 0.43.62 token");
}
#[tokio::test]
async fn quarantined_combines_with_regex_and_counts() {
let pool = seeded_pool().await;
for pc in ["PC001", "PC002"] {
sqlx::query("UPDATE agents SET quarantined_versions = ? WHERE pc_id = ?")
.bind(r#"["0.43.62"]"#)
.bind(pc)
.execute(&pool)
.await
.unwrap();
}
sqlx::query("UPDATE agents SET agent_version = ? WHERE pc_id = 'PC001'")
.bind("0.43.61")
.execute(&pool)
.await
.unwrap();
let (headers, Json(rows)) = call(
pool,
ListParams {
quarantined: Some("0.43.62".into()),
version: Some("0.43.61".into()),
limit: Some(10),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(
rows.iter().map(|r| r.pc_id.as_str()).collect::<Vec<_>>(),
vec!["PC001"]
);
assert_eq!(get_header(&headers, "X-Total-Count"), 1);
}
#[tokio::test]
async fn limit_caps_row_count() {
let got = ids(seeded_pool().await, None, Some(2)).await;
assert_eq!(got.len(), 2);
}
#[tokio::test]
async fn empty_last_logon_display_name_normalises_to_none() {
let pool = seeded_pool().await;
sqlx::query(
"UPDATE agents SET last_logon_user = ?, last_logon_display_name = ? WHERE pc_id = 'PC001'",
)
.bind(r".\yukimemi")
.bind("")
.execute(&pool)
.await
.unwrap();
let (_h, Json(rows)) = call(pool, ListParams::default()).await.unwrap();
let a = rows.iter().find(|r| r.pc_id == "PC001").unwrap();
assert_eq!(a.last_logon_user.as_deref(), Some(r".\yukimemi"));
assert_eq!(
a.last_logon_display_name, None,
"empty display name must normalise to None"
);
}
async fn set_meta(pool: &SqlitePool, pc_id: &str, pairs: &[(&str, &str)]) {
for (k, v) in pairs {
sqlx::query("INSERT INTO agent_meta (pc_id, key, value) VALUES (?, ?, ?)")
.bind(pc_id)
.bind(k)
.bind(v)
.execute(pool)
.await
.unwrap();
}
}
#[tokio::test]
async fn meta_is_decorated_onto_rows() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("氏名", "山田"), ("所属", "経理")]).await;
let (_h, Json(rows)) = call(pool, ListParams::default()).await.unwrap();
let pc001 = rows.iter().find(|r| r.pc_id == "PC001").unwrap();
let map: std::collections::HashMap<_, _> = pc001
.meta
.iter()
.map(|e| (e.key.as_str(), e.value.as_str()))
.collect();
assert_eq!(map.get("氏名"), Some(&"山田"));
assert_eq!(map.get("所属"), Some(&"経理"));
assert!(
rows.iter()
.find(|r| r.pc_id == "PC002")
.unwrap()
.meta
.is_empty()
);
}
#[tokio::test]
async fn meta_key_filter_requires_the_key() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("所属", "経理")]).await;
set_meta(&pool, "PC002", &[("氏名", "田中")]).await;
let got = ids_of(
pool,
ListParams {
meta_key: Some("所属".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn meta_value_is_a_contains_filter() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("氏名", "山田太郎")]).await;
set_meta(&pool, "PC002", &[("氏名", "田中花子")]).await;
let got = ids_of(
pool.clone(),
ListParams {
meta_key: Some("氏名".into()),
meta_value: Some("山田".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
let got = ids_of(
pool,
ListParams {
meta_key: Some("氏名".into()),
meta_value: Some("鈴木".into()),
..Default::default()
},
)
.await;
assert!(got.is_empty());
}
#[tokio::test]
async fn meta_filter_counts_are_correct_with_paging() {
let pool = seeded_pool().await;
for pc in ["PC001", "PC002", "WS-9"] {
set_meta(&pool, pc, &[("dept", "eng")]).await;
}
let (headers, Json(rows)) = call(
pool,
ListParams {
meta_key: Some("dept".into()),
meta_value: Some("eng".into()),
limit: Some(2),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(rows.len(), 2, "page capped by limit");
assert_eq!(
get_header(&headers, "X-Total-Count"),
3,
"count reflects the meta filter"
);
}
#[tokio::test]
async fn meta_filter_composes_with_regex_path() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("dept", "eng")]).await;
set_meta(&pool, "PC002", &[("dept", "sales")]).await;
let got = ids_of(
pool,
ListParams {
q: Some("^PC".into()),
meta_key: Some("dept".into()),
meta_value: Some("eng".into()),
..Default::default()
},
)
.await;
assert_eq!(got, vec!["PC001".to_string()]);
}
#[tokio::test]
async fn attach_meta_chunks_past_the_bind_parameter_limit() {
let pool = seeded_pool().await;
let n = META_IN_CHUNK + 50;
for i in 0..n {
sqlx::query("INSERT INTO agents (pc_id) VALUES (?)")
.bind(format!("BULK-{i:04}"))
.execute(&pool)
.await
.unwrap();
}
set_meta(&pool, "BULK-0300", &[("dept", "eng")]).await;
let (_h, Json(rows)) = call(pool, ListParams::default()).await.unwrap();
assert!(
rows.len() >= n,
"whole fleet returned without a bind-limit error"
);
let bulk = rows.iter().find(|r| r.pc_id == "BULK-0300").unwrap();
assert_eq!(bulk.meta.first().map(|e| e.key.as_str()), Some("dept"));
}
#[tokio::test]
async fn meta_keys_returns_distinct_sorted_keys() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("氏名", "a"), ("所属", "b")]).await;
set_meta(&pool, "PC002", &[("氏名", "c"), ("メール", "d")]).await;
let Json(keys) = meta_keys(State(pool)).await.unwrap();
assert_eq!(keys, vec!["メール", "所属", "氏名"]);
}
fn pcids(rows: Vec<AgentRow>) -> Vec<String> {
rows.into_iter().map(|r| r.pc_id).collect()
}
async fn raw_ids(pool: SqlitePool, raw: &[(&str, &str)]) -> Vec<String> {
let (_h, Json(rows)) = call_raw(pool, ListParams::default(), raw).await.unwrap();
let mut ids = pcids(rows);
ids.sort();
ids
}
#[tokio::test]
async fn meta_ops_eq_contains_starts_neq() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("dept", "Finance")]).await;
set_meta(&pool, "PC002", &[("dept", "Engineering")]).await;
assert_eq!(
raw_ids(pool.clone(), &[("meta.dept.eq", "Finance")]).await,
vec!["PC001"]
);
assert_eq!(
raw_ids(pool.clone(), &[("meta.dept.contains", "ineer")]).await,
vec!["PC002"]
);
assert_eq!(
raw_ids(pool.clone(), &[("meta.dept.starts", "Fin")]).await,
vec!["PC001"]
);
assert_eq!(
raw_ids(pool, &[("meta.dept.neq", "Finance")]).await,
vec!["PC002"]
);
}
#[tokio::test]
async fn meta_ops_set_empty_absent() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("dept", "Finance")]).await;
set_meta(&pool, "PC002", &[("dept", "")]).await; assert_eq!(
raw_ids(pool.clone(), &[("meta.dept.set", "")]).await,
vec!["PC001", "PC002"]
);
assert_eq!(
raw_ids(pool.clone(), &[("meta.dept.empty", "")]).await,
vec!["PC002"]
);
assert_eq!(
raw_ids(pool, &[("meta.dept.absent", "")]).await,
vec!["WS-9", "web%01"]
);
}
#[tokio::test]
async fn multiple_meta_conditions_are_anded() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("dept", "eng"), ("site", "tokyo")]).await;
set_meta(&pool, "PC002", &[("dept", "eng"), ("site", "osaka")]).await;
assert_eq!(
raw_ids(pool, &[("meta.dept.eq", "eng"), ("meta.site.eq", "tokyo")]).await,
vec!["PC001"]
);
}
#[tokio::test]
async fn meta_any_matches_any_value() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("氏名", "山田太郎")]).await;
set_meta(&pool, "PC002", &[("所属", "山田製作所")]).await;
set_meta(&pool, "WS-9", &[("メール", "a@example.com")]).await;
let (_h, Json(rows)) = call(
pool,
ListParams {
meta_any: Some("山田".into()),
..Default::default()
},
)
.await
.unwrap();
let mut got = pcids(rows);
got.sort();
assert_eq!(got, vec!["PC001", "PC002"]);
}
#[tokio::test]
async fn unknown_meta_op_is_a_bad_request() {
let pool = seeded_pool().await;
match call_raw(pool, ListParams::default(), &[("meta.dept.bogus", "x")]).await {
Err((code, _)) => assert_eq!(code, StatusCode::BAD_REQUEST),
Ok(_) => panic!("an unknown meta operator must be a 400"),
}
}
#[tokio::test]
async fn sort_by_column_asc_and_desc() {
let pool = seeded_pool().await;
for (pc, v) in [("PC001", "0.1"), ("PC002", "0.3"), ("WS-9", "0.2")] {
sqlx::query("UPDATE agents SET agent_version = ? WHERE pc_id = ?")
.bind(v)
.bind(pc)
.execute(&pool)
.await
.unwrap();
}
let (_h, Json(rows)) = call(
pool.clone(),
ListParams {
sort: Some("agent_version".into()),
dir: Some("asc".into()),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(pcids(rows), vec!["PC001", "WS-9", "PC002", "web%01"]);
let (_h, Json(rows)) = call(
pool,
ListParams {
sort: Some("agent_version".into()),
dir: Some("desc".into()),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(pcids(rows), vec!["PC002", "WS-9", "PC001", "web%01"]);
}
#[tokio::test]
async fn sort_by_meta_value_nulls_last() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("dept", "beta")]).await;
set_meta(&pool, "PC002", &[("dept", "alpha")]).await;
let (_h, Json(rows)) = call(
pool,
ListParams {
sort: Some("meta:dept".into()),
dir: Some("asc".into()),
limit: Some(2),
..Default::default()
},
)
.await
.unwrap();
assert_eq!(pcids(rows), vec!["PC002", "PC001"]);
}
#[tokio::test]
async fn invalid_sort_field_is_a_bad_request() {
let pool = seeded_pool().await;
match call(
pool,
ListParams {
sort: Some("drop table".into()),
..Default::default()
},
)
.await
{
Err((code, _)) => assert_eq!(code, StatusCode::BAD_REQUEST),
Ok(_) => panic!("an unknown sort field must be a 400"),
}
}
#[tokio::test]
async fn duplicate_conditions_on_same_key_are_both_applied() {
let pool = seeded_pool().await;
set_meta(&pool, "PC001", &[("tag", "alpha-beta")]).await;
set_meta(&pool, "PC002", &[("tag", "alpha-only")]).await;
let got = raw_ids(
pool,
&[
("meta.tag.contains", "alpha"),
("meta.tag.contains", "beta"),
],
)
.await;
assert_eq!(got, vec!["PC001"], "both contains must AND, not collapse");
}
#[tokio::test]
async fn explicit_updated_at_sort_defaults_desc() {
let pool = seeded_pool().await;
for (pc, t) in [
("PC001", "2026-01-01 00:00:00"),
("PC002", "2026-01-03 00:00:00"),
] {
sqlx::query("UPDATE agents SET updated_at = ? WHERE pc_id = ?")
.bind(t)
.bind(pc)
.execute(&pool)
.await
.unwrap();
}
let (_h, Json(rows)) = call(
pool,
ListParams {
sort: Some("updated_at".into()),
..Default::default()
},
)
.await
.unwrap();
let ids = pcids(rows);
let p2 = ids.iter().position(|x| x == "PC002").unwrap();
let p1 = ids.iter().position(|x| x == "PC001").unwrap();
assert!(p2 < p1, "explicit sort=updated_at should default DESC");
}
}