housekeeping-p2: security, correctness, and performance pass before 1.51 (#39)

* chore: dead-code sweep

- Delete socktop_connector/src/connector.rs: orphaned since 08f248c removed
  'pub mod connector;' during the modularization refactor. Never compiled
  (verified under default, wasm, and workspace feature combos) but shipped
  in the crates.io tarball and contained an outdated copy of the TLS
  verifier — a trap for anyone patching the pinning bug in the dead copy.
- Delete empty socktop/src/ws.rs, tracked editor backup ui/.modal.rs.backup,
  and stray test_thiserror.rs at the repo root.
- Delete the two LEGACY #[allow(dead_code)] process input handlers; the
  header-click render test now exercises the live _with_selection handler
  instead (better coverage of the real path).
- Drop unused sysinfo dependency from the socktop client.
- Replace stale 'temporarily increased for testing' comment on
  COMPRESSION_THRESHOLD (it already held the production value).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(connector): make certificate pinning real; disable Nagle

Security: with --verify-hostname off (the default), the old NoVerify
verifier accepted ANY server certificate — the CA loaded from --tls-ca
was never consulted, so the documented pinning was a no-op and the
connection was trivially MITM-able. Replace it with PinnedCertVerifier:
the presented end-entity cert must be byte-identical to a cert in the
--tls-ca file (any cert in a multi-cert PEM matches, supporting
rotation). Signature validation now uses the ring provider's full
algorithm set instead of a hardcoded 3-scheme list. Empty PEM files
fail fast instead of failing closed per-handshake.

The --verify-hostname path is unchanged (WebPki root-store validation).

Also: the third argument of connect_async_tls_with_config is
tungstenite's disable_nagle flag, not a verification toggle — we were
passing verify_hostname there, leaving Nagle ON for default users. Pass
true unconditionally, and disable Nagle on the plain ws:// path too;
socktop exchanges small request/response frames where Nagle only adds
latency.

Client now consumes the connector via a dual path+version dep so these
fixes are in local builds and CI before the crates.io publish (cargo
strips the path on publish). Connector version -> 1.51.0.

Verified E2E: agent A's cert connects to agent A; agent B's cert
against agent A fails the handshake (the rpi-worker-1 wrong-PEM
scenario); --verify-hostname against a 127.0.0.1 SAN still connects.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(agent): GPU worker thread, async journalctl, correctness + cache fixes

Lightweight:
- GPU collection moves to a dedicated worker thread that owns the gfxinfo
  handle for the process lifetime. gfxinfo's active_gpu() runs a full NVML
  init/teardown (~20ms, blocking) and we were paying it on the async
  runtime for every collect — measured at ~80% of the agent's entire
  active CPU on a GPU machine. The handle holds Rc<Nvml> (not Send), so a
  thread + mpsc/oneshot channel pair confines it; a zero-total-VRAM reply
  is treated as a dead session (driver reload) and re-probed.
- journalctl now runs via tokio::process instead of blocking one of the
  two runtime workers for the duration of the subprocess.
- TtlCell (state.rs) replaces the four hand-rolled static TTL caches; a
  cached negative result now counts as fresh, so hosts with no matching
  temp sensor or GPU stop rescanning every request. Single lock+clone on
  the GPU cache hit path (was two).

Correctness:
- Process/child CPU times are now microseconds as documented; they were
  milliseconds, rendering 1000x too small next to (correct) thread times.
- Non-Linux per-process CPU%% clamps AFTER dividing by core count; a
  4-cores-busy process on an 8-core box reported 12.5% instead of 50%.
- Journal timestamps are real RFC 3339 UTC plus an additive timestamp_us
  field (sorting is now numeric); the old strings were Debug-formatted
  SystemTime mangled by string replace.
- Partition detection uses /sys/block on Linux: whole-disk filesystems on
  names like nvme0n1 or zram1 are no longer misclassified as partitions.
  One shared parent_disk_name() replaces two inline copies.
- New sampled_at_ms on the metrics payload (additive) records when the
  snapshot was actually collected, so clients can compute exact rates
  across the agent's TTL cache.

Security/robustness:
- key.pem is created 0600 (was umask default 0644, world-readable) and
  pre-1.51 keys are tightened on startup.
- Per-PID detail/journal caches now evict (60s max age, 64 entries max);
  they previously grew without bound under PID-walking clients.
- The two per-PID ws handlers collapse into one generic helper.
- /proc/<pid>/stat parsing unified in one comm-safe module.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(tui): responsive input, request timeouts, poisoned-stream reconnect

R1 — input latency: the event loop drained input once per iteration, then
slept the whole metrics interval; keys and wheel events queued for up to
500ms (or the full interval at slower rates) and applied in bursts. The
input block is extracted to drain_input() and the tail sleep replaced by
a deadline wait in <=33ms poll slices that handles and repaints input the
moment it arrives. Verified: help modal opens <150ms into a 2000ms tick.

R2 — freeze-proofing: metrics/processes/disks requests had no timeout; a
half-dead connection left ws.next() pending forever and froze the TUI
with no way to quit (raw mode eats Ctrl+C as an unread key event). All
requests now carry a 5s budget.

C3 — desync: replies are matched to requests by order alone, so a timed-
out request's late reply would shift every subsequent reply off by one.
Any timeout now treats the stream as poisoned and goes through the
reconnect flow — a fresh stream is aligned by construction. The modal
endpoints additionally mark process details unsupported (flag resets on
modal close/selection change) so a detail-less agent doesn't cause a
reconnect loop. While disconnected the fetch path idles: recovery belongs
to the manual/auto retry paths instead of 5s-timeout hammering.

C7 — fit::truncate_middle_cols replaces util::truncate_middle: display-
width aware and char-boundary safe; the byte-slicing version panicked the
draw loop on non-ASCII device names.

Verified live: agent kill -9 mid-session -> error modal in <3s, q exits
while disconnected, r reconnects and resumes.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* chore: version 1.51.0, path-dep the wasm examples, README notes

- socktop, socktop_agent, socktop_connector -> 1.51.0.
- socktop_wasm_test and zellij_socktop_plugin consume the in-repo
  connector via path deps so wasm-feature API drift is caught at PR time
  instead of after publish. Immediately proved out: the wasm requests
  module needed the new sampled_at_ms field, invisible to native builds.
- zellij plugin gains the standalone [workspace] marker (it could not be
  cargo-checked in-tree at all before). NOTE: its lib.rs has pre-existing
  compile errors unrelated to the connector (static mut STATE conflicts
  with register_plugin!, missing BTreeMap import) — needs its own rework,
  out of scope here.
- README: sampled_at_ms in the example payload.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(review): restore Agent Update Required flow, command field, axis alignment

Fixes from Jason's hands-on verification of the branch:

1. Old-agent messaging regression (this branch): a detail-request timeout
   went through the loud poison/reconnect flow, burying the ProcessDetails
   modal's 'Agent Update Required' message under a connection-error modal.
   Old agents IGNORE unknown messages (no late reply, no desync), so the
   optional per-PID endpoints now use quiet_reconnect(): swap the stream
   silently (still safe against merely-slow agents) and let the modal show
   its message. Only a failed reconnect surfaces loudly. Verified against
   a real v1.40.0 agent: message shows, session stays healthy.

2. Draw starvation (this branch): an agent that never answers get_metrics
   put the loop in fetch->timeout->poison->restart cycles that never
   reached the draw call — permanently blank TUI. The iteration now paints
   before fetching, and a second consecutive metrics timeout trips a
   circuit breaker: persistent 'Agent is not responding' error, recovery
   left to the manual/30s retry paths. Verified against a 0.9-era agent.

3. Command & Details pane blank (pre-existing on master): the minimal-
   refresh optimization dropped cmd/exe/cwd from the detail endpoint's
   ProcessRefreshKind, so process.cmd() had nothing to return. Restored
   with UpdateKind::OnlyIfNotSet — immutable values, read once per PID.
   Regression test added; journal E2E re-verified (100 entries render).

4. Scatter-plot axis misalignment: Y labels used a fixed 4-char field from
   the era when CPU times were 1000x too small; honest millisecond values
   (e.g. 136114) blew through it. Labels now right-align to the widest
   value per frame and X labels/titles share the dynamic padding.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* feat: journal access notice, 1.60.0, install script, changelog, riscv protoc fallback

- Journal pane now distinguishes 'no entries' from 'no journal access':
  journalctl exits 0 with empty output when the agent's user simply can't
  see the target's entries (demo mode / user-run agents), explaining
  itself only on stderr. The agent forwards that hint as an additive
  JournalResponse.notice and the client renders it with practical advice.
  Verified E2E via a stub journalctl emulating the unprivileged case.
- Version 1.60.0 across all crates (1.51 would read fine, but the repo's
  scheme is 1.40/1.50/…, and a literal 1.6.0 would sort BELOW 1.50.0 in
  semver). All user-facing version strings already come from
  CARGO_PKG_VERSION — a stale binary was the only way to see an old one.
- scripts/install.sh: build-from-source install/upgrade for the test
  fleet (Linux + macOS). Detects in-repo checkouts, installs rustup when
  missing, replaces a systemd socktop-agent service binary in place and
  restarts it, requires system protoc on riscv64.
- build.rs (agent + connector): fall back to $PROTOC / PATH when
  protoc-bin-vendored has no binary for the host (riscv64) — native SBC
  builds previously panicked in the build script.
- CHANGELOG.md covering v1.50.0 -> 1.60.0.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* test: add notice field to cache test initializer

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix: untrack zellij plugin build dir; installer updates all PATH copies

- Remove zellij_socktop_plugin/target from git (3,577 files committed by
  accident in bf6ac87): the root .gitignore anchors /target to the repo
  root, so the standalone plugin's own build dir wasn't covered. Ignore
  target/ at any depth (also fixes the pre-existing
  '/socktop-wasm-test/target' entry, which pointed at a hyphenated path
  that doesn't exist).

- install.sh now updates EVERY copy of socktop/socktop_agent on PATH,
  not just $PREFIX: a stale 'cargo install' in ~/.cargo/bin shadows
  /usr/local/bin on most PATHs, so an install could 'succeed' while
  'socktop --version' kept reporting the old release. Extra copies that
  can't be written are warned about, not fatal, and the script now
  prints which binary is actually active on PATH at the end.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* feat(installer): manage the socktop-agent systemd service

Upgrade path (unit already present): NEVER touch the unit file — it is
the operator's config (SSL, tokens, ports live there as Environment=
lines). Only the binary at the unit's own ExecStart path is replaced,
then the service restarts. Flags/args preserved by construction.

Fresh path (no unit): full first-time setup mirroring the deb postinst
and the agent-service docs — create the socktop system user/group and
/var/lib/socktop, install docs/socktop-agent.service (ExecStart rewritten
to wherever this run installed the agent; embedded fallback for old
refs), daemon-reload, enable --now, and print how to turn on TLS/token.

Also: system-level operations get their own sudo decision (SYS_SUDO) —
previously they inherited the PREFIX sudo flag, so a writable --prefix
made the service section run groupadd/systemctl unprivileged and die.
No sudo at all now skips service management with a warning instead of
failing the install.

Both branches dry-run verified with stubbed systemctl/sudo.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(installer): don't bind fresh agent services onto occupied ports

The Orange Pi install put the new service straight into a crash-restart
loop: the unit's default --port 3000 collided with a Docker service
already publishing 3000 (Umami; Gitea and friends default there too).
Fresh installs now scan 3000/3001/3010/3231/3232 via ss and configure
the unit on the first free port, warning loudly when 3000 was taken and
printing the resulting ws:// URL. Upgrades still never touch the unit.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(agent): detect NVIDIA GPUs on distros without the unversioned NVML soname

On Debian and derivatives the NVIDIA driver ships only libnvidia-ml.so.1
(the unversioned symlink belongs to the dev package), and nvml-wrapper's
default init dlopens the unversioned name — so gfxinfo reported 'No GPU
found' on a fully functional RTX A2000 host while nvidia-smi worked
fine. Arch-family distros ship the symlink, which is why the desktop
never showed this.

The GPU worker now falls back to initializing NVML directly with the
versioned soname when gfxinfo's probe fails, collecting name/util/vram
through the same handle-caching path. nvml-wrapper was already in the
tree via gfxinfo — same version, no new build cost.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(agent): box the NVML handle variant (clippy large_enum_variant)

CI clippy runs with -D warnings; Nvml is a large struct next to the
16-byte Box<dyn Gpu> variant.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(installer): survive self-modification mid-run; sturdier unit detection

Root cause of the mixed-up second install on the A2000 host: when run
from the clone it manages, the script's own git checkout/merge REPLACES
scripts/install.sh while bash is still executing it. Bash reads scripts
lazily by byte offset, so it resumed parsing the NEW file at the OLD
offset and executed an arbitrary tail of it — observed as the fresh-
service path running on a host whose unit already existed: the port scan
saw the still-running old service on 3000 and silently wrote a new unit
on 3001, while enable --now on the already-active service changed
nothing until a manual daemon-reload.

Fix: the whole script now runs inside main(), invoked as
'main "$@"; exit $?' so bash parses everything up front and never
reads the file again after main returns (the exit lives in the same
parse unit — demonstrated necessary: with a bare 'main "$@"' ending,
bash still executed the swapped file's trailing content after main
returned).

Also: unit existence is now checked with 'systemctl cat' instead of
grepping the full list-unit-files output.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* docs(installer): use a durable ref in the usage example

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-21 12:45:10 -07:00
committed by GitHub
co-authored by Claude Fable 5
parent ebda3c51af
commit 0322308896
38 changed files with 6243 additions and 4062 deletions
+10 -7
View File
@@ -1,6 +1,6 @@
[package]
name = "socktop_agent"
version = "1.50.2"
version = "1.60.0"
authors = ["Jason Witty <jasonpwitty+socktop@proton.me>"]
description = "Socktop agent daemon. Serves host metrics over WebSocket."
edition = "2024"
@@ -10,11 +10,10 @@ homepage = "https://github.com/jasonwitty/socktop"
repository = "https://github.com/jasonwitty/socktop"
[dependencies]
# Tokio: Use minimal features instead of "full" to reduce binary size
# Only include: rt-multi-thread (async runtime), net (WebSocket), sync (Mutex/RwLock), macros (#[tokio::test])
# Excluded: io, fs, process, signal, time (not needed for this workload)
# Savings: ~200-300KB binary size, faster compile times
tokio = { version = "1", features = ["rt-multi-thread", "net", "sync", "macros"] }
# Tokio: minimal features instead of "full" to reduce binary size.
# rt-multi-thread (runtime), net (WebSocket), sync (Mutex/oneshot),
# macros (#[tokio::test]), process (async journalctl).
tokio = { version = "1", features = ["rt-multi-thread", "net", "sync", "macros", "process"] }
axum = { version = "0.7", features = ["ws", "macros"] }
sysinfo = { version = "0.37", features = ["network", "disk", "component"] }
serde = { version = "1", features = ["derive"] }
@@ -24,6 +23,10 @@ futures-util = "0.3.31"
tracing = { version = "0.1", optional = true }
tracing-subscriber = { version = "0.3", features = ["env-filter"], optional = true }
gfxinfo = { version = "0.1.2", optional = true }
# Direct NVML fallback for distros that ship only libnvidia-ml.so.1 (Debian
# and derivatives) — gfxinfo's default init dlopens the unversioned name.
# Same version gfxinfo already pulls in, so this adds no new build cost.
nvml-wrapper = { version = "0.10", optional = true }
once_cell = "1.19"
axum-server = { version = "0.7", features = ["tls-rustls"] }
rustls = { version = "0.23", features = ["aws-lc-rs"] }
@@ -36,7 +39,7 @@ time = { version = "0.3", default-features = false, features = ["formatting", "m
[features]
default = ["gpu"]
gpu = ["gfxinfo"]
gpu = ["gfxinfo", "nvml-wrapper"]
logging = ["tracing", "tracing-subscriber"]
[build-dependencies]
+6 -4
View File
@@ -1,13 +1,15 @@
fn main() {
// Vendored protoc for reproducible builds
let protoc = protoc_bin_vendored::protoc_bin_path().expect("protoc");
println!("cargo:rerun-if-changed=proto/processes.proto");
// Compile protobuf definitions for processes
let mut cfg = prost_build::Config::new();
cfg.out_dir(std::env::var("OUT_DIR").unwrap());
cfg.protoc_executable(protoc); // Use the vendored protoc directly
// Vendored protoc for reproducible builds where available. It ships no
// riscv64 binary, so on such hosts fall through to $PROTOC / PATH
// (prost-build's default lookup) — apt: protobuf-compiler.
if let Ok(protoc) = protoc_bin_vendored::protoc_bin_path() {
cfg.protoc_executable(protoc);
}
// Use local path (ensures file is inside published crate tarball)
cfg.compile_protos(&["proto/processes.proto"], &["proto"]) // relative to CARGO_MANIFEST_DIR
.expect("compile protos");
+110 -17
View File
@@ -1,6 +1,4 @@
// gpu.rs
#[cfg(feature = "gpu")]
use gfxinfo::active_gpu;
#[derive(Debug, Clone, serde::Serialize)]
pub struct GpuMetrics {
@@ -10,23 +8,118 @@ pub struct GpuMetrics {
pub mem_total_bytes: u64,
}
/// Collect metrics for the active GPU. `None` when there is no usable GPU.
///
/// Runs on a dedicated worker thread (see `worker`): gfxinfo's handle holds
/// an `Rc<Nvml>` (not `Send`), and *creating* it runs a full NVML library
/// init — ~20ms of blocking work that used to execute on the async runtime
/// for every collection. The worker owns one handle for the process lifetime,
/// so steady-state collection is just NVML queries. Measured on an RTX 5080
/// box, re-initing per collect was ~80% of the agent's entire active CPU.
#[cfg(feature = "gpu")]
pub fn collect_all_gpus() -> Result<Vec<GpuMetrics>, Box<dyn std::error::Error>> {
let gpu = active_gpu()?; // Use ? to unwrap Result
let info = gpu.info();
let metrics = GpuMetrics {
name: gpu.model().to_string(),
utilization_gpu_pct: info.load_pct() as u32,
mem_used_bytes: info.used_vram(),
mem_total_bytes: info.total_vram(),
};
Ok(vec![metrics])
pub async fn collect_all_gpus() -> Option<Vec<GpuMetrics>> {
worker::collect().await
}
#[cfg(not(feature = "gpu"))]
pub fn collect_all_gpus() -> Result<Vec<GpuMetrics>, Box<dyn std::error::Error>> {
// GPU support not available on this platform
Ok(vec![])
pub async fn collect_all_gpus() -> Option<Vec<GpuMetrics>> {
None
}
#[cfg(feature = "gpu")]
mod worker {
use super::GpuMetrics;
use once_cell::sync::OnceCell;
use std::sync::mpsc;
type Reply = tokio::sync::oneshot::Sender<Option<Vec<GpuMetrics>>>;
static TX: OnceCell<mpsc::Sender<Reply>> = OnceCell::new();
pub async fn collect() -> Option<Vec<GpuMetrics>> {
let tx = TX.get_or_init(spawn);
let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
tx.send(reply_tx).ok()?;
reply_rx.await.ok().flatten()
}
fn spawn() -> mpsc::Sender<Reply> {
let (tx, rx) = mpsc::channel::<Reply>();
std::thread::Builder::new()
.name("socktop-gpu".into())
.spawn(move || run(rx))
.expect("spawn gpu worker thread");
tx
}
enum Handle {
/// gfxinfo's own detection (AMD sysfs, NVIDIA via unversioned NVML).
Gfx(Box<dyn gfxinfo::Gpu>),
/// Direct NVML with an explicit versioned soname. Debian & friends
/// ship only libnvidia-ml.so.1 (the unversioned symlink lives in the
/// dev package), so gfxinfo's default dlopen fails there even though
/// the driver is fully functional.
Nvml(Box<nvml_wrapper::Nvml>),
}
fn probe() -> Option<Handle> {
if let Ok(g) = gfxinfo::active_gpu() {
return Some(Handle::Gfx(g));
}
nvml_wrapper::Nvml::builder()
.lib_path(std::ffi::OsStr::new("libnvidia-ml.so.1"))
.init()
.ok()
.map(|nvml| Handle::Nvml(Box::new(nvml)))
}
fn collect_from(handle: &Handle) -> Option<Vec<GpuMetrics>> {
match handle {
Handle::Gfx(gpu) => {
let info = gpu.info();
Some(vec![GpuMetrics {
name: gpu.model().to_string(),
utilization_gpu_pct: info.load_pct().clamp(0, 100),
mem_used_bytes: info.used_vram(),
mem_total_bytes: info.total_vram(),
}])
}
Handle::Nvml(nvml) => {
let device = nvml.device_by_index(0).ok()?;
let mem = device.memory_info().ok()?;
Some(vec![GpuMetrics {
name: device.name().unwrap_or_else(|_| "NVIDIA GPU".into()),
utilization_gpu_pct: device
.utilization_rates()
.map(|u| u.gpu.clamp(0, 100))
.unwrap_or(0),
mem_used_bytes: mem.used,
mem_total_bytes: mem.total,
}])
}
}
}
fn run(rx: mpsc::Receiver<Reply>) {
let mut handle: Option<Handle> = None;
// Probing failed: remember and answer None without re-initing the GPU
// stack per request. The agent's negative cache stops asking anyway.
let mut probe_failed = false;
while let Ok(reply) = rx.recv() {
if handle.is_none() && !probe_failed {
handle = probe();
probe_failed = handle.is_none();
}
let out = handle.as_ref().and_then(collect_from);
// A live GPU cannot report 0 total VRAM; zeros mean the session
// died (e.g. driver reload). Drop the handle so the next request
// re-probes.
if let Some(v) = &out
&& !v.is_empty()
&& v.iter().all(|g| g.mem_total_bytes == 0)
{
handle = None;
}
let _ = reply.send(out.filter(|v| !v.is_empty()));
}
}
}
+270 -270
View File
@@ -13,49 +13,60 @@ use std::collections::HashMap;
use std::fs;
#[cfg(target_os = "linux")]
use std::io;
use std::process::Command;
use std::sync::Mutex;
use std::time::Duration as StdDuration;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use sysinfo::{ProcessRefreshKind, ProcessesToUpdate};
#[cfg(feature = "logging")]
use tracing::warn;
// NOTE: CPU normalization env removed; non-Linux now always reports per-process share (0..100) as given by sysinfo.
// Read (utime, stime) in milliseconds from /proc/{pid}/stat in one go.
// Returns (0, 0) if the file can't be read.
//
// We use `rfind(')')` to step past the `comm` field, which can contain
// arbitrary characters (including spaces and parens), then index the
// post-comm fields by position. This is the same trick `read_proc_jiffies`
// uses below — `split_whitespace().collect::<Vec<_>>()` from the start of
// the file would mis-parse process names with spaces, and also wastes an
// allocation per call. Two callers used to read this file twice (once for
// user, once for system); now it's one syscall per detailed-process record.
/// Shared parsing for `/proc/<pid>/stat` (and per-thread `task/<tid>/stat`).
///
/// The second field, `comm`, can contain arbitrary bytes including spaces and
/// parentheses, so naive whitespace splitting mis-parses such names. All
/// callers step past the LAST `')'` and index the remaining space-separated
/// fields from there: 0 = state, 1 = ppid, 11 = utime, 12 = stime,
/// 19 = starttime.
#[cfg(target_os = "linux")]
fn get_cpu_times_ms(pid: u32) -> (u64, u64) {
mod procstat {
/// Everything after `") "` — the post-comm fields.
pub fn after_comm(stat: &str) -> Option<&str> {
stat.get(stat.rfind(')')? + 2..)
}
pub fn field(stat: &str, n: usize) -> Option<&str> {
after_comm(stat)?.split_whitespace().nth(n)
}
/// (utime, stime) in clock ticks.
pub fn utime_stime(stat: &str) -> Option<(u64, u64)> {
let mut it = after_comm(stat)?.split_whitespace();
let utime = it.nth(11)?.parse().ok()?;
let stime = it.next()?.parse().ok()?;
Some((utime, stime))
}
/// One clock tick at USER_HZ=100 (universal on Linux) in microseconds.
pub const TICK_US: u64 = 10_000;
}
// Read (utime, stime) in MICROSECONDS from /proc/{pid}/stat in one syscall.
// Returns (0, 0) if the file can't be read. Units match the wire contract
// (`DetailedProcessInfo.cpu_time_user` is documented as µs) and the thread
// records — this used to return ms, making process/child CPU times render
// 1000x too small next to thread times.
#[cfg(target_os = "linux")]
fn get_cpu_times_us(pid: u32) -> (u64, u64) {
let Ok(s) = fs::read_to_string(format!("/proc/{pid}/stat")) else {
return (0, 0);
};
let Some(rpar) = s.rfind(')') else {
let Some((utime, stime)) = procstat::utime_stime(&s) else {
return (0, 0);
};
let Some(after) = s.get(rpar + 2..) else {
return (0, 0);
};
let mut it = after.split_whitespace();
// Post-comm field offsets: state, ppid, pgrp, session, tty_nr, tpgid,
// flags, minflt, cminflt, majflt, cmajflt, utime, stime, ...
// utime is offset 11; stime follows.
let utime = it.nth(11).and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
let stime = it.next().and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
// 1 tick = 10ms at 100 Hz (USER_HZ).
(utime * 10, stime * 10)
(utime * procstat::TICK_US, stime * procstat::TICK_US)
}
#[cfg(not(target_os = "linux"))]
fn get_cpu_times_ms(_pid: u32) -> (u64, u64) {
fn get_cpu_times_us(_pid: u32) -> (u64, u64) {
(0, 0)
}
// Runtime toggles (read once)
@@ -117,46 +128,32 @@ fn name_cache_cleanup_threshold() -> usize {
})
}
// Tiny TTL caches to avoid rescanning sensors every 500ms
// Tiny TTL caches to avoid rescanning sensors every 500ms.
//
// The cached type is Option<...>: a fresh `None` means "we looked recently
// and found nothing" — machines with no matching sensor/GPU no longer rescan
// on every request, only once per TTL.
const TTL: Duration = Duration::from_millis(1500);
struct TempCache {
at: Option<Instant>,
v: Option<f32>,
}
static TEMP: OnceCell<Mutex<TempCache>> = OnceCell::new();
static TEMP: crate::state::TtlCell<Option<f32>> = crate::state::TtlCell::new();
static GPUS: crate::state::TtlCell<Option<Vec<crate::gpu::GpuMetrics>>> =
crate::state::TtlCell::new();
// Last time `state.components` was refreshed (by any caller). Both
// Gate on `state.components` refreshes (hwmon scans). Both
// collect_fast_metrics and collect_disks need fresh sensor values; without
// this gate they were each doing their own `Components::refresh` on their
// own cadence, paying the hwmon syscall cost twice per polling cycle.
// 1s is short enough that disk temps stay accurate (they change slowly) and
// long enough to suppress back-to-back refreshes from concurrent endpoints.
// this they each paid the hwmon syscall cost on their own cadence. 1s keeps
// disk temps accurate (they change slowly) while suppressing back-to-back
// refreshes from concurrent endpoints.
const COMPONENTS_REFRESH_TTL: Duration = Duration::from_millis(1000);
static COMPONENTS_LAST_REFRESH: OnceCell<Mutex<Option<Instant>>> = OnceCell::new();
static COMPONENTS_STAMP: crate::state::TtlCell<()> = crate::state::TtlCell::new();
/// Refresh `state.components` only if the cached refresh timestamp is older
/// than `COMPONENTS_REFRESH_TTL`. Caller must already hold the components
/// lock.
/// Refresh `state.components` at most once per `COMPONENTS_REFRESH_TTL`.
/// Caller must already hold the components lock.
fn refresh_components_if_stale(components: &mut sysinfo::Components) {
let lock = COMPONENTS_LAST_REFRESH.get_or_init(|| Mutex::new(None));
let mut last = match lock.lock() {
Ok(g) => g,
Err(_) => return, // Poisoned — skip; values stay as-is until next call
};
let now = Instant::now();
let stale = last.is_none_or(|t| now.duration_since(t) >= COMPONENTS_REFRESH_TTL);
if stale {
if COMPONENTS_STAMP.claim_stale(COMPONENTS_REFRESH_TTL) {
components.refresh(false);
*last = Some(now);
}
}
struct GpuCache {
at: Option<Instant>,
v: Option<Vec<crate::gpu::GpuMetrics>>,
}
static GPUC: OnceCell<Mutex<GpuCache>> = OnceCell::new();
// Static caches for unchanging data
static HOSTNAME: OnceCell<String> = OnceCell::new();
struct NetworkNameCache {
@@ -166,54 +163,6 @@ struct NetworkNameCache {
static NETWORK_CACHE: OnceCell<Mutex<NetworkNameCache>> = OnceCell::new();
static CPU_VEC: OnceCell<Mutex<Vec<f32>>> = OnceCell::new();
fn cached_temp() -> Option<f32> {
if !temp_enabled() {
return None;
}
let now = Instant::now();
let lock = TEMP.get_or_init(|| Mutex::new(TempCache { at: None, v: None }));
let mut c = lock.lock().ok()?;
if c.at.is_none_or(|t| now.duration_since(t) >= TTL) {
c.at = Some(now);
// caller will fill this; we just hold a slot
c.v = None;
}
c.v
}
fn set_temp(v: Option<f32>) {
if let Some(lock) = TEMP.get()
&& let Ok(mut c) = lock.lock()
{
c.v = v;
c.at = Some(Instant::now());
}
}
fn cached_gpus() -> Option<Vec<crate::gpu::GpuMetrics>> {
if !gpu_enabled() {
return None;
}
let now = Instant::now();
let lock = GPUC.get_or_init(|| Mutex::new(GpuCache { at: None, v: None }));
let mut c = lock.lock().ok()?;
if c.at.is_none_or(|t| now.duration_since(t) >= TTL) {
// mark stale; caller will refresh
c.at = Some(now);
c.v = None;
}
c.v.clone()
}
fn set_gpus(v: Option<Vec<crate::gpu::GpuMetrics>>) {
if let Some(lock) = GPUC.get()
&& let Ok(mut c) = lock.lock()
{
c.v = v.clone();
c.at = Some(Instant::now());
}
}
// Collect only fast-changing metrics (CPU/mem/net + optional temps/gpus).
pub async fn collect_fast_metrics(state: &AppState) -> Metrics {
let ttl = StdDuration::from_millis(metrics_ttl_ms());
@@ -253,10 +202,13 @@ pub async fn collect_fast_metrics(state: &AppState) -> Metrics {
let swap_used = sys.used_swap();
drop(sys);
// CPU temperature: only refresh sensors if cache is stale
let cpu_temp_c = if cached_temp().is_some() {
cached_temp()
} else if temp_enabled() {
// CPU temperature: only rescan sensors when the cached result (even a
// cached "no sensor found") goes stale.
let cpu_temp_c = if !temp_enabled() {
None
} else if let Some(cached) = TEMP.get_fresh(TTL) {
cached
} else {
let val = {
let mut components = state.components.lock().await;
refresh_components_if_stale(&mut components);
@@ -273,10 +225,8 @@ pub async fn collect_fast_metrics(state: &AppState) -> Metrics {
}
})
};
set_temp(val);
TEMP.set(val);
val
} else {
None
};
// Networks with reusable name cache
@@ -320,47 +270,37 @@ pub async fn collect_fast_metrics(state: &AppState) -> Metrics {
cache.infos.clone()
};
// GPUs: if we already determined none exist, short-circuit (no repeated probing)
let gpus = if gpu_enabled() {
if state.gpu_checked.load(std::sync::atomic::Ordering::Acquire)
&& !state.gpu_present.load(std::sync::atomic::Ordering::Relaxed)
{
None
} else if cached_gpus().is_some() {
cached_gpus()
} else {
let v = match collect_all_gpus() {
Ok(v) if !v.is_empty() => Some(v),
Ok(_) => None,
Err(_e) => {
#[cfg(feature = "logging")]
warn!("gpu collection failed: {_e}");
None
}
};
// First probe records presence; subsequent calls rely on cache flags.
if !state
.gpu_checked
.swap(true, std::sync::atomic::Ordering::AcqRel)
{
if v.is_some() {
state
.gpu_present
.store(true, std::sync::atomic::Ordering::Release);
} else {
state
.gpu_present
.store(false, std::sync::atomic::Ordering::Release);
}
}
set_gpus(v.clone());
v
}
} else {
// GPUs: negative-probe cache short-circuits GPU-less hosts; otherwise the
// TTL cache answers, and only a stale miss reaches the worker thread.
let gpus = if !gpu_enabled()
|| (state.gpu_checked.load(std::sync::atomic::Ordering::Acquire)
&& !state.gpu_present.load(std::sync::atomic::Ordering::Relaxed))
{
None
} else if let Some(cached) = GPUS.get_fresh(TTL) {
cached
} else {
let v = collect_all_gpus().await;
// First probe records presence; subsequent calls rely on the flags.
if !state
.gpu_checked
.swap(true, std::sync::atomic::Ordering::AcqRel)
{
state
.gpu_present
.store(v.is_some(), std::sync::atomic::Ordering::Release);
}
GPUS.set(v.clone());
v
};
let sampled_at_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let metrics = Metrics {
sampled_at_ms,
cpu_total,
cpu_per_core,
mem_total,
@@ -381,6 +321,48 @@ pub async fn collect_fast_metrics(state: &AppState) -> Metrics {
metrics
}
/// Best-effort parent-disk name for a partition device name:
/// "nvme0n1p1" -> "nvme0n1", "mmcblk0p2" -> "mmcblk0", "sda1" -> "sda".
/// Works with or without a "/dev/" prefix.
fn parent_disk_name(name: &str) -> &str {
if let Some(pos) = name.rfind('p') {
let suffix = &name[pos + 1..];
if !suffix.is_empty() && suffix.chars().all(|c| c.is_ascii_digit()) {
return &name[..pos];
}
}
name.trim_end_matches(|c: char| c.is_ascii_digit())
}
/// Whether a device name refers to a partition rather than a whole disk.
///
/// On Linux, whole-disk devices are directories under /sys/block and
/// partitions are not, so "is a partition" = "not in /sys/block, but the
/// derived parent is". This gets right the cases the old name heuristic got
/// wrong: a whole-disk filesystem on nvme0n1 (ends in a digit but IS in
/// /sys/block) and zram1 (a whole device). Non-Linux keeps the heuristic.
fn is_partition_name(name: &str) -> bool {
let bare = name.strip_prefix("/dev/").unwrap_or(name);
#[cfg(target_os = "linux")]
{
let sys_block = std::path::Path::new("/sys/block");
if sys_block.is_dir() {
return !sys_block.join(bare).is_dir()
&& sys_block.join(parent_disk_name(bare)).is_dir();
}
}
is_partition_heuristic(bare)
}
/// Name-based fallback for platforms without /sys/block: a p<digits> marker
/// or a trailing non-zero digit.
fn is_partition_heuristic(bare: &str) -> bool {
bare.contains("p1")
|| bare.contains("p2")
|| bare.contains("p3")
|| bare.ends_with(|c: char| c.is_ascii_digit() && c != '0')
}
// Cached disks
pub async fn collect_disks(state: &AppState) -> Vec<DiskInfo> {
let ttl = StdDuration::from_millis(disks_ttl_ms());
@@ -445,19 +427,7 @@ pub async fn collect_disks(state: &AppState) -> Vec<DiskInfo> {
return None;
}
// Determine if this is a partition
let is_partition = name.contains("p1")
|| name.contains("p2")
|| name.contains("p3")
|| name.ends_with('1')
|| name.ends_with('2')
|| name.ends_with('3')
|| name.ends_with('4')
|| name.ends_with('5')
|| name.ends_with('6')
|| name.ends_with('7')
|| name.ends_with('8')
|| name.ends_with('9');
let is_partition = is_partition_name(&name);
// Try to find temperature for this disk
let temperature = disk_temps.iter().find_map(|(key, &temp)| {
@@ -491,25 +461,7 @@ pub async fn collect_disks(state: &AppState) -> Vec<DiskInfo> {
for partition in &partitions {
if partition.is_partition {
// Extract parent disk name
// nvme0n1p1 -> nvme0n1, sda1 -> sda, mmcblk0p1 -> mmcblk0
let parent_name = if let Some(pos) = partition.name.rfind('p') {
// Check if character after 'p' is a digit
if partition
.name
.chars()
.nth(pos + 1)
.is_some_and(|c| c.is_ascii_digit())
{
&partition.name[..pos]
} else {
// Handle sda1, sdb2, etc (just trim trailing digit)
partition.name.trim_end_matches(char::is_numeric)
}
} else {
// Handle sda1, sdb2, etc (just trim trailing digit)
partition.name.trim_end_matches(char::is_numeric)
};
let parent_name = parent_disk_name(&partition.name);
// Look up temperature for the PARENT disk, not the partition
// Strip /dev/ prefix if present for matching
@@ -553,21 +505,7 @@ pub async fn collect_disks(state: &AppState) -> Vec<DiskInfo> {
// Add partitions after their parent disk
for partition in partitions {
if partition.is_partition {
// Find parent disk index
let parent_name = if let Some(pos) = partition.name.rfind('p') {
if partition
.name
.chars()
.nth(pos + 1)
.is_some_and(|c| c.is_ascii_digit())
{
&partition.name[..pos]
} else {
partition.name.trim_end_matches(char::is_numeric)
}
} else {
partition.name.trim_end_matches(char::is_numeric)
};
let parent_name = parent_disk_name(&partition.name);
// Find where to insert this partition (after its parent)
if let Some(parent_idx) = disks.iter().position(|d| d.name == parent_name) {
@@ -619,15 +557,8 @@ fn read_total_jiffies() -> io::Result<u64> {
#[cfg(target_os = "linux")]
#[inline]
fn read_proc_jiffies(pid: u32) -> Option<u64> {
let path = format!("/proc/{pid}/stat");
let s = fs::read_to_string(path).ok()?;
// Find the right parenthesis that terminates comm; everything after is space-separated fields starting at "state"
let rpar = s.rfind(')')?;
let after = s.get(rpar + 2..)?; // skip ") "
let mut it = after.split_whitespace();
// utime (14th field) is offset 11 from "state", stime (15th) is next
let utime = it.nth(11)?.parse::<u64>().ok()?;
let stime = it.next()?.parse::<u64>().ok()?;
let s = fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
let (utime, stime) = procstat::utime_stime(&s)?;
Some(utime.saturating_add(stime))
}
@@ -818,9 +749,12 @@ pub async fn collect_processes_all(state: &AppState) -> ProcessesPayload {
};
// Convert to percentage of total CPU capacity
// e.g., 100% on 2 cores of 8 core system = 25% total CPU
let raw = p.cpu_usage(); // This is per-core percentage
let total_cpu = raw.clamp(0.0, 100.0) / cpu_count;
// e.g., 100% on 2 cores of 8 core system = 25% total CPU.
// sysinfo reports per-core percentage which EXCEEDS 100 for
// multi-threaded processes, so clamp AFTER dividing — clamping
// first truncated e.g. 400%-on-8-cores to 12.5% instead of 50%.
let raw = p.cpu_usage();
let total_cpu = (raw / cpu_count.max(1.0)).clamp(0.0, 100.0);
proc_cache.reusable_vec.push(ProcessInfo {
pid,
@@ -984,15 +918,8 @@ fn proc_state_label(c: char) -> &'static str {
#[cfg(target_os = "linux")]
fn read_parent_pid_from_proc(pid: u32) -> Option<u32> {
let stat = fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
// Format: pid (comm) state ppid ... — comm can contain spaces/parens,
// so we step past the closing paren first.
let ppid_start = stat.rfind(')')?;
// After ") ": state, ppid, ... — ppid is the second field.
stat[ppid_start + 1..]
.split_whitespace()
.nth(1)?
.parse::<u32>()
.ok()
// Post-comm field 1 is ppid.
procstat::field(&stat, 1)?.parse::<u32>().ok()
}
/// Collect process information from /proc files
@@ -1035,16 +962,9 @@ fn collect_process_info_from_proc(
let thread_count = st.threads;
let status = proc_state_label(st.state_ch).to_string();
// Read start time from stat — comm-safe via rfind(')').
// starttime is post-comm field 19.
let start_time = if let Ok(stat) = fs::read_to_string(format!("/proc/{pid}/stat")) {
let stat_end = stat.rfind(')')?;
// After ") ": state, ppid, ..., starttime — starttime is the 20th
// post-comm field (index 19).
stat[stat_end + 1..]
.split_whitespace()
.nth(19)?
.parse::<u64>()
.ok()?
procstat::field(&stat, 19)?.parse::<u64>().ok()?
} else {
0
};
@@ -1079,7 +999,7 @@ fn collect_process_info_from_proc(
.map(|p| p.to_string_lossy().to_string());
// One read of /proc/{pid}/stat covers both user + system CPU times.
let (cpu_time_user, cpu_time_system) = get_cpu_times_ms(pid);
let (cpu_time_user, cpu_time_system) = get_cpu_times_us(pid);
Some(DetailedProcessInfo {
pid,
@@ -1192,18 +1112,7 @@ fn collect_thread_info(pid: u32) -> Vec<crate::types::ThreadInfo> {
continue;
};
// Thread/comm names can contain spaces or parens, so step past the
// last ')' before parsing post-comm fields. Post-comm offsets:
// 0: state, 1: ppid, 2: pgrp, ..., 11: utime, 12: stime
let Some(rpar) = stat_content.rfind(')') else {
continue;
};
let Some(after) = stat_content.get(rpar + 1..) else {
continue;
};
let mut it = after.split_whitespace();
let status = it
.next()
let status = procstat::field(&stat_content, 0)
.and_then(|s| s.chars().next())
.map(|c| match c {
'R' => "Running",
@@ -1218,14 +1127,9 @@ fn collect_thread_info(pid: u32) -> Vec<crate::types::ThreadInfo> {
.unwrap_or("Unknown")
.to_string();
// 10 fields between state and utime (ppid..cmajflt).
let utime = it.nth(10).and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
let stime = it.next().and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
// Convert clock ticks to microseconds (assuming 100 Hz)
// 1 tick = 10ms = 10,000 microseconds
let cpu_time_user = utime * 10_000;
let cpu_time_system = stime * 10_000;
let (utime, stime) = procstat::utime_stime(&stat_content).unwrap_or((0, 0));
let cpu_time_user = utime * procstat::TICK_US;
let cpu_time_system = stime * procstat::TICK_US;
threads.push(crate::types::ThreadInfo {
tid,
@@ -1257,10 +1161,17 @@ pub async fn collect_process_metrics(
system.refresh_processes_specifics(
ProcessesToUpdate::Some(&[sysinfo::Pid::from_u32(pid)]),
false,
// cmd/exe/cwd feed the modal's Command & Details pane. They're
// immutable per process, so OnlyIfNotSet reads them once per PID and
// serves the cache afterwards — the "minimal refresh" optimization
// had dropped them entirely, leaving the pane blank.
ProcessRefreshKind::nothing()
.with_memory()
.with_cpu()
.with_disk_usage(),
.with_disk_usage()
.with_cmd(sysinfo::UpdateKind::OnlyIfNotSet)
.with_exe(sysinfo::UpdateKind::OnlyIfNotSet)
.with_cwd(sysinfo::UpdateKind::OnlyIfNotSet),
);
let process = system
@@ -1350,7 +1261,7 @@ pub async fn collect_process_metrics(
let threads = collect_thread_info(pid);
// One read of /proc/{pid}/stat covers both user + system CPU times.
let (cpu_time_user, cpu_time_system) = get_cpu_times_ms(pid);
let (cpu_time_user, cpu_time_system) = get_cpu_times_us(pid);
// Now construct the detailed info without holding the lock
let detailed_info = DetailedProcessInfo {
@@ -1384,9 +1295,26 @@ pub async fn collect_process_metrics(
})
}
/// Collect journal entries for a specific process
pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
let output = Command::new("journalctl")
/// Epoch microseconds -> RFC 3339 UTC for display. The old code
/// Debug-formatted a SystemTime and string-replaced it into a raw epoch
/// string that was neither ISO 8601 nor what the field documented.
fn format_journal_timestamp(timestamp_us: u64) -> String {
time::OffsetDateTime::from_unix_timestamp_nanos(timestamp_us as i128 * 1000)
.ok()
.and_then(|t| {
t.format(&time::format_description::well_known::Rfc3339)
.ok()
})
.unwrap_or_else(|| timestamp_us.to_string())
}
/// Collect journal entries for a specific process.
///
/// Async via tokio::process — journalctl can take hundreds of ms on slow
/// storage, and the old std::process call blocked one of the runtime's two
/// worker threads for the duration.
pub async fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
let output = tokio::process::Command::new("journalctl")
.args([
&format!("_PID={pid}"),
"--output=json",
@@ -1394,6 +1322,7 @@ pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
"--no-pager",
])
.output()
.await
.map_err(|e| format!("Failed to execute journalctl: {e}"))?;
if !output.status.success() {
@@ -1415,27 +1344,14 @@ pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
let json: serde_json::Value =
serde_json::from_str(line).map_err(|e| format!("Failed to parse journal JSON: {e}"))?;
// Extract relevant fields
let timestamp_str = json
// __REALTIME_TIMESTAMP is epoch microseconds as a string.
let timestamp_us = json
.get("__REALTIME_TIMESTAMP")
.and_then(|v| v.as_str())
.unwrap_or("0");
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(0);
// Convert timestamp to ISO 8601 format
let timestamp = if let Ok(ts_micros) = timestamp_str.parse::<u64>() {
let ts_secs = ts_micros / 1_000_000;
let ts_nanos = (ts_micros % 1_000_000) * 1000;
let time = SystemTime::UNIX_EPOCH
+ Duration::from_secs(ts_secs)
+ Duration::from_nanos(ts_nanos);
// Simple ISO 8601 format - we can improve this if needed
format!("{time:?}")
.replace("SystemTime { tv_sec: ", "")
.replace(", tv_nsec: ", ".")
.replace(" }", "")
} else {
timestamp_str.to_string()
};
let timestamp = format_journal_timestamp(timestamp_us);
let priority = match json.get("PRIORITY").and_then(|v| v.as_str()) {
Some("0") => LogLevel::Emergency,
@@ -1482,6 +1398,7 @@ pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
entries.push(JournalEntry {
timestamp,
timestamp_us,
priority,
message,
unit,
@@ -1493,7 +1410,26 @@ pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
}
// Sort by timestamp (newest first)
entries.sort_by(|a, b| b.timestamp.cmp(&a.timestamp));
entries.sort_by_key(|e| std::cmp::Reverse(e.timestamp_us));
// journalctl exits 0 with no output when the invoking user simply cannot
// SEE the process's entries (e.g. a user-run agent asking about a system
// service) — but it explains itself on stderr ("You are currently not
// seeing messages from other users and the system…"). Pass that hint
// along so the client can distinguish "no logs" from "no access".
let notice = if entries.is_empty() {
let err = String::from_utf8_lossy(&output.stderr);
let hint: String = err
.lines()
.map(str::trim)
.filter(|l| !l.is_empty())
.take(2)
.collect::<Vec<_>>()
.join(" ");
if hint.is_empty() { None } else { Some(hint) }
} else {
None
};
let response_timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
@@ -1507,6 +1443,70 @@ pub fn collect_journal_entries(pid: u32) -> Result<JournalResponse, String> {
entries,
total_count,
truncated,
notice,
cached_at: response_timestamp,
})
}
#[cfg(test)]
mod tests {
use super::*;
/// comm can contain spaces and parens; parsing must key off the LAST ')'.
#[cfg(target_os = "linux")]
#[test]
fn procstat_handles_hostile_comm_names() {
let stat = "1234 (weird name) (2)) R 1 2 3 4 5 6 7 8 9 10 700 800 0 0 20";
assert_eq!(procstat::field(stat, 0), Some("R"));
assert_eq!(procstat::field(stat, 1), Some("1"));
assert_eq!(procstat::utime_stime(stat), Some((700, 800)));
}
/// USER_HZ ticks convert to MICROSECONDS — the wire contract. This used
/// to be *10 (ms), rendering process CPU times 1000x too small next to
/// thread times.
#[cfg(target_os = "linux")]
#[test]
fn cpu_times_are_microseconds() {
assert_eq!(procstat::TICK_US, 10_000);
}
#[test]
fn parent_disk_name_strips_partition_suffixes() {
assert_eq!(parent_disk_name("nvme0n1p1"), "nvme0n1");
assert_eq!(parent_disk_name("nvme1n1p12"), "nvme1n1");
assert_eq!(parent_disk_name("mmcblk0p2"), "mmcblk0");
assert_eq!(parent_disk_name("sda1"), "sda");
assert_eq!(parent_disk_name("/dev/nvme0n1p1"), "/dev/nvme0n1");
// 'p' inside a word is not a partition marker.
assert_eq!(parent_disk_name("mapper/vg-lv"), "mapper/vg-lv");
}
/// The old heuristic flagged whole-disk names ending in a digit
/// (nvme0n1, zram1) as partitions. On Linux /sys/block decides; this
/// pins the real-machine behavior for devices every Linux box has.
#[cfg(target_os = "linux")]
#[test]
fn sys_block_devices_are_not_partitions() {
let sys_block = std::path::Path::new("/sys/block");
if !sys_block.is_dir() {
return; // exotic environment; nothing to assert
}
for entry in std::fs::read_dir(sys_block).unwrap().flatten() {
let name = entry.file_name().to_string_lossy().into_owned();
assert!(
!is_partition_name(&name),
"{name} is a whole disk but was flagged as a partition"
);
}
}
#[test]
fn journal_timestamps_are_rfc3339() {
let s = format_journal_timestamp(1_786_752_000_000_000);
assert_eq!(s, "2026-08-15T00:00:00Z");
// Sub-second precision survives.
let s = format_journal_timestamp(1_786_752_000_123_456);
assert!(s.starts_with("2026-08-15T00:00:00.123456"), "{s}");
}
}
+50 -1
View File
@@ -74,6 +74,55 @@ pub struct AppState {
pub cache_journal_entries: Arc<Mutex<HashMap<u32, CacheEntry<crate::types::JournalResponse>>>>,
}
/// TTL-gated value behind a std Mutex, for `static` caches on hot paths.
/// Replaces the hand-rolled TempCache/GpuCache/refresh-timestamp statics
/// that each reimplemented the same at/value pair.
pub struct TtlCell<T> {
inner: std::sync::Mutex<CacheEntry<T>>,
}
impl<T: Clone> Default for TtlCell<T> {
fn default() -> Self {
Self::new()
}
}
impl<T: Clone> TtlCell<T> {
pub const fn new() -> Self {
Self {
inner: std::sync::Mutex::new(CacheEntry::new()),
}
}
/// The stored value, only while fresh. Poisoned lock reads as a miss.
pub fn get_fresh(&self, ttl: Duration) -> Option<T> {
let g = self.inner.lock().ok()?;
if g.is_fresh(ttl) {
g.value.clone()
} else {
None
}
}
pub fn set(&self, v: T) {
if let Ok(mut g) = self.inner.lock() {
g.set(v);
}
}
/// True exactly once per TTL window: restamps and tells the caller to do
/// the refresh. Atomic check-and-stamp so concurrent callers don't both
/// refresh.
pub fn claim_stale(&self, ttl: Duration) -> bool {
let Ok(mut g) = self.inner.lock() else {
return false;
};
if g.at.is_none_or(|t| t.elapsed() >= ttl) {
g.at = Some(Instant::now());
true
} else {
false
}
}
}
#[derive(Clone, Debug)]
pub struct CacheEntry<T> {
pub at: Option<Instant>,
@@ -87,7 +136,7 @@ impl<T> Default for CacheEntry<T> {
}
impl<T> CacheEntry<T> {
pub fn new() -> Self {
pub const fn new() -> Self {
Self {
at: None,
value: None,
+21 -1
View File
@@ -24,6 +24,17 @@ pub fn cert_paths() -> (PathBuf, PathBuf) {
pub fn ensure_self_signed_cert() -> anyhow::Result<(PathBuf, PathBuf)> {
let (cert_path, key_path) = cert_paths();
if cert_path.exists() && key_path.exists() {
// Keys generated by agents older than 1.60 were written with the
// default umask (typically 0644): tighten them on startup.
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
if let Ok(meta) = fs::metadata(&key_path)
&& meta.permissions().mode() & 0o077 != 0
{
let _ = fs::set_permissions(&key_path, fs::Permissions::from_mode(0o600));
}
}
return Ok((cert_path, key_path));
}
fs::create_dir_all(cert_path.parent().unwrap())?;
@@ -79,7 +90,16 @@ pub fn ensure_self_signed_cert() -> anyhow::Result<(PathBuf, PathBuf)> {
let mut f = fs::File::create(&cert_path)?;
f.write_all(cert_pem.as_bytes())?;
let mut k = fs::File::create(&key_path)?;
// The private key must not be world-readable (File::create honors the
// umask, which typically yields 0644).
let mut key_opts = fs::OpenOptions::new();
key_opts.write(true).create(true).truncate(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
key_opts.mode(0o600);
}
let mut k = key_opts.open(&key_path)?;
k.write_all(key_pem.as_bytes())?;
println!(
+11 -1
View File
@@ -30,6 +30,11 @@ pub struct ProcessInfo {
#[derive(Debug, Clone, Serialize)]
pub struct Metrics {
/// Epoch ms when this snapshot was actually collected. The agent serves
/// TTL-cached snapshots, so the client needs the AGENT's sample time to
/// compute rates — measuring against client receive time turned cache
/// hits into a 0-then-2x sawtooth in the network graphs.
pub sampled_at_ms: u64,
pub cpu_total: f32,
pub cpu_per_core: Vec<f32>,
pub mem_total: u64,
@@ -93,7 +98,8 @@ pub struct ProcessMetricsResponse {
#[derive(Debug, Clone, Serialize)]
pub struct JournalEntry {
pub timestamp: String, // ISO 8601 formatted timestamp
pub timestamp: String, // RFC 3339 UTC, for display
pub timestamp_us: u64, // epoch microseconds, for sorting/formatting
pub priority: LogLevel,
pub message: String,
pub unit: Option<String>, // systemd unit name
@@ -120,5 +126,9 @@ pub struct JournalResponse {
pub entries: Vec<JournalEntry>,
pub total_count: u32,
pub truncated: bool,
/// journalctl's own explanation when the result is empty because of
/// journal ACCESS (not absence of logs) — e.g. a user-run agent asking
/// about a system service. None when entries exist or nothing to say.
pub notice: Option<String>,
pub cached_at: u64, // Unix timestamp when this data was cached
}
+79 -75
View File
@@ -16,9 +16,7 @@ use crate::metrics::{collect_disks, collect_fast_metrics, collect_processes_all}
use crate::proto::pb;
use crate::state::AppState;
// Compression threshold based on typical payload size
// Temporarily increased for testing - revert to 768 for production
//const COMPRESSION_THRESHOLD: usize = 50_000;
// Payloads at or below this many bytes are sent as-is; larger ones are gzipped.
const COMPRESSION_THRESHOLD: usize = 768;
// Reusable buffer for compression to avoid allocations
@@ -52,6 +50,66 @@ pub async fn ws_handler(
ws.on_upgrade(move |socket| handle_socket(socket, state))
}
/// Per-PID cache limits: entries older than MAX_AGE are swept on every
/// insert and the map is capped at MAX_ENTRIES (oldest evicted first), so a
/// client walking PIDs cannot grow agent memory without bound.
const PER_PID_CACHE_MAX_AGE: std::time::Duration = std::time::Duration::from_secs(60);
const PER_PID_CACHE_MAX_ENTRIES: usize = 64;
/// Serve a per-PID request from a TTL cache, collecting on miss. One home
/// for the logic that get_process_metrics and get_journal_entries used to
/// duplicate ~50 lines apiece.
async fn respond_per_pid_cached<T, Fut>(
socket: &mut WebSocket,
cache: &Mutex<HashMap<u32, crate::state::CacheEntry<T>>>,
pid: u32,
ttl: std::time::Duration,
request_name: &str,
collect: impl FnOnce() -> Fut,
) where
T: serde::Serialize + Clone,
Fut: std::future::Future<Output = Result<T, String>>,
{
{
let cache = cache.lock().await;
if let Some(entry) = cache.get(&pid)
&& entry.is_fresh(ttl)
&& let Some(v) = entry.get()
{
let _ = send_json(socket, v).await;
return;
}
}
match collect().await {
Ok(resp) => {
{
let mut cache = cache.lock().await;
cache.retain(|_, e| e.at.is_some_and(|t| t.elapsed() < PER_PID_CACHE_MAX_AGE));
while cache.len() >= PER_PID_CACHE_MAX_ENTRIES {
let oldest = cache.iter().min_by_key(|(_, e)| e.at).map(|(k, _)| *k);
match oldest {
Some(k) => cache.remove(&k),
None => break,
};
}
cache
.entry(pid)
.or_insert_with(crate::state::CacheEntry::new)
.set(resp.clone());
}
let _ = send_json(socket, &resp).await;
}
Err(err) => {
let error_response = serde_json::json!({
"error": err,
"request": request_name,
"pid": pid
});
let _ = send_json(socket, &error_response).await;
}
}
}
async fn handle_socket(mut socket: WebSocket, state: AppState) {
state
.client_count
@@ -126,84 +184,30 @@ async fn handle_socket(mut socket: WebSocket, state: AppState) {
if let Some(pid_str) = text.strip_prefix("get_process_metrics:")
&& let Ok(pid) = pid_str.parse::<u32>()
{
let ttl = std::time::Duration::from_millis(250); // 250ms TTL
// Check cache first
{
let cache = state.cache_process_metrics.lock().await;
if let Some(entry) = cache.get(&pid)
&& entry.is_fresh(ttl)
&& let Some(cached_response) = entry.get()
{
let _ = send_json(&mut socket, cached_response).await;
continue;
}
}
// Collect fresh data
match crate::metrics::collect_process_metrics(pid, &state).await {
Ok(response) => {
// Cache the response
{
let mut cache = state.cache_process_metrics.lock().await;
cache
.entry(pid)
.or_insert_with(crate::state::CacheEntry::new)
.set(response.clone());
}
let _ = send_json(&mut socket, &response).await;
}
Err(err) => {
let error_response = serde_json::json!({
"error": err,
"request": "get_process_metrics",
"pid": pid
});
let _ = send_json(&mut socket, &error_response).await;
}
}
respond_per_pid_cached(
&mut socket,
&state.cache_process_metrics,
pid,
std::time::Duration::from_millis(250),
"get_process_metrics",
|| crate::metrics::collect_process_metrics(pid, &state),
)
.await;
}
}
Message::Text(ref text) if text.starts_with("get_journal_entries:") => {
if let Some(pid_str) = text.strip_prefix("get_journal_entries:")
&& let Ok(pid) = pid_str.parse::<u32>()
{
let ttl = std::time::Duration::from_secs(1); // 1s TTL
// Check cache first
{
let cache = state.cache_journal_entries.lock().await;
if let Some(entry) = cache.get(&pid)
&& entry.is_fresh(ttl)
&& let Some(cached_response) = entry.get()
{
let _ = send_json(&mut socket, cached_response).await;
continue;
}
}
// Collect fresh data
match crate::metrics::collect_journal_entries(pid) {
Ok(response) => {
// Cache the response
{
let mut cache = state.cache_journal_entries.lock().await;
cache
.entry(pid)
.or_insert_with(crate::state::CacheEntry::new)
.set(response.clone());
}
let _ = send_json(&mut socket, &response).await;
}
Err(err) => {
let error_response = serde_json::json!({
"error": err,
"request": "get_journal_entries",
"pid": pid
});
let _ = send_json(&mut socket, &error_response).await;
}
}
respond_per_pid_cached(
&mut socket,
&state.cache_journal_entries,
pid,
std::time::Duration::from_secs(1),
"get_journal_entries",
|| crate::metrics::collect_journal_entries(pid),
)
.await;
}
}
Message::Close(_) => break,
+1
View File
@@ -42,6 +42,7 @@ async fn test_process_cache_ttl() {
};
let journal_response = JournalResponse {
notice: None,
entries: vec![],
total_count: 0,
truncated: false,
+18 -2
View File
@@ -33,7 +33,7 @@ async fn test_collect_journal_entries_self() {
// Test collecting journal entries for our own process
let pid = process::id();
match collect_journal_entries(pid) {
match collect_journal_entries(pid).await {
Ok(response) => {
assert!(response.cached_at > 0);
println!(
@@ -74,7 +74,7 @@ async fn test_collect_journal_entries_invalid_pid() {
// Test with an invalid PID - journalctl might still return empty results
let invalid_pid = 999999;
match collect_journal_entries(invalid_pid) {
match collect_journal_entries(invalid_pid).await {
Ok(response) => {
println!(
"✓ Journal query completed for invalid PID {} (empty result expected): {} entries",
@@ -87,3 +87,19 @@ async fn test_collect_journal_entries_invalid_pid() {
}
}
}
/// The Command & Details pane went blank when the minimal-refresh
/// optimization dropped cmd from the detail endpoint's refresh kind.
#[tokio::test]
async fn test_process_metrics_include_command() {
let state = AppState::new();
let pid = std::process::id();
let resp = collect_process_metrics(pid, &state)
.await
.expect("collect self");
assert!(
!resp.process.command.is_empty(),
"command should not be empty for self (cmdline is always readable)"
);
println!("command = {}", resp.process.command);
}