Compare commits

...

15 Commits

Author SHA1 Message Date
jasonwitty 83c5f6ebcf docs(installer): use a durable ref in the usage example
CI / build (ubuntu-latest) (push) Has been cancelled
CI / build (windows-latest) (push) Has been cancelled
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-21 12:28:48 -07:00
jasonwitty 43ce4f1aaa 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>
2026-08-21 11:42:51 -07:00
jasonwitty d8cceb1795 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>
2026-08-21 11:19:04 -07:00
jasonwitty 2a7951cf65 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>
2026-08-21 11:16:41 -07:00
jasonwitty 8e0effe361 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>
2026-08-20 12:42:49 -07:00
jasonwitty f5286008b2 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>
2026-08-20 12:25:46 -07:00
jasonwitty f61b232e42 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>
2026-08-20 10:23:45 -07:00
jasonwitty 2c773eaa6c test: add notice field to cache test initializer
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-19 20:25:41 -07:00
jasonwitty d653cbcadc 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>
2026-08-19 20:24:56 -07:00
jasonwitty 746ca4cf58 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>
2026-08-19 16:52:31 -07:00
jasonwitty bf6ac877c2 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>
2026-08-19 14:58:39 -07:00
jasonwitty 0c800f83f9 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>
2026-08-19 14:52:35 -07:00
jasonwitty e4b0d9b9f6 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>
2026-08-19 14:43:37 -07:00
jasonwitty fbc788c799 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>
2026-08-19 14:31:11 -07:00
jasonwitty 679a50b2e8 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>
2026-08-19 14:27:40 -07:00
38 changed files with 6243 additions and 4062 deletions
+3 -2
View File
@@ -1,6 +1,7 @@
/target
# Any crate's build directory, including standalone sub-crates
# (zellij_socktop_plugin, socktop_wasm_test) that live outside the workspace.
target/
.vscode/
/socktop-wasm-test/target
/.cargo/
# Documentation files from development sessions (context-specific, not for public repo)
+57
View File
@@ -0,0 +1,57 @@
# Changelog
## 1.60.0 — unreleased
Everything since `v1.50.0`. Applies to all three crates (`socktop`, `socktop_agent`, `socktop_connector`), which move to 1.60.0 together.
### Security
- **Certificate pinning is now real.** With `--verify-hostname` off (the default), the client previously accepted *any* server certificate — the `--tls-ca` file was never consulted. The presented certificate must now be byte-identical to one in the pinned PEM (multi-cert files supported for rotation). If you use TLS, update the client: earlier versions are MITM-able despite the pinning documentation. (housekeeping-p2)
- `key.pem` is created with mode 0600 (was world-readable 0644); agents also tighten existing keys on startup. (housekeeping-p2)
- The agent's per-PID caches now evict (60s age / 64 entries); previously they grew without bound. (housekeeping-p2)
### Performance
- Agent CPU on GPU machines cut ~6× (measured 23.5 → 4.0 ms/s at default polling): GPU collection moved to a dedicated worker thread that keeps the NVML session open instead of re-initializing it every 1.5 s on the async runtime. (housekeeping-p2)
- `journalctl` no longer blocks the agent's async workers. (housekeeping-p2)
- Cached "no temp sensor / no GPU" results count as fresh — no more per-request rescans on hosts without them. (housekeeping-p2)
- Nagle disabled on all connection paths (small request/response frames). (housekeeping-p2)
### TUI
- **Compact layout for small windows**: when the window is too short for the Disks pane, Disks is dropped, Memory/Swap go side by side, GPU collapses to one line (omitted if absent), and the reclaimed rows keep the CPU graph and per-core bars visible. `--compact` pins it. (#37)
- **Width-aware text**: header, CPU title, and process table shed detail by priority as the terminal narrows instead of overwriting each other; process Name column is now the last to go, not the first. Fixed sort-header clicks landing up to 4 columns off. (#38)
- **Responsive input**: keys and mouse are handled within ~30 ms instead of queueing for a full metrics interval. (housekeeping-p2)
- **No more freezes**: all requests carry a 5 s timeout; a dead connection shows the reconnect modal (with working `q`) instead of hanging the UI. Consecutive timeouts surface a persistent "agent not responding" error. (housekeeping-p2)
- Old agents without the per-process endpoints once again show "Agent Update Required" instead of a reconnect loop. (housekeeping-p2)
- Journal pane distinguishes "no entries" from "no journal access" (e.g. user-run/demo agents) and shows journalctl's hint plus the fix. (housekeeping-p2)
- Scatter-plot axes align correctly for large CPU-time values. (housekeeping-p2)
- Demo mode explains how to install `socktop_agent` when the binary is missing. (#36)
### Correctness
- Process/child CPU times were sent as ms but displayed as µs — values rendered 1000× too small in the details modal. (housekeeping-p2)
- Non-Linux per-process CPU% no longer truncates multi-core usage (clamp after divide). (housekeeping-p2)
- Journal timestamps are real RFC 3339 UTC with numeric sorting (additive `timestamp_us`). (housekeeping-p2)
- Partition detection uses `/sys/block` on Linux — whole-disk filesystems (`nvme0n1`, `zram1`) are no longer misclassified as partitions. (housekeeping-p2)
- Network rates use agent-side sample timestamps (additive `sampled_at_ms`), eliminating rate sawtooth from TTL-cached snapshots; falls back to the client clock with older agents. (housekeeping-p2)
- The details modal's Command/exe/cwd fields are populated again (dropped by an earlier refresh optimization). (housekeeping-p2)
- Non-ASCII device names no longer panic the disk pane. (housekeeping-p2)
### Wire format (additive only — old/new client-agent pairs keep working)
- `Metrics.sampled_at_ms` (epoch ms of actual collection)
- `JournalEntry.timestamp_us` (epoch µs), `JournalEntry.timestamp` now RFC 3339
- `JournalResponse.notice` (journal-access hint)
### Internal / packaging
- ratatui 0.28 → 0.30 (#33); aws-lc-rs advisories patched (#34); Debian packaging for the agent (#25); assorted dependabot bumps.
- ~3,100 lines of dead code removed, including an orphaned pre-refactor copy of the connector.
- `socktop` consumes `socktop_connector` via a path+version dep — connector changes are testable in-repo before publishing.
- wasm examples build against the in-repo connector; note `zellij_socktop_plugin` has pre-existing compile errors and needs its own rework.
### Upgrade notes
- **Release/publish order**: `socktop_connector``socktop` → agent packages.
- Clients older than 1.60 work against 1.60 agents and vice versa; the security fix is client-side, so prioritize client updates where TLS is used.
Generated
+5 -26
View File
@@ -2412,7 +2412,7 @@ dependencies = [
[[package]]
name = "socktop"
version = "1.50.0"
version = "1.60.0"
dependencies = [
"anyhow",
"assert_cmd",
@@ -2422,8 +2422,7 @@ dependencies = [
"ratatui",
"serde",
"serde_json",
"socktop_connector 1.50.0 (registry+https://github.com/rust-lang/crates.io-index)",
"sysinfo",
"socktop_connector",
"tempfile",
"tokio",
"unicode-width",
@@ -2432,7 +2431,7 @@ dependencies = [
[[package]]
name = "socktop_agent"
version = "1.50.2"
version = "1.60.0"
dependencies = [
"anyhow",
"assert_cmd",
@@ -2442,6 +2441,7 @@ dependencies = [
"futures-util",
"gfxinfo",
"hostname",
"nvml-wrapper",
"once_cell",
"prost",
"prost-build",
@@ -2463,7 +2463,7 @@ dependencies = [
[[package]]
name = "socktop_connector"
version = "1.50.0"
version = "1.60.0"
dependencies = [
"flate2",
"futures-util",
@@ -2484,27 +2484,6 @@ dependencies = [
"web-sys",
]
[[package]]
name = "socktop_connector"
version = "1.50.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "61ea6a5733e71da6d5c94d23265b85f7041305bca51e6c33e7104464444047bc"
dependencies = [
"flate2",
"futures-util",
"prost",
"prost-build",
"protoc-bin-vendored",
"rustls",
"rustls-pemfile",
"serde",
"serde_json",
"thiserror 2.0.17",
"tokio",
"tokio-tungstenite 0.24.0",
"url",
]
[[package]]
name = "stable_deref_trait"
version = "1.2.1"
+2 -1
View File
@@ -416,6 +416,7 @@ Tip: If only the binary changed, restart is enough. If the unit file changed, ru
```json
{
"sampled_at_ms": 1786752000123,
"cpu_total": 12.4,
"cpu_per_core": [11.2, 15.7],
"mem_total": 33554432,
@@ -475,7 +476,7 @@ socktop --tls-ca /path/to/agent/cert.pem wss://HOST:8443/ws
Notes:
- Do not copy the private key off the server; only the cert.pem is needed by clients.
- When --tls-ca/-t is supplied, the client autoupgrades ws:// to wss:// to avoid protocol mismatch.
- Hostname (SAN) verification is DISABLED by default (the cert is still pinned). Use `--verify-hostname` to enable strict SAN checking.
- Hostname (SAN) verification is DISABLED by default; instead the client PINS the certificate: the agent must present a cert byte-identical to one in your `--tls-ca` file (expiry is ignored in this mode — you pinned that exact cert). Use `--verify-hostname` to switch to strict chain + SAN validation instead.
- You can run multiple clients with different cert paths by passing --tls-ca per invocation.
---
+246
View File
@@ -0,0 +1,246 @@
#!/usr/bin/env bash
# Build socktop + socktop_agent from source and install them.
#
# Works on Linux (x86_64, arm64/armv7, riscv64) and macOS. Handles fresh
# installs and upgrades; if a systemd socktop-agent service is present, its
# binary is replaced in place and the service restarted.
#
# ./scripts/install.sh # build HEAD of the repo you're in
# ./scripts/install.sh --ref v1.60.0 # build a tag/branch (clones if needed)
# ./scripts/install.sh --ref master # or any branch
# ./scripts/install.sh --prefix ~/.local/bin --no-service
#
set -euo pipefail
REPO_URL="https://github.com/jasonwitty/socktop.git"
REF=""
PREFIX=""
NO_SERVICE=0
SRC_DIR="${SOCKTOP_SRC_DIR:-$HOME/.cache/socktop-src}"
while [ $# -gt 0 ]; do
case "$1" in
--ref) REF="$2"; shift 2 ;;
--prefix) PREFIX="$2"; shift 2 ;;
--no-service) NO_SERVICE=1; shift ;;
-h|--help) grep '^#' "$0" | sed 's/^# \{0,1\}//'; exit 0 ;;
*) echo "unknown argument: $1" >&2; exit 2 ;;
esac
done
say() { printf '\033[1;36m==>\033[0m %s\n' "$*"; }
warn() { printf '\033[1;33mwarn:\033[0m %s\n' "$*" >&2; }
die() { printf '\033[1;31merror:\033[0m %s\n' "$*" >&2; exit 1; }
# The entire remainder runs inside main(), invoked on the LAST line. This
# makes the script safe against being MODIFIED WHILE RUNNING: when executed
# from the clone it manages, the git checkout below replaces this very file,
# and bash reads scripts lazily by byte offset — without this wrapper it
# resumes parsing the NEW file at the OLD offset and executes an arbitrary
# tail of it (observed: the fresh-service path ran on a host whose unit
# already existed). With main(), the whole script is parsed before any of
# it executes.
main() {
OS="$(uname -s)"
ARCH="$(uname -m)"
# ---------- toolchain ----------
command -v git >/dev/null || die "git is required"
if ! command -v cargo >/dev/null; then
# rustup may be installed but not on PATH in this shell
[ -f "$HOME/.cargo/env" ] && . "$HOME/.cargo/env"
fi
if ! command -v cargo >/dev/null; then
say "Rust toolchain not found — installing via rustup (stable, default profile)"
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y --profile minimal
. "$HOME/.cargo/env"
fi
command -v cc >/dev/null || warn "no C compiler found (apt: build-essential / brew: xcode-select --install) — the build may fail"
case "$ARCH" in
riscv64*)
# protoc-bin-vendored ships no riscv64 binary; the build falls back to
# the system protoc (see build.rs).
command -v protoc >/dev/null || die "riscv64 needs a system protoc: sudo apt install protobuf-compiler"
;;
esac
# ---------- source ----------
# If run from inside a socktop checkout and no --ref given, build that tree
# as-is (whatever is checked out, including local changes).
if [ -z "$REF" ] && git rev-parse --show-toplevel >/dev/null 2>&1 \
&& grep -qs '^name = "socktop"' "$(git rev-parse --show-toplevel)/socktop/Cargo.toml" 2>/dev/null; then
SRC_DIR="$(git rev-parse --show-toplevel)"
say "Building the current checkout: $SRC_DIR ($(git -C "$SRC_DIR" describe --always --dirty 2>/dev/null))"
else
REF="${REF:-master}"
if [ ! -d "$SRC_DIR/.git" ]; then
say "Cloning $REPO_URL -> $SRC_DIR"
git clone "$REPO_URL" "$SRC_DIR"
fi
say "Checking out $REF"
git -C "$SRC_DIR" fetch --tags origin
git -C "$SRC_DIR" checkout -q "$REF"
# fast-forward when REF is a branch
git -C "$SRC_DIR" merge --ff-only "origin/$REF" >/dev/null 2>&1 || true
fi
# ---------- build ----------
say "Building release binaries (this can take a while on SBCs)"
( cd "$SRC_DIR" && cargo build --release -p socktop -p socktop_agent )
CLIENT="$SRC_DIR/target/release/socktop"
AGENT="$SRC_DIR/target/release/socktop_agent"
# ---------- install ----------
if [ -z "$PREFIX" ]; then
PREFIX="/usr/local/bin"
fi
SUDO=""
if [ ! -w "$PREFIX" ]; then
if command -v sudo >/dev/null; then SUDO="sudo"; else
PREFIX="$HOME/.local/bin"; mkdir -p "$PREFIX"
warn "no sudo — installing to $PREFIX (ensure it is on your PATH)"
fi
fi
say "Installing to $PREFIX"
$SUDO install -m 755 "$CLIENT" "$PREFIX/socktop"
$SUDO install -m 755 "$AGENT" "$PREFIX/socktop_agent"
# Update every other copy on PATH as well. A stale `cargo install` in
# ~/.cargo/bin would otherwise SHADOW the fresh binary (~/.cargo/bin
# usually precedes /usr/local/bin on PATH), leaving `socktop --version`
# stuck on the old release after a "successful" install.
update_path_copies() {
local name="$1" src="$2" copy dir
# type -ap lists every match on PATH (bash builtin, symlinks not resolved)
for copy in $(type -ap "$name" | sort -u); do
[ "$copy" = "$PREFIX/$name" ] && continue
dir="$(dirname "$copy")"
say "Updating additional copy on PATH: $copy"
if [ -w "$copy" ] || [ -w "$dir" ]; then
install -m 755 "$src" "$copy"
else
# Non-fatal: an un-updatable extra copy shouldn't kill the install,
# but the user must know it may shadow the fresh binary.
$SUDO install -m 755 "$src" "$copy" || warn "could not update $copy — it may shadow $PREFIX/$name"
fi
done
}
update_path_copies socktop "$CLIENT"
update_path_copies socktop_agent "$AGENT"
# ---------- systemd service (Linux only) ----------
# System-level operations (unit files, users, service control) need root no
# matter where the binaries were installed — decide independently of PREFIX.
SYS_SUDO=""
if [ "$(id -u)" -ne 0 ]; then
if command -v sudo >/dev/null; then SYS_SUDO="sudo"; else SYS_SUDO="__none__"; fi
fi
if [ "$SYS_SUDO" = "__none__" ] && [ "$NO_SERVICE" -eq 0 ]; then
warn "no sudo available — skipping systemd service management"
NO_SERVICE=1
fi
if [ "$OS" = "Linux" ] && [ "$NO_SERVICE" -eq 0 ] && command -v systemctl >/dev/null; then
if systemctl cat socktop-agent.service >/dev/null 2>&1; then
# UPGRADE: the unit file is the operator's (SSL, tokens, ports may be
# configured there) — never overwrite it. Only the binary it points at
# is replaced, then the service is restarted.
say "Existing socktop-agent.service found — preserving unit file, refreshing binary"
UNIT_BIN="$(systemctl show -p ExecStart socktop-agent.service 2>/dev/null \
| sed -n 's/.*path=\([^ ;]*\).*/\1/p' | head -1)"
if [ -n "$UNIT_BIN" ] && [ "$UNIT_BIN" != "$PREFIX/socktop_agent" ]; then
$SYS_SUDO systemctl stop socktop-agent.service
$SYS_SUDO install -m 755 "$AGENT" "$UNIT_BIN"
$SYS_SUDO systemctl start socktop-agent.service
else
$SYS_SUDO systemctl restart socktop-agent.service
fi
else
# FRESH INSTALL: unit + the system user it runs as + its state dir,
# then enable and start. Mirrors the deb package's postinst and
# https://www.socktop.io/assets/docs/installation/agent-service.html
say "No socktop-agent.service found — installing and enabling it"
if ! getent group socktop >/dev/null; then
$SYS_SUDO groupadd --system socktop
fi
if ! getent passwd socktop >/dev/null; then
NOLOGIN="$(command -v nologin || echo /usr/sbin/nologin)"
$SYS_SUDO useradd --system -g socktop -d /var/lib/socktop -M -s "$NOLOGIN" socktop
fi
$SYS_SUDO mkdir -p /var/lib/socktop
$SYS_SUDO chown socktop:socktop /var/lib/socktop
$SYS_SUDO chmod 755 /var/lib/socktop
UNIT_TMP="$(mktemp)"
if [ -f "$SRC_DIR/docs/socktop-agent.service" ]; then
cp "$SRC_DIR/docs/socktop-agent.service" "$UNIT_TMP"
else
# Fallback for refs that predate docs/socktop-agent.service
cat > "$UNIT_TMP" <<'UNIT'
[Unit]
Description=Socktop agent
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
ExecStart=/usr/local/bin/socktop_agent --port 3000
Environment=RUST_LOG=info
# Optional auth:
# Environment=SOCKTOP_TOKEN=changeme
# TLS (self-signed cert on first run, default port 8443):
# Environment=SOCKTOP_ENABLE_SSL=1
Restart=on-failure
User=socktop
Group=socktop
NoNewPrivileges=true
[Install]
WantedBy=multi-user.target
UNIT
fi
# Pick the agent port: 3000 by default, but NEVER bind onto a port that
# something else already holds (e.g. Gitea/Umami and friends love 3000)
# — that puts the fresh service straight into a crash-restart loop.
AGENT_PORT=""
for p in 3000 3001 3010 3231 3232; do
if ! ss -tln 2>/dev/null | awk '{print $4}' | grep -q ":${p}\$"; then
AGENT_PORT="$p"
break
fi
done
if [ -z "$AGENT_PORT" ]; then
AGENT_PORT=3000
warn "no free port among the defaults — using 3000; edit the unit if the service fails to start"
elif [ "$AGENT_PORT" != "3000" ]; then
warn "port 3000 is already in use by another service — configuring the agent on port $AGENT_PORT"
fi
# Point ExecStart at wherever this run installed the agent, on the chosen port.
sed -i.bak -e "s|^ExecStart=[^ ]*socktop_agent|ExecStart=$PREFIX/socktop_agent|" \
-e "s|--port [0-9]*|--port $AGENT_PORT|" "$UNIT_TMP"
rm -f "$UNIT_TMP.bak"
$SYS_SUDO install -o root -g root -m 0644 "$UNIT_TMP" /etc/systemd/system/socktop-agent.service
rm -f "$UNIT_TMP"
$SYS_SUDO systemctl daemon-reload
$SYS_SUDO systemctl enable --now socktop-agent.service
say "Service installed — agent URL: ws://$(hostname):$AGENT_PORT/ws"
say "To enable TLS or a token, edit /etc/systemd/system/socktop-agent.service, then: sudo systemctl daemon-reload && sudo systemctl restart socktop-agent"
fi
sleep 1
systemctl --no-pager -l status socktop-agent.service | head -5 || true
fi
say "Installed:"
"$PREFIX/socktop" --version
"$PREFIX/socktop_agent" --version
say "Active on PATH: $(type -p socktop || true) / $(type -p socktop_agent || true)"
socktop --version
}
# exit in the same parse unit as the call: after main returns, bash must not
# read another byte from this (possibly replaced) file.
main "$@"; exit $?
+2 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "socktop"
version = "1.50.0"
version = "1.60.0"
authors = ["Jason Witty <jasonpwitty+socktop@proton.me>"]
description = "Remote system monitor over WebSocket, TUI like top"
edition = "2024"
@@ -11,7 +11,7 @@ repository = "https://github.com/jasonwitty/socktop"
[dependencies]
# socktop connector for agent communication
socktop_connector = "1.50.0"
socktop_connector = { version = "1.60.0", path = "../socktop_connector" }
tokio = { workspace = true }
futures-util = { workspace = true }
@@ -23,7 +23,6 @@ crossterm = { workspace = true }
unicode-width = { workspace = true }
anyhow = { workspace = true }
dirs-next = { workspace = true }
sysinfo = { workspace = true }
[dev-dependencies]
assert_cmd = "2.0"
+247 -51
View File
@@ -51,6 +51,11 @@ use socktop_connector::{
const MIN_METRICS_INTERVAL_MS: u64 = 100;
const MIN_PROCESSES_INTERVAL_MS: u64 = 200;
/// Budget for one request/response round trip. Replies are matched to
/// requests by order, so a request that never answers would otherwise hang
/// `ws.next()` forever and freeze the TUI (raw mode even eats Ctrl+C).
const REQUEST_TIMEOUT: Duration = Duration::from_secs(5);
/// Drop duplicate-name entries from a disks payload (the agent occasionally
/// reports a partition twice). Done once when fresh disk data arrives so the
/// per-frame draw path doesn't have to rebuild a HashSet.
@@ -60,6 +65,13 @@ fn dedup_disks(disks: &mut Vec<socktop_connector::DiskInfo>) {
disks.retain(|d| seen.insert(d.name.clone()));
}
/// Outcome of draining input: keep going, or restart the event loop because a
/// reconnect installed a replacement connection.
enum InputFlow {
Continue,
RestartConnection,
}
#[derive(Debug, Clone, PartialEq)]
pub enum ConnectionState {
Connected,
@@ -80,6 +92,14 @@ pub struct App {
// Network totals snapshot + histories of KB/s
last_net_totals: Option<(u64, u64, Instant)>,
// Agent-side sample timestamp of the previous snapshot (1.60+ agents).
last_net_sampled_at_ms: Option<u64>,
// Consecutive metrics-request timeouts. One timeout gets a silent stream
// refresh; a second in a row means the agent accepts connections but
// never answers, and deserves a persistent error instead of an invisible
// reconnect loop that starves the UI.
consecutive_request_timeouts: u32,
rx_hist: VecDeque<u64>,
tx_hist: VecDeque<u64>,
rx_peak: u64,
@@ -180,6 +200,8 @@ impl App {
cpu_hist_sum: 0,
per_core_hist: PerCoreHistory::new(60),
last_net_totals: None,
last_net_sampled_at_ms: None,
consecutive_request_timeouts: 0,
rx_hist: VecDeque::with_capacity(600),
tx_hist: VecDeque::with_capacity(600),
rx_peak: 0,
@@ -349,6 +371,40 @@ impl App {
}
}
/// A request produced no reply in time. Any late reply would desync every
/// subsequent request/response on this stream (replies are matched to
/// requests purely by order), so treat the connection as poisoned and go
/// through the standard reconnect flow — a fresh stream is realigned by
/// construction.
async fn poison_connection(&mut self, what: &str) {
self.show_connection_error(format!("{what}; reconnecting…"));
self.retry_connection().await;
}
/// Replace the connection WITHOUT any modal or state churn.
///
/// For timeouts on the optional per-process endpoints: an old agent
/// ignores those messages entirely (no late reply, so no desync), but a
/// merely-slow agent would desync the stream — indistinguishable at
/// timeout time, so we still swap to a fresh stream, silently. The
/// ProcessDetails modal keeps showing its "Agent Update Required"
/// message instead of being buried under a connection-error modal.
/// Only a failed reconnect (connection genuinely dead) surfaces loudly.
async fn quiet_reconnect(&mut self) {
let tls_ca_ref = self.tls_ca.as_deref();
match self
.try_connect(&self.ws_url, tls_ca_ref, self.verify_hostname)
.await
{
Ok(ws) => {
self.replacement_connection = Some(ws);
}
Err(e) => {
self.show_connection_error(format!("Reconnect failed: {e}"));
}
}
}
/// Mark connection as successful and dismiss any error modals
pub fn mark_connected(&mut self) {
if self.connection_state != ConnectionState::Connected {
@@ -666,17 +722,21 @@ impl App {
}
}
async fn run_event_loop_iteration<B: ratatui::backend::Backend>(
/// Drains and handles every queued terminal event (keys, mouse). Returns
/// whether the caller must restart the event loop on a replacement
/// connection. Extracted from the loop body so the tick wait can process
/// input at ~30ms latency instead of letting it queue for a whole
/// metrics interval.
async fn drain_input<B: ratatui::backend::Backend>(
&mut self,
terminal: &mut Terminal<B>,
ws: &mut SocktopConnector,
) -> Result<(), Box<dyn std::error::Error>>
) -> Result<InputFlow, Box<dyn std::error::Error>>
where
<B as ratatui::backend::Backend>::Error: 'static,
{
loop {
// Input (non-blocking)
while event::poll(Duration::from_millis(10))? {
// Drain everything already queued; the caller has verified (or will
// verify via poll) that input is or may be pending.
while event::poll(Duration::ZERO)? {
match event::read()? {
Event::Key(k) => {
// Handle modal input first - if a modal consumes the input, don't process normal keys
@@ -691,9 +751,8 @@ impl App {
self.retry_connection().await;
// Check if retry succeeded and we have a replacement connection
if self.replacement_connection.is_some() {
// Signal that we want to restart with new connection
// Return from this iteration so the outer loop can restart
return Ok(());
// Restart the outer loop on the new connection
return Ok(InputFlow::RestartConnection);
}
continue; // Skip normal key processing
}
@@ -846,11 +905,7 @@ impl App {
// If process selection didn't handle it, use CPU scrolling
if !process_handled {
per_core_handle_key(
&mut self.per_core_scroll,
k,
content.height as usize,
);
per_core_handle_key(&mut self.per_core_scroll, k, content.height as usize);
}
// Auto-scroll to keep selected process visible
@@ -861,8 +916,7 @@ impl App {
let idxs = &self.procs_filtered;
// Find the display position of the selected process in filtered list
if let Some(display_pos) =
idxs.iter().position(|&idx| idx == selected_idx)
if let Some(display_pos) = idxs.iter().position(|&idx| idx == selected_idx)
{
// Calculate viewport size
// Account for: borders (2) + header (1) + search box if active (3)
@@ -996,6 +1050,26 @@ impl App {
}
}
Ok(InputFlow::Continue)
}
async fn run_event_loop_iteration<B: ratatui::backend::Backend>(
&mut self,
terminal: &mut Terminal<B>,
ws: &mut SocktopConnector,
) -> Result<(), Box<dyn std::error::Error>>
where
<B as ratatui::backend::Backend>::Error: 'static,
{
loop {
// Input: drain anything already queued
if matches!(
self.drain_input(terminal).await?,
InputFlow::RestartConnection
) {
return Ok(());
}
// Check for automatic retry (every 30 seconds)
if self.should_auto_retry() {
self.auto_retry_connection().await;
@@ -1010,29 +1084,66 @@ impl App {
break;
}
// Fetch and update
match ws.request(AgentRequest::Metrics).await {
Ok(AgentResponse::Metrics(m)) => {
// Paint the current state BEFORE fetching: a request can stall for
// the full 5s timeout, and an iteration that ends in a poisoned-
// stream restart never reaches the draw at the bottom — without
// this, an agent that never answers left the screen permanently
// blank. (ratatui diffs make an unchanged repaint nearly free.)
terminal.draw(|f| self.draw(f))?;
// Fetch and update. Skipped while disconnected — the retry paths
// (manual 'r' or the 30s auto-retry) own recovery, and hammering a
// dead socket with 5s-timeout requests would stall the loop. The
// shared draw + responsive wait below still run, so the error
// modal stays live and input stays snappy.
if self.connection_state == ConnectionState::Connected {
match timeout(REQUEST_TIMEOUT, ws.request(AgentRequest::Metrics)).await {
Err(_) => {
self.consecutive_request_timeouts += 1;
if self.consecutive_request_timeouts >= 2 {
// The agent accepts connections but never answers
// (wrong protocol era, or wedged): reconnecting
// can't help, so surface a persistent error and
// leave recovery to the manual/auto retry paths.
self.show_connection_error(
"Agent is not responding to requests".to_string(),
);
} else {
self.poison_connection("Metrics request timed out").await;
}
}
Ok(Ok(AgentResponse::Metrics(m))) => {
self.mark_connected(); // Mark as connected on successful request
self.consecutive_request_timeouts = 0;
self.update_with_metrics(m);
// Only poll processes every 2s
if self.last_procs_poll.elapsed() >= self.procs_interval {
let mut updated = false;
if let Ok(AgentResponse::Processes(procs)) =
ws.request(AgentRequest::Processes).await
&& let Some(mm) = self.last_metrics.as_mut()
match timeout(REQUEST_TIMEOUT, ws.request(AgentRequest::Processes))
.await
{
Err(_) => {
self.poison_connection("Processes request timed out").await;
}
Ok(Ok(AgentResponse::Processes(procs))) => {
if let Some(mm) = self.last_metrics.as_mut() {
mm.top_processes = procs.top_processes;
mm.process_count = Some(procs.process_count);
updated = true;
}
}
// Request error or wrong type: keep stale rows; a
// broken socket surfaces on the next metrics tick.
Ok(_) => {}
}
if updated {
self.invalidate_procs_filter();
// Rebuild the pre-formatted row cache for the next
// ~N frames. Done once per poll, not per frame.
if let Some(mm) = self.last_metrics.as_ref() {
self.procs_row_peak_cpu = crate::ui::processes::rebuild_row_cache(
self.procs_row_peak_cpu =
crate::ui::processes::rebuild_row_cache(
mm,
&mut self.procs_row_cache,
);
@@ -1042,25 +1153,40 @@ impl App {
}
// Only poll disks every 5s
if self.last_disks_poll.elapsed() >= self.disks_interval {
if let Ok(AgentResponse::Disks(mut disks)) =
ws.request(AgentRequest::Disks).await
&& let Some(mm) = self.last_metrics.as_mut()
if self.connection_state == ConnectionState::Connected
&& self.last_disks_poll.elapsed() >= self.disks_interval
{
match timeout(REQUEST_TIMEOUT, ws.request(AgentRequest::Disks)).await {
Err(_) => {
self.poison_connection("Disks request timed out").await;
}
Ok(Ok(AgentResponse::Disks(mut disks))) => {
if let Some(mm) = self.last_metrics.as_mut() {
dedup_disks(&mut disks);
mm.disks = disks;
}
}
Ok(_) => {}
}
self.last_disks_poll = Instant::now();
}
// Poll process details when modal is active and process is selected
if let Some(pid) = self.selected_process_pid {
if let Some(pid) = self.selected_process_pid
&& self.connection_state == ConnectionState::Connected
{
// Check if ProcessDetails modal is currently active
if let Some(crate::ui::modal::ModalType::ProcessDetails { .. }) =
self.modal_manager.current_modal()
{
// Poll process details every 500ms when modal is active
if self.last_process_details_poll.elapsed()
// Poll process details every 500ms when modal is
// active. Skipped once the agent is known not to
// support the endpoint (flag resets when the modal
// closes or the selection changes, so a one-off
// timeout doesn't disable details for the session).
if self.connection_state == ConnectionState::Connected
&& !self.process_details_unsupported
&& self.last_process_details_poll.elapsed()
>= self.process_details_interval
{
// Use timeout to prevent blocking the event loop
@@ -1078,12 +1204,16 @@ impl App {
cpu_usage,
600,
);
self.process_cpu_history_sum = self.process_cpu_history_sum
+ cpu_usage
self.process_cpu_history_sum =
self.process_cpu_history_sum + cpu_usage
- evicted_cpu.unwrap_or(0.0);
let mem_bytes = details.process.mem_bytes;
push_capped(&mut self.process_mem_history, mem_bytes, 600);
push_capped(
&mut self.process_mem_history,
mem_bytes,
600,
);
// Track maximum memory usage
if mem_bytes > self.max_process_mem_bytes {
@@ -1092,8 +1222,8 @@ impl App {
// I/O bytes from agent are cumulative, calculate deltas
if let Some(read) = details.process.read_bytes {
let delta = if let Some(last) = self.last_io_read_bytes
{
let delta =
if let Some(last) = self.last_io_read_bytes {
read.saturating_sub(last)
} else {
0 // First sample, no delta available
@@ -1106,8 +1236,8 @@ impl App {
self.last_io_read_bytes = Some(read);
}
if let Some(write) = details.process.write_bytes {
let delta = if let Some(last) = self.last_io_write_bytes
{
let delta =
if let Some(last) = self.last_io_write_bytes {
write.saturating_sub(last)
} else {
0 // First sample, no delta available
@@ -1123,11 +1253,21 @@ impl App {
self.process_details = Some(details);
self.process_details_unsupported = false;
}
Ok(Err(_)) | Err(_) => {
// Agent doesn't support this feature or timeout occurred
// Mark as unsupported so we can show appropriate message
Ok(Err(_)) => {
// Agent responded with an error: endpoint
// not supported.
self.process_details_unsupported = true;
}
Err(_) => {
// No reply at all: old agents IGNORE
// this message, so show the "Agent
// Update Required" state and refresh
// the stream quietly (a merely-slow
// agent's late reply would otherwise
// desync it).
self.process_details_unsupported = true;
self.quiet_reconnect().await;
}
Ok(Ok(_)) => {
// Wrong response type
self.process_details_unsupported = true;
@@ -1136,8 +1276,13 @@ impl App {
self.last_process_details_poll = Instant::now();
}
// Poll journal entries every 5s when modal is active
if self.last_journal_poll.elapsed() >= self.journal_interval {
// Poll journal entries every 5s when modal is active.
// Gated on the same unsupported flag: agents that lack
// process details lack the journal endpoint too.
if self.connection_state == ConnectionState::Connected
&& !self.process_details_unsupported
&& self.last_journal_poll.elapsed() >= self.journal_interval
{
// Use timeout to prevent blocking the event loop
match timeout(
Duration::from_millis(2000),
@@ -1148,9 +1293,14 @@ impl App {
Ok(Ok(AgentResponse::JournalEntries(journal))) => {
self.journal_entries = Some(journal);
}
Ok(Err(_)) | Err(_) | Ok(Ok(_)) => {
// Agent doesn't support this feature, error occurred, or wrong response type
// Keep journal_entries as None
Err(_) => {
// No reply: same quiet stream refresh
// as the details endpoint above.
self.quiet_reconnect().await;
}
Ok(Err(_)) | Ok(Ok(_)) => {
// Endpoint unsupported or wrong type;
// keep journal_entries as None
}
}
self.last_journal_poll = Instant::now();
@@ -1158,16 +1308,23 @@ impl App {
}
}
}
Err(e) => {
Ok(Err(e)) => {
// Connection error - show modal if not already shown
let error_message = format!("Failed to fetch metrics: {e}");
self.show_connection_error(error_message);
}
_ => {
Ok(_) => {
// Unexpected response type
self.show_connection_error("Unexpected response from agent".to_string());
}
}
}
// A poisoned connection may have been replaced mid-iteration:
// restart on the fresh stream before issuing any more requests.
if self.replacement_connection.is_some() {
return Ok(());
}
// Update countdown for connection error modal if active
if self.modal_manager.is_active() {
@@ -1178,8 +1335,27 @@ impl App {
// Draw
terminal.draw(|f| self.draw(f))?;
// Tick rate
sleep(self.metrics_interval).await;
// Tick wait, kept responsive: instead of sleeping the whole
// metrics interval (which queued keys/wheel events for up to
// 500ms and applied them in bursts), wait in ≤33ms slices and
// handle + repaint input the moment it arrives.
let deadline = Instant::now() + self.metrics_interval;
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() || self.should_quit {
break;
}
if !event::poll(remaining.min(Duration::from_millis(33)))? {
continue;
}
if matches!(
self.drain_input(terminal).await?,
InputFlow::RestartConnection
) {
return Ok(());
}
terminal.draw(|f| self.draw(f))?;
}
}
Ok(())
@@ -1249,19 +1425,39 @@ impl App {
self.per_core_hist.ensure_cores(m.cpu_per_core.len());
self.per_core_hist.push_samples(&m.cpu_per_core);
// NET: sum across all ifaces, compute KB/s via elapsed time
// NET: sum across all ifaces, compute KB/s. Prefer the agent's sample
// timestamps (the agent serves TTL-cached snapshots, so client receive
// time overstates dt on a cache hit and produces a 0-then-2x sawtooth);
// fall back to the client clock against pre-1.60 agents.
let now = Instant::now();
let rx_total = m.networks.iter().map(|n| n.received).sum::<u64>();
let tx_total = m.networks.iter().map(|n| n.transmitted).sum::<u64>();
let (rx_kb, tx_kb) = if let Some((prx, ptx, pts)) = self.last_net_totals {
let dt = now.duration_since(pts).as_secs_f64().max(1e-6);
// None = identical agent snapshot (cache hit): repeat the previous
// rates so the timeline advances without a fake dip to zero.
let dt = match (m.sampled_at_ms, self.last_net_sampled_at_ms) {
(Some(a), Some(b)) if a == b => None,
(Some(a), Some(b)) if a > b => Some((a - b) as f64 / 1000.0),
// Agent restarted or clock stepped backwards: client clock.
_ => Some(now.duration_since(pts).as_secs_f64().max(1e-6)),
};
match dt {
None => (
self.rx_hist.back().copied().unwrap_or(0),
self.tx_hist.back().copied().unwrap_or(0),
),
Some(dt) => {
let dt = dt.max(1e-6);
let rx = ((rx_total.saturating_sub(prx)) as f64 / dt / 1024.0).round() as u64;
let tx = ((tx_total.saturating_sub(ptx)) as f64 / dt / 1024.0).round() as u64;
(rx, tx)
}
}
} else {
(0, 0)
};
self.last_net_totals = Some((rx_total, tx_total, now));
self.last_net_sampled_at_ms = m.sampled_at_ms;
push_capped(&mut self.rx_hist, rx_kb, 600);
push_capped(&mut self.tx_hist, tx_kb, 600);
self.rx_peak = self.rx_peak.max(rx_kb);
File diff suppressed because it is too large Load Diff
+1
View File
@@ -589,6 +589,7 @@ mod render_tests {
fn fake_metrics(cores: Vec<f32>) -> Metrics {
Metrics {
sampled_at_ms: None,
cpu_total: 0.0,
cpu_per_core: cores,
mem_total: 1024,
+3 -2
View File
@@ -1,7 +1,8 @@
//! Disk cards with per-device gauge and title line.
use crate::types::Metrics;
use crate::ui::util::{disk_icon, human, truncate_middle};
use crate::ui::fit::truncate_middle_cols;
use crate::ui::util::{disk_icon, human};
use ratatui::{
layout::{Constraint, Direction, Layout, Rect},
style::Style,
@@ -69,7 +70,7 @@ pub fn draw_disks(f: &mut ratatui::Frame<'_>, area: Rect, m: Option<&Metrics>) {
"{}{}{}{} {} / {} ({}%)",
indent,
disk_icon(&d.name),
truncate_middle(&d.name, (slot.width.saturating_sub(6)) as usize / 2),
truncate_middle_cols(&d.name, slot.width.saturating_sub(6) / 2),
temp_str,
human(used),
human(d.total),
+57
View File
@@ -43,6 +43,46 @@ pub fn truncate_cols(s: &str, max: u16) -> String {
out
}
/// Shortens `s` to at most `max` columns by cutting the MIDDLE, marking the
/// cut with `…` — device names like `/dev/nvme0n1p1` keep their distinctive
/// prefix and suffix. Column- and char-boundary-safe; the byte-slicing
/// predecessor in `util.rs` panicked on non-ASCII names.
pub fn truncate_middle_cols(s: &str, max: u16) -> String {
if cols(s) <= max {
return s.to_string();
}
if max <= 1 {
return truncate_cols(s, max);
}
// Reserve one column for the ellipsis; split the rest left/right.
let left_budget = (max - 1) / 2;
let right_budget = max - 1 - left_budget;
let mut left_end = 0; // byte index
let mut used = 0u16;
for (i, ch) in s.char_indices() {
let w = cols(ch.encode_utf8(&mut [0u8; 4]));
if used + w > left_budget {
break;
}
used += w;
left_end = i + ch.len_utf8();
}
let mut right_start = s.len();
let mut used = 0u16;
for (i, ch) in s.char_indices().rev() {
let w = cols(ch.encode_utf8(&mut [0u8; 4]));
if used + w > right_budget || i < left_end {
break;
}
used += w;
right_start = i;
}
format!("{}{}", &s[..left_end], &s[right_start..])
}
/// Picks the first (richest) candidate pair that fits side by side in `width` columns
/// with at least `gap` columns between them.
///
@@ -108,6 +148,23 @@ mod tests {
assert_eq!(truncate_cols("🔒ab", 2), "");
}
/// Middle truncation keeps both ends — the parts that identify a device —
/// and must never exceed the budget or split a character.
#[test]
fn truncate_middle_keeps_both_ends_within_budget() {
assert_eq!(truncate_middle_cols("/dev/nvme0n1p1", 20), "/dev/nvme0n1p1");
let out = truncate_middle_cols("/dev/nvme0n1p1", 9);
assert_eq!(cols(&out), 9);
assert!(out.starts_with("/dev"), "{out}");
assert!(out.ends_with("1p1"), "{out}");
assert!(out.contains('…'), "{out}");
// Non-ASCII names must not panic (the old byte-slicing version did).
for max in 0..12u16 {
let out = truncate_middle_cols("диск-🗄️-данные", max);
assert!(cols(&out) <= max.max(1), "{out:?} exceeds {max}");
}
}
#[test]
fn pick_pair_takes_the_richest_that_fits() {
let candidates = [
+1
View File
@@ -249,6 +249,7 @@ mod render_tests {
fn metrics(gpus: Option<Vec<GpuInfo>>) -> Metrics {
Metrics {
sampled_at_ms: None,
cpu_total: 0.0,
cpu_per_core: vec![],
mem_total: 1024,
+46 -17
View File
@@ -472,13 +472,28 @@ impl ModalManager {
.borders(Borders::ALL);
let content_lines: Vec<Line> = if journal.entries.is_empty() {
vec![
let mut lines = vec![
Line::from(""),
Line::from(Span::styled(
"No journal entries found for this process",
Style::default().add_modifier(Modifier::DIM),
)),
]
];
// Access limits, not absence of logs: show journalctl's own hint
// (typical when the agent runs as an unprivileged user, e.g. demo
// mode) plus the practical fix.
if let Some(notice) = &journal.notice {
lines.push(Line::from(""));
lines.push(Line::from(Span::styled(
format!("{notice}"),
Style::default().fg(Color::Yellow),
)));
lines.push(Line::from(Span::styled(
" Run the agent as a service (or a user in the systemd-journal group) for full journal access.",
Style::default().add_modifier(Modifier::DIM),
)));
}
lines
} else {
journal
.entries
@@ -624,17 +639,30 @@ impl ModalManager {
// labels + axis title + (top) Y-axis title + legend + spacing.
let mut lines: Vec<Line> = Vec::with_capacity(plot_height + 6);
// Y-axis labels and plot content
let mut row_buf = String::with_capacity(plot_width);
for y in 0..plot_height {
let y_value = params.max_system * (1.0 - (y as f64 / (plot_height - 1).max(1) as f64));
// 4-char fixed-width label so the axis doesn't shift as digits change.
let y_label = if y_value >= 100.0 {
format!("{y_value:>4.0}")
// Format a CPU-time value: whole ms once past 100, one decimal below.
let fmt_ms = |v: f64| {
if v >= 100.0 {
format!("{v:.0}")
} else {
format!("{y_value:>4.1}")
format!("{v:.1}")
}
};
// Y-axis labels, right-aligned to the widest value this frame so the
// axis stays a straight line. The old fixed 4-char field predates the
// CPU-time unit fix; honest millisecond values (e.g. 136114) blew
// through it and skewed the whole axis.
let y_values: Vec<String> = (0..plot_height)
.map(|y| {
fmt_ms(params.max_system * (1.0 - (y as f64 / (plot_height - 1).max(1) as f64)))
})
.collect();
let y_label_w = y_values.iter().map(|s| s.len()).max().unwrap_or(4).max(4);
let mut row_buf = String::with_capacity(plot_width);
for (y, y_value) in y_values.iter().enumerate() {
let y_label = format!("{y_value:>y_label_w$}");
// Build the row's char slice into a reusable String buffer.
row_buf.clear();
let start = y * plot_width;
@@ -650,8 +678,8 @@ impl ModalManager {
]));
}
// Add X-axis
let x_axis_padding = " ".to_string(); // Match Y-axis label width
// Add X-axis (padding = Y label width + the space before the bar)
let x_axis_padding = " ".repeat(y_label_w + 1);
let x_axis_line = "".repeat(plot_width + 1);
lines.push(Line::from(vec![
Span::styled(x_axis_padding, Style::default()),
@@ -659,13 +687,14 @@ impl ModalManager {
]));
// Add X-axis labels
let x_label_start = "0.0".to_string();
let x_label_mid = format!("{:.1}", params.max_user / 2.0);
let x_label_end = format!("{:.1}", params.max_user);
let x_label_start = fmt_ms(0.0);
let x_label_mid = fmt_ms(params.max_user / 2.0);
let x_label_end = fmt_ms(params.max_user);
let spacing = plot_width / 3;
let x_labels = format!(
" {}{}{}{}{}",
"{}{}{}{}{}{}",
" ".repeat(y_label_w + 1),
x_label_start,
" ".repeat(spacing.saturating_sub(x_label_start.len())),
x_label_mid,
@@ -677,7 +706,7 @@ impl ModalManager {
// Add axis titles with better visibility
lines.push(Line::from(vec![Span::styled(
" User CPU Time (ms) →",
format!("{}User CPU Time (ms) →", " ".repeat(y_label_w + 1)),
Style::default()
.fg(Color::Yellow)
.add_modifier(Modifier::BOLD),
+17 -95
View File
@@ -499,7 +499,6 @@ pub fn draw_top_processes(f: &mut ratatui::Frame<'_>, area: Rect, params: Proces
}
}
/// Handle keyboard scrolling (Up/Down/PageUp/PageDown/Home/End)
/// Parameters for process key event handling
pub struct ProcessKeyParams<'a> {
pub selected_process_pid: &'a mut Option<u32>,
@@ -509,16 +508,6 @@ pub struct ProcessKeyParams<'a> {
pub filtered_indices: &'a [usize],
}
/// LEGACY: Use processes_handle_key_with_selection for enhanced functionality
#[allow(dead_code)]
pub fn processes_handle_key(
scroll_offset: &mut usize,
key: crossterm::event::KeyEvent,
page_size: usize,
) {
crate::ui::cpu::per_core_handle_key(scroll_offset, key, page_size);
}
pub fn processes_handle_key_with_selection(params: ProcessKeyParams) -> bool {
use crossterm::event::KeyCode;
@@ -598,83 +587,6 @@ pub fn processes_handle_key_with_selection(params: ProcessKeyParams) -> bool {
}
}
/// Handle mouse for content scrolling and scrollbar dragging.
/// Returns Some(new_sort) if the header "CPU %" or "Mem" was clicked.
/// LEGACY: Use processes_handle_mouse_with_selection for enhanced functionality
#[allow(dead_code)]
pub fn processes_handle_mouse(
scroll_offset: &mut usize,
drag: &mut Option<crate::ui::cpu::PerCoreScrollDrag>,
mouse: MouseEvent,
area: Rect,
total_rows: usize,
) -> Option<ProcSortBy> {
// Inner and content areas (match draw_top_processes)
let inner = Rect {
x: area.x + 1,
y: area.y + 1,
width: area.width.saturating_sub(2),
height: area.height.saturating_sub(2),
};
if inner.height == 0 || inner.width <= 2 {
return None;
}
let content = Rect {
x: inner.x,
y: inner.y,
width: inner.width.saturating_sub(2),
height: inner.height,
};
// Scrollbar interactions (click arrows/page/drag)
per_core_handle_scrollbar_mouse(scroll_offset, drag, mouse, area, total_rows);
// Wheel scrolling when inside the content
crate::ui::cpu::per_core_handle_mouse(scroll_offset, mouse, content, content.height as usize);
// Header click to change sort
let header_area = Rect {
x: content.x,
y: content.y,
width: content.width,
height: 1,
};
let inside_header = mouse.row == header_area.y
&& mouse.column >= header_area.x
&& mouse.column < header_area.x + header_area.width;
if inside_header && matches!(mouse.kind, MouseEventKind::Down(MouseButton::Left)) {
// Split the header the same way the draw path did, so a click lands on the
// column actually on screen even when PID has been dropped.
let columns = ProcColumns::for_width(header_area.width);
let cols = Layout::default()
.direction(Direction::Horizontal)
.constraints(columns.constraints())
.spacing(COL_SPACING) // must match Table::column_spacing in the draw path
.split(header_area);
if let Some(cpu) = columns.cpu_index().map(|i| cols[i])
&& mouse.column >= cpu.x
&& mouse.column < cpu.x + cpu.width
{
return Some(ProcSortBy::CpuDesc);
}
if let Some(mem) = columns.mem_index().map(|i| cols[i])
&& mouse.column >= mem.x
&& mouse.column < mem.x + mem.width
{
return Some(ProcSortBy::MemDesc);
}
}
// Clamp to valid range
per_core_clamp(
scroll_offset,
total_rows,
(content.height.saturating_sub(1)) as usize,
);
None
}
/// Parameters for process mouse event handling
pub struct ProcessMouseParams<'a> {
pub scroll_offset: &'a mut usize,
@@ -948,6 +860,7 @@ mod click_tests {
fn metrics() -> Metrics {
Metrics {
sampled_at_ms: None,
cpu_total: 0.0,
cpu_per_core: vec![],
mem_total: 32_000_000_000,
@@ -1003,20 +916,29 @@ mod click_tests {
}
fn click(width: u16, column: u16) -> Option<ProcSortBy> {
let m = metrics();
let mut scroll = 0usize;
let mut drag = None;
processes_handle_mouse(
&mut scroll,
&mut drag,
MouseEvent {
let mut sel_pid = None;
let mut sel_idx = None;
let idxs = [0usize];
processes_handle_mouse_with_selection(ProcessMouseParams {
scroll_offset: &mut scroll,
selected_process_pid: &mut sel_pid,
selected_process_index: &mut sel_idx,
drag: &mut drag,
mouse: MouseEvent {
kind: MouseEventKind::Down(MouseButton::Left),
column,
row: 1,
modifiers: KeyModifiers::NONE,
},
Rect::new(0, 0, width, 8),
1,
)
area: Rect::new(0, 0, width, 8),
total_rows: 1,
metrics: Some(&m),
search_box_visible: false,
filtered_indices: &idxs,
})
}
/// The hit-test rects are computed by a separate `Layout` call from the one `Table`
-13
View File
@@ -22,19 +22,6 @@ pub fn human(b: u64) -> String {
format!("{tb:.2}TB")
}
pub fn truncate_middle(s: &str, max: usize) -> String {
if s.len() <= max {
return s.to_string();
}
if max <= 3 {
return "...".into();
}
let keep = max - 3;
let left = keep / 2;
let right = keep - left;
format!("{}...{}", &s[..left], &s[s.len() - right..])
}
pub fn disk_icon(name: &str) -> &'static str {
let n = name.to_ascii_lowercase();
if n.contains(':') {
View File
+1
View File
@@ -8,6 +8,7 @@ static ENV_LOCK: Mutex<()> = Mutex::new(());
#[allow(dead_code)] // touch crate
fn touch() {
let _ = socktop::types::Metrics {
sampled_at_ms: None,
cpu_total: 0.0,
cpu_per_core: vec![],
mem_total: 0,
+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()));
}
}
}
+260 -260
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)
// 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 cached_gpus().is_some() {
cached_gpus()
} else if let Some(cached) = GPUS.get_fresh(TTL) {
cached
} 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.
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)
{
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);
.store(v.is_some(), std::sync::atomic::Ordering::Release);
}
}
set_gpus(v.clone());
GPUS.set(v.clone());
v
}
} else {
None
};
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);
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "socktop_connector"
version = "1.50.0"
version = "1.60.0"
edition = "2024"
license = "MIT"
description = "WebSocket connector library for socktop agent communication"
+7 -3
View File
@@ -1,8 +1,12 @@
fn main() -> Result<(), Box<dyn std::error::Error>> {
// Set the protoc binary path to use the vendored version for CI compatibility
// SAFETY: We're only setting PROTOC in a build script environment, which is safe
// Vendored protoc for reproducible builds where available. It ships no
// riscv64 binary, so on such hosts leave $PROTOC / PATH lookup to
// prost-build (apt: protobuf-compiler).
// SAFETY: We're only setting PROTOC in a build script environment.
if let Ok(protoc) = protoc_bin_vendored::protoc_bin_path() {
unsafe {
std::env::set_var("PROTOC", protoc_bin_vendored::protoc_bin_path()?);
std::env::set_var("PROTOC", protoc);
}
}
prost_build::compile_protos(&["processes.proto"], &["."])?;
File diff suppressed because it is too large Load Diff
+142 -34
View File
@@ -6,7 +6,7 @@ use crate::error::{ConnectorError, Result};
use std::io::BufReader;
use std::sync::Arc;
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async};
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream};
use url::Url;
#[cfg(feature = "tls")]
@@ -15,7 +15,7 @@ use {
rustls::{
DigitallySignedStruct, RootCertStore, SignatureScheme,
client::danger::{HandshakeSignatureValid, ServerCertVerified, ServerCertVerifier},
crypto::ring,
crypto::{WebPkiSupportedAlgorithms, ring},
pki_types::{CertificateDer, ServerName, UnixTime},
},
rustls_pemfile::Item,
@@ -64,7 +64,8 @@ async fn connect_without_ca_and_config(url: &str, config: &ConnectorConfig) -> R
);
}
let (ws, _) = connect_async(req).await?;
// `true` disables Nagle: small request/response frames, latency matters.
let (ws, _) = tokio_tungstenite::connect_async_with_config(req, None, true).await?;
Ok(ws)
}
@@ -85,7 +86,12 @@ async fn connect_with_ca_and_config(
der_certs.push(der);
}
}
root.add_parsable_certificates(der_certs);
if der_certs.is_empty() {
return Err(ConnectorError::protocol_error(format!(
"no certificates found in --tls-ca file: {ca_path}"
)));
}
root.add_parsable_certificates(der_certs.iter().cloned());
let mut cfg = ClientConfig::builder()
.with_root_certificates(root)
@@ -114,55 +120,87 @@ async fn connect_with_ca_and_config(
}
if !config.verify_hostname {
#[derive(Debug)]
struct NoVerify;
impl ServerCertVerifier for NoVerify {
// Default mode: certificate PINNING without hostname verification.
// The server must present a certificate byte-identical to one in the
// --tls-ca file. This intentionally ignores expiry and chain building
// (the operator pinned this exact cert), but unlike a blanket accept
// it makes MITM certs fail the handshake.
cfg.dangerous()
.set_certificate_verifier(Arc::new(PinnedCertVerifier::new(der_certs)));
}
let cfg = Arc::new(cfg);
// Third argument is tungstenite's `disable_nagle`: always true — socktop
// exchanges small request/response frames where Nagle only adds latency.
let (ws, _) = tokio_tungstenite::connect_async_tls_with_config(
req,
None,
true,
Some(Connector::Rustls(cfg)),
)
.await?;
Ok(ws)
}
/// Accepts exactly the certificates the user pinned via `--tls-ca`, nothing else.
///
/// Used when hostname verification is off (the default for self-signed
/// home-lab certs). Signature validation still runs with the ring provider's
/// full algorithm set; only the certificate identity check is replaced —
/// by an exact DER comparison against the pinned certificate(s).
#[cfg(feature = "tls")]
#[derive(Debug)]
struct PinnedCertVerifier {
pinned: Vec<CertificateDer<'static>>,
algorithms: WebPkiSupportedAlgorithms,
}
#[cfg(feature = "tls")]
impl PinnedCertVerifier {
fn new(pinned: Vec<CertificateDer<'static>>) -> Self {
Self {
pinned,
algorithms: ring::default_provider().signature_verification_algorithms,
}
}
}
#[cfg(feature = "tls")]
impl ServerCertVerifier for PinnedCertVerifier {
fn verify_server_cert(
&self,
_end_entity: &CertificateDer<'_>,
end_entity: &CertificateDer<'_>,
_intermediates: &[CertificateDer<'_>],
_server_name: &ServerName,
_ocsp_response: &[u8],
_now: UnixTime,
) -> std::result::Result<ServerCertVerified, rustls::Error> {
if self.pinned.iter().any(|p| p == end_entity) {
Ok(ServerCertVerified::assertion())
} else {
Err(rustls::Error::InvalidCertificate(
rustls::CertificateError::ApplicationVerificationFailure,
))
}
}
fn verify_tls12_signature(
&self,
_message: &[u8],
_cert: &CertificateDer<'_>,
_dss: &DigitallySignedStruct,
message: &[u8],
cert: &CertificateDer<'_>,
dss: &DigitallySignedStruct,
) -> std::result::Result<HandshakeSignatureValid, rustls::Error> {
Ok(HandshakeSignatureValid::assertion())
rustls::crypto::verify_tls12_signature(message, cert, dss, &self.algorithms)
}
fn verify_tls13_signature(
&self,
_message: &[u8],
_cert: &CertificateDer<'_>,
_dss: &DigitallySignedStruct,
message: &[u8],
cert: &CertificateDer<'_>,
dss: &DigitallySignedStruct,
) -> std::result::Result<HandshakeSignatureValid, rustls::Error> {
Ok(HandshakeSignatureValid::assertion())
rustls::crypto::verify_tls13_signature(message, cert, dss, &self.algorithms)
}
fn supported_verify_schemes(&self) -> Vec<SignatureScheme> {
vec![
SignatureScheme::ECDSA_NISTP256_SHA256,
SignatureScheme::ED25519,
SignatureScheme::RSA_PSS_SHA256,
]
self.algorithms.supported_schemes()
}
}
cfg.dangerous().set_certificate_verifier(Arc::new(NoVerify));
// Note: hostname verification disabled (default). Set SOCKTOP_VERIFY_NAME=1 to enable strict SAN checking.
}
let cfg = Arc::new(cfg);
let (ws, _) = tokio_tungstenite::connect_async_tls_with_config(
req,
None,
config.verify_hostname,
Some(Connector::Rustls(cfg)),
)
.await?;
Ok(ws)
}
#[cfg(not(feature = "tls"))]
@@ -181,3 +219,73 @@ async fn connect_with_ca_and_config(
fn ensure_crypto_provider() {
let _ = ring::default_provider().install_default();
}
#[cfg(all(test, feature = "tls"))]
mod tests {
use super::*;
fn verifier(pinned: &[&[u8]]) -> PinnedCertVerifier {
let _ = ring::default_provider().install_default();
PinnedCertVerifier::new(
pinned
.iter()
.map(|b| CertificateDer::from(b.to_vec()))
.collect(),
)
}
fn verify(v: &PinnedCertVerifier, presented: &[u8]) -> bool {
v.verify_server_cert(
&CertificateDer::from(presented.to_vec()),
&[],
&ServerName::try_from("agent.test").unwrap(),
&[],
UnixTime::now(),
)
.is_ok()
}
/// The regression this verifier exists to prevent: the old NoVerify
/// accepted ANY certificate when hostname verification was off, so the
/// documented pinning was a no-op. The pinned cert must be accepted and
/// every other cert rejected.
#[test]
fn only_the_pinned_certificate_is_accepted() {
let v = verifier(&[b"pinned-cert-der"]);
assert!(verify(&v, b"pinned-cert-der"));
assert!(!verify(&v, b"some-mitm-cert"), "unpinned cert accepted");
assert!(!verify(&v, b""), "empty cert accepted");
}
/// A --tls-ca file may hold several certs (e.g. during rotation); any of
/// them must satisfy the pin.
#[test]
fn any_cert_in_a_multi_cert_pem_satisfies_the_pin() {
let v = verifier(&[b"old-cert", b"new-cert"]);
assert!(verify(&v, b"old-cert"));
assert!(verify(&v, b"new-cert"));
assert!(!verify(&v, b"third-party-cert"));
}
/// Fail closed: an empty pin set must reject everything rather than
/// falling back to accept-all.
#[test]
fn an_empty_pin_set_rejects_all_certificates() {
let v = verifier(&[]);
assert!(!verify(&v, b"anything"));
}
/// Signature schemes come from the real provider, not a hardcoded list —
/// an agent using e.g. RSA-PKCS1 must still be able to handshake.
#[test]
fn signature_schemes_come_from_the_provider() {
let v = verifier(&[b"x"]);
let schemes = v.supported_verify_schemes();
assert!(
schemes.len() > 3,
"suspiciously short scheme list: {schemes:?}"
);
assert!(schemes.contains(&SignatureScheme::RSA_PKCS1_SHA256));
assert!(schemes.contains(&SignatureScheme::ECDSA_NISTP256_SHA256));
}
}
+8
View File
@@ -54,6 +54,10 @@ pub struct GpuInfo {
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Metrics {
/// Epoch ms when the agent actually collected this snapshot (agents may
/// serve TTL-cached data). Absent on agents older than 1.60.
#[serde(default)]
pub sampled_at_ms: Option<u64>,
pub cpu_total: f32,
pub cpu_per_core: Vec<f32>,
pub mem_total: u64,
@@ -147,6 +151,10 @@ pub struct JournalResponse {
pub entries: Vec<JournalEntry>,
pub total_count: u32,
pub truncated: bool,
/// Agent-side explanation for an empty result (journal access limits).
/// Absent on agents older than 1.60.
#[serde(default)]
pub notice: Option<String>,
pub cached_at: u64, // Unix timestamp when this data was cached
}
+1
View File
@@ -46,6 +46,7 @@ pub async fn send_request_and_wait(
// For now, return a placeholder metrics response indicating binary data received
// TODO: Implement proper protobuf decoding for binary data
let placeholder_metrics = Metrics {
sampled_at_ms: None,
cpu_total: 0.0,
cpu_per_core: vec![0.0],
mem_total: 0,
+1 -3
View File
@@ -475,9 +475,7 @@ dependencies = [
[[package]]
name = "socktop_connector"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a63dadaa5105df11b0684759a829012257d48e72a469cc554c0cf4394605f5a"
version = "1.51.0"
dependencies = [
"flate2",
"js-sys",
+2 -2
View File
@@ -10,8 +10,8 @@ edition = "2021"
crate-type = ["cdylib"]
[dependencies]
# Use WASM features for WebSocket connectivity (published version)
socktop_connector = { version = "0.1.5", default-features = false, features = ["wasm"] }
# Use WASM features for WebSocket connectivity (in-repo connector via path)
socktop_connector = { path = "../socktop_connector", default-features = false, features = ["wasm"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
wasm-bindgen = "0.2"
View File
+4384
View File
File diff suppressed because it is too large Load Diff
+4 -1
View File
@@ -3,6 +3,9 @@ name = "zellij_socktop_plugin"
version = "0.1.0"
edition = "2021"
# Standalone package, not part of the parent workspace (same as socktop_wasm_test)
[workspace]
[lib]
crate-type = ["cdylib"]
@@ -10,7 +13,7 @@ crate-type = ["cdylib"]
zellij-tile = "0.40.0"
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
socktop_connector = { version = "0.1.5", default-features = false, features = ["wasm"] }
socktop_connector = { path = "../socktop_connector", default-features = false, features = ["wasm"] }
futures = "0.3"
[dependencies.chrono]