diff --git a/server/src/http/gateway/mcp.rs b/server/src/http/gateway/mcp.rs index bd315272..e3ab2a1b 100644 --- a/server/src/http/gateway/mcp.rs +++ b/server/src/http/gateway/mcp.rs @@ -29,6 +29,8 @@ const GATEWAY_TOOLS: &[&str] = &[ "prompt_permission", // One-shot container rebuild: build fresh image, swap container, preserve state. "project_rebuild", + // On-demand host + per-container disk/load/CPU/mem snapshot (story 1207). + "fleet_resources", // Read sled identity pins vs. live signed identity, and TOFU re-pin. "fleet_identity", ]; @@ -156,6 +158,31 @@ pub(crate) fn gateway_tool_definitions() -> Vec { "required": ["name"] } }), + json!({ + "name": "fleet_resources", + "description": "On-demand host and per-container resource snapshot: host disk free/total, load averages, core count, and memory; per-container CPU%/mem sourced via `docker stats`; and per-project `target/`+`.huskies/worktrees/` directory sizes sourced via `docker exec ... find` (bounded, skips node_modules/.git, cached briefly per container+path). Flags disk/load conditions past thresholds first so problems lead the response.", + "inputSchema": { + "type": "object", + "properties": { + "disk_warn_gb": { + "type": "integer", + "description": "Host free-disk warn threshold in GB (default: 50)." + }, + "disk_critical_gb": { + "type": "integer", + "description": "Host free-disk critical threshold in GB (default: 20)." + }, + "load_warn_per_core": { + "type": "number", + "description": "1-minute load average per core above which a warn flag fires (default: 1.0)." + }, + "load_critical_per_core": { + "type": "number", + "description": "1-minute load average per core above which a critical flag fires (default: 2.0)." + } + } + } + }), json!({ "name": "fleet_identity", "description": "Read mode (default): for every registered sled, report project, url, connected, the recorded pin (expected_node_id), the live signature-verified node_id from a signed challenge-response (never the unsigned /identity display field), and whether they match. Repin mode: capture a sled's live verified identity via TOFU and persist it as the new pin, refusing when the signature is missing or does not verify.", @@ -440,6 +467,7 @@ async fn handle_gateway_tool( "agents.list" => handle_agents_list_tool(id), "prompt_permission" => handle_prompt_permission_tool(params, state, id).await, "project_rebuild" => handle_project_rebuild_tool(params, state, id).await, + "fleet_resources" => handle_fleet_resources_tool(params, state, id).await, "fleet_identity" => handle_fleet_identity_tool(params, state, id).await, _ => JsonRpcResponse::error(id, -32601, format!("Unknown gateway tool: {tool_name}")), } @@ -943,6 +971,65 @@ async fn handle_project_rebuild_tool( ) } +/// Handle the `fleet_resources` gateway tool (story 1207). +/// +/// Collects host disk/load/cpu/mem, per-container CPU%/mem (via `docker +/// stats`), and per-project `target/`/`worktrees/` sizes (via `docker exec +/// ... find`) — all sourced gateway-side rather than requiring each project's +/// sled to self-report. Thresholds for the leading `flags` list are supplied +/// as optional arguments (see `gateway_tool_definitions`), defaulting to +/// `ResourceThresholds::default()`. +async fn handle_fleet_resources_tool( + params: &Value, + state: &GatewayState, + id: Option, +) -> JsonRpcResponse { + use crate::service::gateway::resources::ResourceThresholds; + + let args = params.get("arguments").unwrap_or(params); + let defaults = ResourceThresholds::default(); + let thresholds = ResourceThresholds { + disk_warn_gb: args + .get("disk_warn_gb") + .and_then(|v| v.as_u64()) + .unwrap_or(defaults.disk_warn_gb), + disk_critical_gb: args + .get("disk_critical_gb") + .and_then(|v| v.as_u64()) + .unwrap_or(defaults.disk_critical_gb), + load_warn_per_core: args + .get("load_warn_per_core") + .and_then(|v| v.as_f64()) + .unwrap_or(defaults.load_warn_per_core), + load_critical_per_core: args + .get("load_critical_per_core") + .and_then(|v| v.as_f64()) + .unwrap_or(defaults.load_critical_per_core), + }; + + let project_names: Vec = state.projects.read().await.keys().cloned().collect(); + + match crate::service::gateway::resources::io::collect_fleet_resources( + &state.config_dir, + &project_names, + &thresholds, + ) + .await + { + Ok(resources) => JsonRpcResponse::success( + id, + json!({ + "content": [{ + "type": "text", + "text": serde_json::to_string_pretty(&resources).unwrap_or_default() + }], + "resources": serde_json::to_value(&resources).unwrap_or(json!(null)), + }), + ), + Err(e) => JsonRpcResponse::error(id, -32603, format!("fleet_resources failed: {e}")), + } +} + /// Handle the `fleet_identity` gateway tool. /// /// Dispatches on the `action` argument: `"read"` (default) reports every diff --git a/server/src/service/gateway/mod.rs b/server/src/service/gateway/mod.rs index 8942b481..682df647 100644 --- a/server/src/service/gateway/mod.rs +++ b/server/src/service/gateway/mod.rs @@ -18,6 +18,8 @@ pub(crate) mod io; pub mod polling; /// Pure signed release-manifest verification (signature, sha256, rollback) — no I/O. pub mod release_manifest; +/// Fleet resource types and collection (`fleet_resources` MCP tool, story 1207). +pub mod resources; pub use aggregation::{ find_project_containing_story, format_aggregate_status_compact, format_identity_reports, diff --git a/server/src/service/gateway/resources/io.rs b/server/src/service/gateway/resources/io.rs new file mode 100644 index 00000000..79214708 --- /dev/null +++ b/server/src/service/gateway/resources/io.rs @@ -0,0 +1,414 @@ +//! Side effects for `fleet_resources` (story 1207): host stat reads (statvfs, +//! getloadavg, `/proc/meminfo`), `docker stats`/`docker exec` subprocess +//! calls, and a short-TTL cache for the `docker exec ... find` directory-size +//! probe. +//! +//! Host stat reads are genuinely blocking syscalls/file reads and run inside +//! `spawn_blocking`. `docker` subprocess calls use `tokio::process::Command`, +//! which is already async and non-blocking — wrapping an awaited async call +//! in `spawn_blocking` would be a no-op, so those are left as plain `.await`s +//! (matching the existing convention in `chat::transport::matrix::project_rebuild`). + +use std::collections::HashMap; +use std::ffi::CString; +use std::os::unix::ffi::OsStrExt as _; +use std::path::Path; +use std::sync::{Mutex, OnceLock}; +use std::time::{Duration, Instant}; + +use futures::future::Future; + +use super::{ + ContainerStats, HostStats, ProjectDirSizes, parse_docker_cpu_percent, parse_docker_mem_usage, + parse_meminfo, +}; + +// ── Host stats ─────────────────────────────────────────────────────────── + +/// Read free/total bytes on the filesystem containing `path` via `statvfs(2)`. +fn disk_stats(path: &Path) -> std::io::Result<(u64, u64)> { + let c_path = CString::new(path.as_os_str().as_bytes()) + .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?; + // SAFETY: `c_path` is a valid NUL-terminated string and `stat` is a + // zero-initialized `libc::statvfs` passed by mutable reference, matching + // the `statvfs(2)` contract. + let mut stat: libc::statvfs = unsafe { std::mem::zeroed() }; + let ret = unsafe { libc::statvfs(c_path.as_ptr(), &mut stat) }; + if ret != 0 { + return Err(std::io::Error::last_os_error()); + } + let free = stat.f_bavail as u64 * stat.f_frsize as u64; + let total = stat.f_blocks as u64 * stat.f_frsize as u64; + Ok((free, total)) +} + +/// Read 1/5/15-minute load averages via `getloadavg(3)`. +fn load_averages() -> Option<(f64, f64, f64)> { + let mut loads = [0f64; 3]; + // SAFETY: `loads` has room for 3 entries, matching the `3` passed in. + let n = unsafe { libc::getloadavg(loads.as_mut_ptr(), 3) }; + if n == 3 { + Some((loads[0], loads[1], loads[2])) + } else { + None + } +} + +/// Read `MemTotal`/`MemAvailable` from `/proc/meminfo`. +fn mem_stats() -> Option<(u64, u64)> { + let contents = std::fs::read_to_string("/proc/meminfo").ok()?; + parse_meminfo(&contents) +} + +/// Collect all host-level stats. Runs on `spawn_blocking` since `statvfs`, +/// `getloadavg`, and the `/proc/meminfo` read are all blocking calls. +pub async fn collect_host_stats(workspace_root: &Path) -> Result { + let root = workspace_root.to_path_buf(); + tokio::task::spawn_blocking(move || { + let (disk_free_bytes, disk_total_bytes) = disk_stats(&root).map_err(|e| e.to_string())?; + let (load_avg_1, load_avg_5, load_avg_15) = load_averages().unwrap_or((0.0, 0.0, 0.0)); + let cpu_cores = std::thread::available_parallelism() + .map(|n| n.get()) + .unwrap_or(1); + let (mem_total_bytes, mem_available_bytes) = mem_stats().unwrap_or((0, 0)); + Ok(HostStats { + disk_free_bytes, + disk_total_bytes, + load_avg_1, + load_avg_5, + load_avg_15, + cpu_cores, + mem_total_bytes, + mem_available_bytes, + }) + }) + .await + .map_err(|e| format!("spawn_blocking panicked: {e}"))? +} + +// ── Docker container stats ────────────────────────────────────────────── + +/// The Docker container name for a registered project, following the +/// convention established by `project_rebuild.rs`. +fn container_name(project: &str) -> String { + format!("huskies-{project}") +} + +/// Query `docker stats --no-stream` for one container's CPU%/mem usage. +/// Tolerates failure (missing/stopped container, no Docker) by returning a +/// `ContainerStats` with `error` set rather than propagating an error. +async fn docker_stats_one(project: &str) -> ContainerStats { + let container = container_name(project); + let output = tokio::process::Command::new("docker") + .args([ + "stats", + "--no-stream", + "--format", + "{{.CPUPerc}}\t{{.MemUsage}}", + &container, + ]) + .output() + .await; + + match output { + Ok(out) if out.status.success() => { + let text = String::from_utf8_lossy(&out.stdout); + let line = text.lines().next().unwrap_or(""); + let mut parts = line.splitn(2, '\t'); + let cpu_percent = parts.next().and_then(parse_docker_cpu_percent); + let mem = parts.next().and_then(parse_docker_mem_usage); + ContainerStats { + project: project.to_string(), + container, + cpu_percent, + mem_usage_bytes: mem.map(|(used, _)| used), + mem_limit_bytes: mem.map(|(_, limit)| limit), + error: None, + } + } + Ok(out) => ContainerStats { + project: project.to_string(), + container, + cpu_percent: None, + mem_usage_bytes: None, + mem_limit_bytes: None, + error: Some(String::from_utf8_lossy(&out.stderr).trim().to_string()), + }, + Err(e) => ContainerStats { + project: project.to_string(), + container, + cpu_percent: None, + mem_usage_bytes: None, + mem_limit_bytes: None, + error: Some(format!("docker stats failed to spawn: {e}")), + }, + } +} + +/// Query `docker stats` for every project's container concurrently, +/// tolerating per-container failures (mirrors `gateway_health`'s per-project +/// tolerance) so one dead container never fails the whole call. +pub async fn collect_container_stats(project_names: &[String]) -> Vec { + use futures::future::join_all; + join_all(project_names.iter().map(|name| docker_stats_one(name))).await +} + +// ── Bounded `du` via `docker exec ... find`, TTL-cached ───────────────── + +/// How long a computed directory size is reused before being recomputed. +/// Bounds repeated `fleet_resources` calls in a short window from re-shelling +/// out to `docker exec ... find` on every call (story 1207 AC4). +const DU_CACHE_TTL: Duration = Duration::from_secs(30); + +/// Memoize an async, fallible value producer for [`DU_CACHE_TTL`], keyed by +/// `key`. `compute` is only invoked on a cache miss/expiry — errors are +/// returned directly without populating the cache, so a transient failure +/// doesn't poison later calls for [`DU_CACHE_TTL`]. +/// +/// Generic over the producer so it can be exercised in tests with a mocked +/// closure (no Docker required) while real callers pass a `docker exec` +/// invocation. +async fn cached_or_compute(key: String, compute: F) -> Result +where + F: FnOnce() -> Fut, + Fut: Future>, +{ + fn cache() -> &'static Mutex> { + static CACHE: OnceLock>> = OnceLock::new(); + CACHE.get_or_init(|| Mutex::new(HashMap::new())) + } + + if let Some((at, bytes)) = cache().lock().unwrap().get(&key) + && at.elapsed() < DU_CACHE_TTL + { + return Ok(*bytes); + } + + let bytes = compute().await?; + cache().lock().unwrap().insert(key, (Instant::now(), bytes)); + Ok(bytes) +} + +/// Sum file sizes under `path_in_container` (inside the named container) via +/// `docker exec find`, skipping `node_modules`/`.git` entirely so +/// the walk stays bounded on large worktree trees. Returns `0` (rather than +/// erroring) when the path doesn't exist in the container — `find` reports +/// that on stderr with a non-zero exit but partial/empty stdout is still +/// valid, matching `dir_size_bytes`'s "missing path returns 0" contract. +async fn du_via_docker_exec(container: &str, path_in_container: &str) -> Result { + let output = tokio::process::Command::new("docker") + .args([ + "exec", + container, + "find", + path_in_container, + "-type", + "f", + "-not", + "-path", + "*/node_modules/*", + "-not", + "-path", + "*/.git/*", + "-printf", + "%s\n", + ]) + .output() + .await + .map_err(|e| format!("docker exec failed to spawn: {e}"))?; + + let bytes = String::from_utf8_lossy(&output.stdout) + .lines() + .filter_map(|l| l.trim().parse::().ok()) + .sum(); + Ok(bytes) +} + +/// Compute `target/` and `.huskies/worktrees/` sizes for one project's +/// container, each cached independently for [`DU_CACHE_TTL`]. +async fn project_dir_sizes_one(project: &str) -> ProjectDirSizes { + let container = container_name(project); + + let target_key = format!("{container}:/workspace/target"); + let target_container = container.clone(); + let target_result = cached_or_compute(target_key, || async move { + du_via_docker_exec(&target_container, "/workspace/target").await + }) + .await; + + let worktrees_key = format!("{container}:/workspace/.huskies/worktrees"); + let worktrees_container = container.clone(); + let worktrees_result = cached_or_compute(worktrees_key, || async move { + du_via_docker_exec(&worktrees_container, "/workspace/.huskies/worktrees").await + }) + .await; + + match (target_result, worktrees_result) { + (Ok(target_bytes), Ok(worktrees_bytes)) => ProjectDirSizes { + project: project.to_string(), + target_bytes: Some(target_bytes), + worktrees_bytes: Some(worktrees_bytes), + error: None, + }, + (Err(e), _) | (_, Err(e)) => ProjectDirSizes { + project: project.to_string(), + target_bytes: None, + worktrees_bytes: None, + error: Some(e), + }, + } +} + +/// Compute per-project directory sizes for every project concurrently, +/// tolerating per-project failures. +pub async fn collect_project_dir_sizes(project_names: &[String]) -> Vec { + use futures::future::join_all; + join_all(project_names.iter().map(|name| project_dir_sizes_one(name))).await +} + +// ── Orchestration ──────────────────────────────────────────────────────── + +/// Collect the full `fleet_resources` response: host stats, per-container +/// CPU/mem, per-project directory sizes, and threshold flags. +/// +/// `workspace_root` is the gateway's own workspace path (used for host disk +/// stats); `project_names` are the currently registered `projects.toml` keys. +pub async fn collect_fleet_resources( + workspace_root: &Path, + project_names: &[String], + thresholds: &super::ResourceThresholds, +) -> Result { + let host = collect_host_stats(workspace_root).await?; + let (containers, projects) = tokio::join!( + collect_container_stats(project_names), + collect_project_dir_sizes(project_names), + ); + let flags = super::flag_host(&host, thresholds); + + Ok(super::FleetResources { + flags, + host, + containers, + projects, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + + // ── collect_host_stats ────────────────────────────────────────────── + + #[tokio::test] + async fn collect_host_stats_reads_real_host() { + let stats = collect_host_stats(Path::new(".")).await.unwrap(); + assert!(stats.disk_total_bytes > 0); + assert!(stats.disk_free_bytes <= stats.disk_total_bytes); + assert!(stats.cpu_cores >= 1); + } + + #[tokio::test] + async fn collect_host_stats_missing_path_is_error() { + let result = collect_host_stats(Path::new("/no/such/path/1207")).await; + assert!(result.is_err()); + } + + // ── cached_or_compute (mocked producer, no Docker needed) ─────────── + + #[tokio::test] + async fn cached_or_compute_calls_producer_once_within_ttl() { + let calls = std::sync::Arc::new(AtomicUsize::new(0)); + let key = "test-key-1".to_string(); + + let make_producer = |calls: std::sync::Arc, value: u64| { + move || { + let calls = calls.clone(); + async move { + calls.fetch_add(1, Ordering::SeqCst); + Ok(value) + } + } + }; + + let first = cached_or_compute(key.clone(), make_producer(calls.clone(), 42)) + .await + .unwrap(); + let second = cached_or_compute(key.clone(), make_producer(calls.clone(), 99)) + .await + .unwrap(); + + assert_eq!(first, 42); + assert_eq!( + second, 42, + "second call within TTL must reuse the cached value" + ); + assert_eq!( + calls.load(Ordering::SeqCst), + 1, + "producer must only run once while the cache entry is still fresh — proves the \ + expensive computation itself (not just a pre-computed value) is skipped on a hit" + ); + } + + #[tokio::test] + async fn cached_or_compute_different_keys_do_not_share_a_cache_entry() { + let calls = std::sync::Arc::new(AtomicUsize::new(0)); + let make_producer = |calls: std::sync::Arc, value: u64| { + move || { + let calls = calls.clone(); + async move { + calls.fetch_add(1, Ordering::SeqCst); + Ok(value) + } + } + }; + + let a = cached_or_compute("key-a".to_string(), make_producer(calls.clone(), 1)) + .await + .unwrap(); + let b = cached_or_compute("key-b".to_string(), make_producer(calls.clone(), 2)) + .await + .unwrap(); + + assert_eq!(a, 1); + assert_eq!(b, 2); + assert_eq!(calls.load(Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn cached_or_compute_errors_are_not_cached() { + let calls = std::sync::Arc::new(AtomicUsize::new(0)); + let key = "test-key-err".to_string(); + + let calls_1 = calls.clone(); + let first = cached_or_compute(key.clone(), || async move { + calls_1.fetch_add(1, Ordering::SeqCst); + Err::("boom".to_string()) + }) + .await; + assert!(first.is_err()); + + let calls_2 = calls.clone(); + let second = cached_or_compute(key.clone(), || async move { + calls_2.fetch_add(1, Ordering::SeqCst); + Ok(7) + }) + .await + .unwrap(); + + assert_eq!(second, 7); + assert_eq!( + calls.load(Ordering::SeqCst), + 2, + "an error result must not populate the cache, so the next call recomputes" + ); + } + + // ── container_name ──────────────────────────────────────────────────── + + #[test] + fn container_name_follows_project_rebuild_convention() { + assert_eq!(container_name("myproj"), "huskies-myproj"); + } +} diff --git a/server/src/service/gateway/resources/mod.rs b/server/src/service/gateway/resources/mod.rs new file mode 100644 index 00000000..5a825d30 --- /dev/null +++ b/server/src/service/gateway/resources/mod.rs @@ -0,0 +1,429 @@ +//! Fleet resource types and pure logic for the `fleet_resources` MCP tool +//! (story 1207): host disk/load/cpu/mem, per-container CPU%/mem, per-project +//! `target/`+`worktrees/` sizes, and threshold-based flagging. +//! +//! Follows the `service/gateway` conventions: this file holds pure types and +//! classification logic (no I/O); `io.rs` performs all side effects (statvfs, +//! getloadavg, `/proc/meminfo`, and `docker` subprocess calls). + +/// Side effects for fleet resource collection: host stat reads, `docker +/// stats`/`docker exec` subprocess calls, and a TTL-cached `du` helper. +pub mod io; + +use serde::Serialize; + +/// Host machine (gateway node) resource snapshot. +#[derive(Debug, Clone, Serialize, PartialEq)] +pub struct HostStats { + /// Free bytes on the filesystem containing the gateway's workspace root. + pub disk_free_bytes: u64, + /// Total bytes on the filesystem containing the gateway's workspace root. + pub disk_total_bytes: u64, + /// 1-minute load average. + pub load_avg_1: f64, + /// 5-minute load average. + pub load_avg_5: f64, + /// 15-minute load average. + pub load_avg_15: f64, + /// Number of logical CPU cores available to the gateway process. + pub cpu_cores: usize, + /// Total physical memory in bytes. + pub mem_total_bytes: u64, + /// Available (not merely free) memory in bytes, per `/proc/meminfo`'s + /// `MemAvailable` estimate. + pub mem_available_bytes: u64, +} + +/// Per-container CPU/mem snapshot, sourced via `docker stats`. +/// +/// `error` is set (with all other fields `None`) when the container is +/// unreachable or `docker stats` fails — one dead container must not fail the +/// whole `fleet_resources` call. +#[derive(Debug, Clone, Serialize, PartialEq)] +pub struct ContainerStats { + /// Registered project name (key in `projects.toml`). + pub project: String, + /// Docker container name (`huskies-{project}`). + pub container: String, + /// CPU usage percentage, as reported by `docker stats`. + pub cpu_percent: Option, + /// Memory currently in use, in bytes. + pub mem_usage_bytes: Option, + /// Memory limit (cgroup limit or host total), in bytes. + pub mem_limit_bytes: Option, + /// Set when `docker stats` failed for this container. + pub error: Option, +} + +/// Per-project on-disk footprint: `target/` build output and +/// `.huskies/worktrees/` (coder agent worktrees), measured inside the +/// project's own container via `docker exec ... find` (see `io.rs`). +#[derive(Debug, Clone, Serialize, PartialEq)] +pub struct ProjectDirSizes { + /// Registered project name. + pub project: String, + /// Bytes used by `/workspace/target`, or `None` if it could not be measured. + pub target_bytes: Option, + /// Bytes used by `/workspace/.huskies/worktrees`, or `None` if it could + /// not be measured. + pub worktrees_bytes: Option, + /// Set when the `docker exec ... find` measurement failed. + pub error: Option, +} + +/// Threshold configuration for [`flag_host`]. Supplied per tool call (see +/// `handle_fleet_resources_tool`) rather than read from `projects.toml`, so an +/// operator can widen/narrow thresholds per invocation. +#[derive(Debug, Clone, PartialEq)] +pub struct ResourceThresholds { + /// Host free disk space (GB) below which a "warn" flag fires. + pub disk_warn_gb: u64, + /// Host free disk space (GB) below which a "critical" flag fires. + pub disk_critical_gb: u64, + /// 1-minute load average per core above which a "warn" flag fires. + pub load_warn_per_core: f64, + /// 1-minute load average per core above which a "critical" flag fires. + pub load_critical_per_core: f64, +} + +impl Default for ResourceThresholds { + fn default() -> Self { + Self { + disk_warn_gb: 50, + disk_critical_gb: 20, + load_warn_per_core: 1.0, + load_critical_per_core: 2.0, + } + } +} + +/// A single flagged resource condition, surfaced ahead of the raw readings so +/// problems "lead" the response (story 1207 AC3). +#[derive(Debug, Clone, Serialize, PartialEq)] +pub struct ResourceFlag { + /// Flag category: `"low_disk"` or `"high_load"`. + pub kind: String, + /// Severity: `"warn"` or `"critical"`. + pub level: String, + /// Human-readable detail (e.g. current reading vs. threshold). + pub detail: String, +} + +/// Full `fleet_resources` response payload. +#[derive(Debug, Clone, Serialize, PartialEq)] +pub struct FleetResources { + /// Flagged problems, first so they lead the response. + pub flags: Vec, + /// Host-level resource snapshot. + pub host: HostStats, + /// Per-container CPU/mem snapshots. + pub containers: Vec, + /// Per-project `target/`/`worktrees/` directory sizes. + pub projects: Vec, +} + +const BYTES_PER_GB: u64 = 1_000_000_000; + +/// Flag host disk/load conditions against `thresholds`. Pure — no I/O. +/// +/// Mirrors `disk_watch::classify_level`'s "below X" boundary semantics for +/// disk (free space strictly below the threshold triggers the flag) and uses +/// "at or above" for load (a load average right at the threshold already +/// indicates saturation, unlike free-space headroom). +pub fn flag_host(host: &HostStats, thresholds: &ResourceThresholds) -> Vec { + let mut flags = Vec::new(); + + if host.disk_free_bytes < thresholds.disk_critical_gb * BYTES_PER_GB { + flags.push(ResourceFlag { + kind: "low_disk".to_string(), + level: "critical".to_string(), + detail: format!( + "{:.1} GB free (critical threshold: {} GB)", + host.disk_free_bytes as f64 / BYTES_PER_GB as f64, + thresholds.disk_critical_gb + ), + }); + } else if host.disk_free_bytes < thresholds.disk_warn_gb * BYTES_PER_GB { + flags.push(ResourceFlag { + kind: "low_disk".to_string(), + level: "warn".to_string(), + detail: format!( + "{:.1} GB free (warn threshold: {} GB)", + host.disk_free_bytes as f64 / BYTES_PER_GB as f64, + thresholds.disk_warn_gb + ), + }); + } + + if host.cpu_cores > 0 { + let per_core = host.load_avg_1 / host.cpu_cores as f64; + if per_core >= thresholds.load_critical_per_core { + flags.push(ResourceFlag { + kind: "high_load".to_string(), + level: "critical".to_string(), + detail: format!( + "load1={:.2} across {} core(s) ({:.2}/core, critical threshold: {:.2}/core)", + host.load_avg_1, host.cpu_cores, per_core, thresholds.load_critical_per_core + ), + }); + } else if per_core >= thresholds.load_warn_per_core { + flags.push(ResourceFlag { + kind: "high_load".to_string(), + level: "warn".to_string(), + detail: format!( + "load1={:.2} across {} core(s) ({:.2}/core, warn threshold: {:.2}/core)", + host.load_avg_1, host.cpu_cores, per_core, thresholds.load_warn_per_core + ), + }); + } + } + + flags +} + +/// Parse a `docker stats` CPU percentage field like `"12.34%"`. +pub fn parse_docker_cpu_percent(s: &str) -> Option { + s.trim().trim_end_matches('%').parse::().ok() +} + +/// Parse a single `docker stats`-formatted size like `"12.3MiB"`, `"1.5GiB"`, +/// or `"512B"` into bytes. Handles the IEC-binary units `docker stats` uses. +pub fn parse_docker_size(s: &str) -> Option { + let s = s.trim(); + let unit_start = s.find(|c: char| !c.is_ascii_digit() && c != '.')?; + let (number, unit) = s.split_at(unit_start); + let value: f64 = number.parse().ok()?; + let multiplier: f64 = match unit { + "B" => 1.0, + "KiB" | "KB" => 1024.0, + "MiB" | "MB" => 1024.0 * 1024.0, + "GiB" | "GB" => 1024.0 * 1024.0 * 1024.0, + "TiB" | "TB" => 1024.0 * 1024.0 * 1024.0 * 1024.0, + _ => return None, + }; + Some((value * multiplier) as u64) +} + +/// Parse a `docker stats` `MemUsage` field like `"58.14MiB / 3.842GiB"` into +/// `(used_bytes, limit_bytes)`. +pub fn parse_docker_mem_usage(s: &str) -> Option<(u64, u64)> { + let (used, limit) = s.split_once(" / ")?; + Some((parse_docker_size(used)?, parse_docker_size(limit)?)) +} + +/// Parse `MemTotal`/`MemAvailable` (in kB) out of `/proc/meminfo` contents +/// into `(total_bytes, available_bytes)`. +pub fn parse_meminfo(contents: &str) -> Option<(u64, u64)> { + let mut total = None; + let mut available = None; + for line in contents.lines() { + if let Some(rest) = line.strip_prefix("MemTotal:") { + total = parse_meminfo_kb_field(rest); + } else if let Some(rest) = line.strip_prefix("MemAvailable:") { + available = parse_meminfo_kb_field(rest); + } + } + Some((total?, available?)) +} + +/// Parse a `/proc/meminfo` value field like `" 16384000 kB"` into bytes. +fn parse_meminfo_kb_field(s: &str) -> Option { + let kb: u64 = s.trim().trim_end_matches("kB").trim().parse().ok()?; + Some(kb * 1024) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn thresholds() -> ResourceThresholds { + ResourceThresholds { + disk_warn_gb: 50, + disk_critical_gb: 20, + load_warn_per_core: 1.0, + load_critical_per_core: 2.0, + } + } + + fn host(disk_free_gb: u64, load_avg_1: f64, cpu_cores: usize) -> HostStats { + HostStats { + disk_free_bytes: disk_free_gb * BYTES_PER_GB, + disk_total_bytes: 500 * BYTES_PER_GB, + load_avg_1, + load_avg_5: load_avg_1, + load_avg_15: load_avg_1, + cpu_cores, + mem_total_bytes: 16 * BYTES_PER_GB, + mem_available_bytes: 8 * BYTES_PER_GB, + } + } + + // ── flag_host: disk ────────────────────────────────────────────────── + + #[test] + fn flag_host_healthy_disk_and_load_produces_no_flags() { + let flags = flag_host(&host(100, 0.5, 4), &thresholds()); + assert!(flags.is_empty(), "expected no flags, got: {flags:?}"); + } + + #[test] + fn flag_host_below_warn_gb_flags_warn() { + let flags = flag_host(&host(49, 0.5, 4), &thresholds()); + assert_eq!(flags.len(), 1); + assert_eq!(flags[0].kind, "low_disk"); + assert_eq!(flags[0].level, "warn"); + } + + #[test] + fn flag_host_at_warn_boundary_is_ok() { + // Exactly at warn_gb is not below it, so no flag (matches disk_watch semantics). + let flags = flag_host(&host(50, 0.5, 4), &thresholds()); + assert!(flags.is_empty()); + } + + #[test] + fn flag_host_below_critical_gb_flags_critical_not_warn() { + let flags = flag_host(&host(19, 0.5, 4), &thresholds()); + assert_eq!(flags.len(), 1); + assert_eq!(flags[0].kind, "low_disk"); + assert_eq!(flags[0].level, "critical"); + } + + // ── flag_host: load ─────────────────────────────────────────────────── + + #[test] + fn flag_host_load_at_warn_per_core_flags_warn() { + // 4 cores, warn threshold 1.0/core => load1 == 4.0 triggers warn. + let flags = flag_host(&host(100, 4.0, 4), &thresholds()); + assert_eq!(flags.len(), 1); + assert_eq!(flags[0].kind, "high_load"); + assert_eq!(flags[0].level, "warn"); + } + + #[test] + fn flag_host_load_at_critical_per_core_flags_critical_not_warn() { + // 4 cores, critical threshold 2.0/core => load1 == 8.0 triggers critical only. + let flags = flag_host(&host(100, 8.0, 4), &thresholds()); + assert_eq!(flags.len(), 1); + assert_eq!(flags[0].kind, "high_load"); + assert_eq!(flags[0].level, "critical"); + } + + #[test] + fn flag_host_zero_cores_does_not_panic_or_flag_load() { + let flags = flag_host(&host(100, 10.0, 0), &thresholds()); + assert!(flags.iter().all(|f| f.kind != "high_load")); + } + + #[test] + fn flag_host_disk_and_load_both_flag_together() { + let flags = flag_host(&host(10, 10.0, 4), &thresholds()); + assert_eq!(flags.len(), 2); + assert!(flags.iter().any(|f| f.kind == "low_disk")); + assert!(flags.iter().any(|f| f.kind == "high_load")); + } + + // ── parse_docker_cpu_percent ────────────────────────────────────────── + + #[test] + fn parse_cpu_percent_typical() { + assert_eq!(parse_docker_cpu_percent("12.34%"), Some(12.34)); + } + + #[test] + fn parse_cpu_percent_zero() { + assert_eq!(parse_docker_cpu_percent("0.00%"), Some(0.0)); + } + + #[test] + fn parse_cpu_percent_invalid_is_none() { + assert_eq!(parse_docker_cpu_percent("n/a"), None); + } + + // ── parse_docker_size / parse_docker_mem_usage ──────────────────────── + + #[test] + fn parse_size_bytes_no_unit_suffix() { + assert_eq!(parse_docker_size("512B"), Some(512)); + } + + #[test] + fn parse_size_mib() { + assert_eq!(parse_docker_size("1MiB"), Some(1024 * 1024)); + } + + #[test] + fn parse_size_fractional_gib() { + assert_eq!( + parse_docker_size("1.5GiB"), + Some((1.5 * 1024.0 * 1024.0 * 1024.0) as u64) + ); + } + + #[test] + fn parse_size_unknown_unit_is_none() { + assert_eq!(parse_docker_size("3.2XB"), None); + } + + #[test] + fn parse_mem_usage_typical() { + let (used, limit) = parse_docker_mem_usage("58.14MiB / 3.842GiB").unwrap(); + assert_eq!(used, (58.14 * 1024.0 * 1024.0) as u64); + assert_eq!(limit, (3.842 * 1024.0 * 1024.0 * 1024.0) as u64); + } + + #[test] + fn parse_mem_usage_missing_separator_is_none() { + assert_eq!(parse_docker_mem_usage("58.14MiB"), None); + } + + // ── parse_meminfo ────────────────────────────────────────────────────── + + #[test] + fn parse_meminfo_extracts_total_and_available() { + let sample = "MemTotal: 16384000 kB\nMemFree: 1000000 kB\nMemAvailable: 8000000 kB\n"; + let (total, available) = parse_meminfo(sample).unwrap(); + assert_eq!(total, 16_384_000 * 1024); + assert_eq!(available, 8_000_000 * 1024); + } + + #[test] + fn parse_meminfo_missing_fields_is_none() { + assert_eq!(parse_meminfo("MemFree: 1000 kB\n"), None); + } + + // ── FleetResources shape ──────────────────────────────────────────────── + + #[test] + fn fleet_resources_serializes_expected_top_level_keys() { + let resources = FleetResources { + flags: vec![ResourceFlag { + kind: "low_disk".to_string(), + level: "warn".to_string(), + detail: "40.0 GB free".to_string(), + }], + host: host(40, 0.5, 4), + containers: vec![ContainerStats { + project: "huskies".to_string(), + container: "huskies-huskies".to_string(), + cpu_percent: Some(12.3), + mem_usage_bytes: Some(100), + mem_limit_bytes: Some(1000), + error: None, + }], + projects: vec![ProjectDirSizes { + project: "huskies".to_string(), + target_bytes: Some(123), + worktrees_bytes: Some(456), + error: None, + }], + }; + let value = serde_json::to_value(&resources).unwrap(); + assert!(value.get("flags").unwrap().is_array()); + assert!(value.get("host").unwrap().get("disk_free_bytes").is_some()); + assert!(value.get("containers").unwrap().is_array()); + assert!(value.get("projects").unwrap().is_array()); + assert_eq!(value["containers"][0]["cpu_percent"], 12.3); + assert_eq!(value["projects"][0]["target_bytes"], 123); + } +}