From b7a2a0fe464e29e179963adcf816018410c8cad8 Mon Sep 17 00:00:00 2001 From: not-matthias Date: Wed, 23 Sep 2026 18:38:52 +0200 Subject: [PATCH 1/5] feat(memtrack): log ring buffer throughput behind CODSPEED_MEMTRACK_RING_STATS Each ring poller can now report how fast its ring is written and read, sampled from libbpf's mmapped producer/consumer positions (ring__producer_pos, ring__consumer_pos, ring__avail_data_size). This needs no BPF changes and no syscalls: a sample is a few memory loads per poll tick. With CODSPEED_MEMTRACK_RING_STATS=1, every poll thread logs a line per second at debug level and a whole-run summary at info level on shutdown: stacks ring (1s): wrote 818.7 MB/s, read 664.5 MB/s, drain 61398.1 MB/s, peak backlog 159.7 MiB (31.2%), busy 1.1% (max tick 1.02 ms) - wrote: producer position delta over wall time. It counts every reserved byte, including record headers and records the BPF side discards, so it measures ring pressure rather than artifact size. - read: consumer position delta over wall time. - drain: bytes consumed per second of poll-thread busy time, i.e. the read rate the poller can sustain. - peak backlog: largest unconsumed backlog seen at the start of a poll tick. - busy / max tick: share of wall time spent consuming, and the slowest tick. When the variable is unset, the poll loop does one Option check per tick. --- crates/memtrack/AGENTS.md | 2 +- crates/memtrack/src/ebpf/mod.rs | 1 + crates/memtrack/src/ebpf/poller.rs | 43 +++++--- crates/memtrack/src/ebpf/ring_stats.rs | 132 +++++++++++++++++++++++++ 4 files changed, 164 insertions(+), 14 deletions(-) create mode 100644 crates/memtrack/src/ebpf/ring_stats.rs diff --git a/crates/memtrack/AGENTS.md b/crates/memtrack/AGENTS.md index 6c7b09d37..f8f1a1ace 100644 --- a/crates/memtrack/AGENTS.md +++ b/crates/memtrack/AGENTS.md @@ -84,7 +84,7 @@ sudo -E cargo test --test c_tests -- --test-threads 1 - **Build toolchain:** `clang` + BTF/vmlinux headers, `libbpf-dev`, `zlib1g-dev`, `pkgconf`, `build-essential`; vendored libbpf also needs `autopoint`/`bison`/`flex`. - `vmlinux.h` is pinned to a specific git rev; `libbpf-rs` uses the `vendored` feature (dist links `libbpf-rs/static`). -Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). +Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_MEMTRACK_RING_STATS` (1 logs per-ring write/read/drain rates and backlog: every second at debug, a run summary at info), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). ### Minimum kernel version diff --git a/crates/memtrack/src/ebpf/mod.rs b/crates/memtrack/src/ebpf/mod.rs index 3aba823ac..4e7cf730b 100644 --- a/crates/memtrack/src/ebpf/mod.rs +++ b/crates/memtrack/src/ebpf/mod.rs @@ -4,6 +4,7 @@ mod memtrack; pub(crate) mod pause; pub(crate) mod poller; mod proc_fs; +mod ring_stats; mod spawn; mod stacks; mod tracker; diff --git a/crates/memtrack/src/ebpf/poller.rs b/crates/memtrack/src/ebpf/poller.rs index 637ed27e7..d34f86fe1 100644 --- a/crates/memtrack/src/ebpf/poller.rs +++ b/crates/memtrack/src/ebpf/poller.rs @@ -1,3 +1,4 @@ +use crate::ebpf::ring_stats::RingStats; use anyhow::{Context, Result}; use libbpf_rs::{AsRawLibbpf, MapCore, RingBuffer, RingBufferBuilder, libbpf_sys}; use parking_lot::Mutex; @@ -69,6 +70,11 @@ fn poll_iteration( } } +fn ring_of(ringbuf: &RingBuffer) -> *mut libbpf_sys::ring { + // SAFETY: a built `RingBuffer` holds exactly the one ring added in `new`. + unsafe { libbpf_sys::ring_buffer__ring(ringbuf.as_libbpf_object().as_ptr(), 0) } +} + /// Polls a BPF ring buffer in a background thread, parsing raw entries with a /// user-supplied closure and forwarding them to an mpsc channel in batches. /// @@ -114,26 +120,34 @@ impl RingBufferPoller { 0 })?; let ringbuf = builder.build()?; + let name = rb_map.name().to_string_lossy().into_owned(); // The control channel doubles as the poll pacing: a received message is // a drain request (acked after a full consume), a timeout is a regular // poll tick, and disconnection is the shutdown signal. let (ctl, ctl_rx) = mpsc::channel::>(); let poll_thread = std::thread::spawn(move || { - // SAFETY: the built `RingBuffer` holds exactly the one ring added above. - let ring = - unsafe { libbpf_sys::ring_buffer__ring(ringbuf.as_libbpf_object().as_ptr(), 0) }; - while poll_iteration( - ctl_rx.recv_timeout(Duration::from_millis(poll_interval_ms)), - || consume_all(&ringbuf, ring), - || { - let _ = ringbuf.poll(Duration::ZERO); - }, - &batch, - &tx, - ) { + let mut stats = RingStats::enabled().then(|| RingStats::new(name, ring_of(&ringbuf))); + loop { + let control = ctl_rx.recv_timeout(Duration::from_millis(poll_interval_ms)); + let tick = stats.as_ref().map(RingStats::begin); + let running = poll_iteration( + control, + || consume_all(&ringbuf, ring_of(&ringbuf)), + || { + let _ = ringbuf.poll(Duration::ZERO); + }, + &batch, + &tx, + ); + if let (Some(stats), Some(tick)) = (&mut stats, tick) { + stats.end(tick); + } + if !running { + break; + } if let Some(on_drained) = &on_drained - && unsafe { libbpf_sys::ring__avail_data_size(ring) } == 0 + && unsafe { libbpf_sys::ring__avail_data_size(ring_of(&ringbuf)) } == 0 { on_drained(); } @@ -141,6 +155,9 @@ impl RingBufferPoller { if let Some(on_drained) = &on_drained { on_drained(); } + if let Some(stats) = &stats { + stats.report_run(); + } }); Ok(Self { diff --git a/crates/memtrack/src/ebpf/ring_stats.rs b/crates/memtrack/src/ebpf/ring_stats.rs new file mode 100644 index 000000000..a97e7ae69 --- /dev/null +++ b/crates/memtrack/src/ebpf/ring_stats.rs @@ -0,0 +1,132 @@ +//! Ring throughput from libbpf's mmapped producer/consumer positions, enabled +//! with `CODSPEED_MEMTRACK_RING_STATS=1`. +//! +//! The producer position counts every reserved byte, including the 8-byte +//! record headers and records the BPF side later discarded, so "wrote" is ring +//! pressure rather than artifact bytes. + +use crate::prelude::*; +use libbpf_rs::libbpf_sys; +use std::time::{Duration, Instant}; + +const REPORT_INTERVAL: Duration = Duration::from_secs(1); + +pub(crate) struct RingStats { + name: String, + ring: *const libbpf_sys::ring, + run: Window, + window: Window, +} + +/// Taken right before a poll iteration: the backlog then is everything the +/// producers wrote since the previous iteration. +pub(crate) struct Tick { + started: Instant, + backlog: u64, +} + +struct Window { + started: Instant, + produced: u64, + consumed: u64, + busy: Duration, + max_tick: Duration, + peak_backlog: u64, +} + +impl RingStats { + pub(crate) fn enabled() -> bool { + std::env::var("CODSPEED_MEMTRACK_RING_STATS").is_ok_and(|v| v == "1") + } + + /// `ring` must stay valid for as long as the stats are used. + pub(crate) fn new(name: String, ring: *const libbpf_sys::ring) -> Self { + Self { + name, + ring, + run: Window::start(ring), + window: Window::start(ring), + } + } + + pub(crate) fn begin(&self) -> Tick { + Tick { + started: Instant::now(), + // SAFETY: `ring` is valid per `new`'s contract. + backlog: unsafe { libbpf_sys::ring__avail_data_size(self.ring) } as u64, + } + } + + pub(crate) fn end(&mut self, tick: Tick) { + let busy = tick.started.elapsed(); + self.run.record(tick.backlog, busy); + self.window.record(tick.backlog, busy); + if self.window.started.elapsed() < REPORT_INTERVAL { + return; + } + debug!( + "{} ring (1s): {}", + self.name, + self.window.summary(self.ring) + ); + self.window = Window::start(self.ring); + } + + pub(crate) fn report_run(&self) { + info!("{} ring (run): {}", self.name, self.run.summary(self.ring)); + } +} + +impl Window { + fn start(ring: *const libbpf_sys::ring) -> Self { + let (produced, consumed) = positions(ring); + Self { + started: Instant::now(), + produced, + consumed, + busy: Duration::ZERO, + max_tick: Duration::ZERO, + peak_backlog: 0, + } + } + + fn record(&mut self, backlog: u64, busy: Duration) { + self.busy += busy; + self.max_tick = self.max_tick.max(busy); + self.peak_backlog = self.peak_backlog.max(backlog); + } + + /// `drain` is the consume rate while the poll thread is busy, i.e. the + /// read speed the thread can sustain; `read` is averaged over wall time. + fn summary(&self, ring: *const libbpf_sys::ring) -> String { + let (produced, consumed) = positions(ring); + let wall = self.started.elapsed(); + let consumed = consumed - self.consumed; + // SAFETY: `ring` is valid per `RingStats::new`'s contract. + let size = unsafe { libbpf_sys::ring__size(ring) } as u64; + format!( + "wrote {:.1} MB/s, read {:.1} MB/s, drain {:.1} MB/s, peak backlog {:.1} MiB ({:.1}%), busy {:.1}% (max tick {:.2} ms)", + mb_per_s(produced - self.produced, wall), + mb_per_s(consumed, wall), + mb_per_s(consumed, self.busy), + self.peak_backlog as f64 / (1024.0 * 1024.0), + self.peak_backlog as f64 * 100.0 / size as f64, + self.busy.as_secs_f64() * 100.0 / wall.as_secs_f64(), + self.max_tick.as_secs_f64() * 1e3, + ) + } +} + +fn positions(ring: *const libbpf_sys::ring) -> (u64, u64) { + // SAFETY: `ring` is valid per `RingStats::new`'s contract. + unsafe { + ( + libbpf_sys::ring__producer_pos(ring), + libbpf_sys::ring__consumer_pos(ring), + ) + } +} + +fn mb_per_s(bytes: u64, over: Duration) -> f64 { + bytes as f64 / over.as_secs_f64() / 1e6 +} From 7fa99b20be04c16fc68caaf603c20b6b3d5d13fd Mon Sep 17 00:00:00 2001 From: not-matthias Date: Thu, 24 Sep 2026 17:46:26 +0200 Subject: [PATCH 2/5] feat(memtrack): store the stop time in pressure_stopped --- crates/memtrack/src/ebpf/c/utils/process_stop.h | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/crates/memtrack/src/ebpf/c/utils/process_stop.h b/crates/memtrack/src/ebpf/c/utils/process_stop.h index 2d8fd1ee5..1ee23d58d 100644 --- a/crates/memtrack/src/ebpf/c/utils/process_stop.h +++ b/crates/memtrack/src/ebpf/c/utils/process_stop.h @@ -8,12 +8,12 @@ #define MEMTRACK_SIGCONT 18 #define MEMTRACK_SIGSTOP 19 -/* tgid -> 1 for every process BPF stopped, one map per reason. A process is - * recorded before its stop can take effect, and userspace resumes only - * recorded processes, so a process stopped for both reasons resumes once - * neither map holds it. Sized like tracked_pids. */ -BPF_HASH_MAP(pressure_stopped, __u32, __u8, 10000); -BPF_HASH_MAP(attach_stopped, __u32, __u8, 10000); +/* tgid -> stop ktime (ns) for every process BPF stopped, one map per reason. + * A process is recorded before its stop can take effect, and userspace + * resumes only recorded processes, so a process stopped for both reasons + * resumes once neither map holds it. Sized like tracked_pids. */ +BPF_HASH_MAP(pressure_stopped, __u32, __u64, 10000); +BPF_HASH_MAP(attach_stopped, __u32, __u64, 10000); /* Stops that could not be recorded because the map was full; userspace warns. */ BPF_ARRAY_MAP(stop_record_failed, __u64, 1); @@ -27,8 +27,8 @@ static __always_inline void memtrack_stop_current(void* map, __u32 tgid) { return; } - __u8 marker = 1; - if (bpf_map_update_elem(map, &tgid, &marker, BPF_ANY) == 0) { + __u64 stopped_at = bpf_ktime_get_ns(); + if (bpf_map_update_elem(map, &tgid, &stopped_at, BPF_ANY) == 0) { return; } From 048fb162de3dbea9131612ce18557dcf1ca41d72 Mon Sep 17 00:00:00 2001 From: not-matthias Date: Thu, 24 Sep 2026 17:46:31 +0200 Subject: [PATCH 3/5] feat(memtrack): write ring and pressure stats as JSONL CODSPEED_MEMTRACK_STATS= records each ring's producer/consumer positions around every poll tick, and each pressure release with the released pids' stop times. Rates, fill levels and pause windows are derived offline by crates/memtrack/scripts/plot_stats.py. Replaces CODSPEED_MEMTRACK_RING_STATS and its periodic log lines. --- crates/memtrack/AGENTS.md | 2 +- crates/memtrack/scripts/plot_stats.py | 136 ++++++++++++++++++ crates/memtrack/src/ebpf/memtrack/maps.rs | 6 +- crates/memtrack/src/ebpf/mod.rs | 2 +- crates/memtrack/src/ebpf/pause.rs | 20 ++- crates/memtrack/src/ebpf/poller.rs | 24 ++-- crates/memtrack/src/ebpf/ring_stats.rs | 132 ----------------- crates/memtrack/src/ebpf/stats.rs | 166 ++++++++++++++++++++++ crates/memtrack/src/main.rs | 6 +- 9 files changed, 343 insertions(+), 151 deletions(-) create mode 100644 crates/memtrack/scripts/plot_stats.py delete mode 100644 crates/memtrack/src/ebpf/ring_stats.rs create mode 100644 crates/memtrack/src/ebpf/stats.rs diff --git a/crates/memtrack/AGENTS.md b/crates/memtrack/AGENTS.md index f8f1a1ace..2ee11bdd3 100644 --- a/crates/memtrack/AGENTS.md +++ b/crates/memtrack/AGENTS.md @@ -84,7 +84,7 @@ sudo -E cargo test --test c_tests -- --test-threads 1 - **Build toolchain:** `clang` + BTF/vmlinux headers, `libbpf-dev`, `zlib1g-dev`, `pkgconf`, `build-essential`; vendored libbpf also needs `autopoint`/`bison`/`flex`. - `vmlinux.h` is pinned to a specific git rev; `libbpf-rs` uses the `vendored` feature (dist links `libbpf-rs/static`). -Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_MEMTRACK_RING_STATS` (1 logs per-ring write/read/drain rates and backlog: every second at debug, a run summary at info), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). +Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_MEMTRACK_STATS` (absolute path; writes per-tick ring positions and pressure episodes as JSONL, plotted by `scripts/plot_stats.py`), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). ### Minimum kernel version diff --git a/crates/memtrack/scripts/plot_stats.py b/crates/memtrack/scripts/plot_stats.py new file mode 100644 index 000000000..f71730ffd --- /dev/null +++ b/crates/memtrack/scripts/plot_stats.py @@ -0,0 +1,136 @@ +# /// script +# requires-python = ">=3.11" +# dependencies = ["polars", "matplotlib"] +# /// +"""Plot memtrack pipeline stats written via CODSPEED_MEMTRACK_STATS.""" + +import argparse +import json +import sys +from pathlib import Path + +import matplotlib + +matplotlib.use("Agg") +import matplotlib.pyplot as plt +import polars as pl + +MB = 1e6 +BIN_NS = 100_000_000 + + +def load(path: Path) -> tuple[pl.DataFrame, pl.DataFrame, dict[str, int]]: + rows = [json.loads(line) for line in path.read_text().splitlines() if line.strip()] + if not rows: + sys.exit(f"{path}: no records") + sizes = {r["ring"]: r["size"] for r in rows if r["k"] == "ring_open"} + ring_rows = [r for r in rows if r["k"] == "ring"] + if not ring_rows: + sys.exit(f"{path}: no ring records") + pressure_rows = [ + {"ring": r["ring"], "t": r["t"], "pid": pid, "stopped_at": at} + for r in rows + if r["k"] == "pressure" + for pid, at in r["pids"] + ] + t_min = min(r.get("t", r.get("t0", 0)) for r in rows) + rings = pl.DataFrame(ring_rows).drop("k") + rings = rings.with_columns(pl.col("t0", "t1") - t_min) + schema = {"ring": pl.Utf8, "t": pl.Int64, "pid": pl.Int64, "stopped_at": pl.Int64} + pressure = pl.DataFrame(pressure_rows, schema=schema) + pressure = pressure.with_columns(pl.col("t", "stopped_at") - t_min) + return rings, pressure, sizes + + +def episodes(pressure: pl.DataFrame) -> pl.DataFrame: + """One row per released pid; episode bounds are shared by all pids released together.""" + return pressure.with_columns( + start=pl.col("stopped_at").min().over("ring", "t"), end=pl.col("t") + ) + + +def write_rate(r: pl.DataFrame) -> pl.DataFrame: + # Positions are cumulative bytes; spread each delta over the real gap since + # the previous sample, since samples are sparse when the ring is idle. + r = r.sort("t1") + dt = pl.col("t1").diff() + return r.select("t1", mbps=pl.col("prod1").diff() / MB / (dt / 1e9)).filter(dt > 0) + + +def drain_rate(r: pl.DataFrame) -> pl.DataFrame: + # Aggregated per bin: ticks of a few us give meaningless per-tick ratios. + return ( + r.select(t1=pl.col("t1") // BIN_NS * BIN_NS, bytes=pl.col("cons1") - pl.col("cons0"), busy=pl.col("t1") - pl.col("t0")) + .group_by("t1").agg(pl.col("bytes", "busy").sum()).sort("t1") + .filter(pl.col("busy") > 0) + .select("t1", mbps=pl.col("bytes") / MB / (pl.col("busy") / 1e9)) + ) + + +def plot(rings: pl.DataFrame, eps: pl.DataFrame, sizes: dict[str, int], out: Path) -> None: + n = 3 if len(eps) else 2 + fig, axes = plt.subplots(n, 1, sharex=True, figsize=(14, 3.5 * n), squeeze=False) + ax_fill, ax_tp = axes[0][0], axes[1][0] + for name, r in rings.sort("t0").group_by("ring", maintain_order=True): + name, size = name[0], sizes.get(name[0]) + if size: + xs = [v for a, b in zip(r["t0"], r["t1"]) for v in (a, b)] + ys = [v for a, b in zip(r["prod0"] - r["cons0"], r["prod1"] - r["cons1"]) for v in (a, b)] + ax_fill.plot([x / 1e9 for x in xs], [100 * y / size for y in ys], lw=0.8, label=name) + w, d = write_rate(r), drain_rate(r) + ax_tp.step(w["t1"] / 1e9, w["mbps"], where="pre", lw=0.8, label=f"{name} write") + ax_tp.plot(d["t1"] / 1e9, d["mbps"], lw=0.8, ls=":", label=f"{name} drain") + ax_fill.axhline(75, ls="--", color="gray") + for e in eps.unique(["ring", "t"]).iter_rows(named=True): + ax_fill.axvspan(e["start"] / 1e9, e["end"] / 1e9, color="red", alpha=0.15) + ax_fill.set_ylabel("fill %") + ax_tp.set_ylabel("MB/s") + if len(eps): + ax = axes[2][0] + pids = sorted(eps["pid"].unique()) + colors = {ring: f"C{i}" for i, ring in enumerate(sorted(eps["ring"].unique()))} + for e in eps.iter_rows(named=True): + ax.barh(pids.index(e["pid"]), (e["end"] - e["stopped_at"]) / 1e9, + left=e["stopped_at"] / 1e9, color=colors[e["ring"]]) + ax.set_yticks(range(len(pids)), [str(p) for p in pids]) + ax.set_ylabel("paused pid") + for row in axes: + if row[0].get_legend_handles_labels()[0]: + row[0].legend(loc="upper right", fontsize="small") + axes[-1][0].set_xlabel("seconds") + fig.tight_layout() + fig.savefig(out, dpi=120) + + +def summary(rings: pl.DataFrame, eps: pl.DataFrame, sizes: dict[str, int]) -> None: + print(f"{'ring':<16} {'MB':>9} {'wr avg':>8} {'wr peak':>8} {'dr avg':>8} {'dr peak':>8}" + f" {'fill%':>6} {'busy%':>6} {'eps':>4} {'paused ms':>10} {'max ms':>8}") + for name, r in rings.group_by("ring", maintain_order=True): + name = name[0] + span = max(r["t1"].max() - r["t0"].min(), 1) + written = r["prod1"].max() - r["prod0"].min() + w, d = write_rate(r), drain_rate(r) + size = sizes.get(name) + fill = 100 * max((r["prod0"] - r["cons0"]).max(), (r["prod1"] - r["cons1"]).max()) / size if size else float("nan") + busy = 100 * (r["t1"] - r["t0"]).sum() / span + e = eps.filter(pl.col("ring") == name) + paused = (e["end"] - e["stopped_at"]) / 1e6 + print(f"{name:<16} {written / MB:>9.1f} {written / MB / (span / 1e9):>8.1f} {w['mbps'].max() or 0:>8.1f}" + f" {d['mbps'].mean() or 0:>8.1f} {d['mbps'].max() or 0:>8.1f} {fill:>6.1f} {busy:>6.1f}" + f" {e.unique('t').height:>4} {paused.sum():>10.1f} {paused.max() or 0:>8.1f}") + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("stats", type=Path) + parser.add_argument("-o", "--output", type=Path, default=Path("stats.png")) + args = parser.parse_args() + rings, pressure, sizes = load(args.stats) + eps = episodes(pressure) + plot(rings, eps, sizes, args.output) + summary(rings, eps, sizes) + print(f"wrote {args.output}") + + +if __name__ == "__main__": + main() diff --git a/crates/memtrack/src/ebpf/memtrack/maps.rs b/crates/memtrack/src/ebpf/memtrack/maps.rs index a78c9a954..8b961a067 100644 --- a/crates/memtrack/src/ebpf/memtrack/maps.rs +++ b/crates/memtrack/src/ebpf/memtrack/maps.rs @@ -94,10 +94,10 @@ impl MemtrackBpf { } /// Callback that resumes every pressure-stopped process. - pub(super) fn on_ring_drained(&self) -> Box { + pub(super) fn on_ring_drained(&self) -> crate::ebpf::poller::OnDrained { let stopped = self.stopped.clone(); - Box::new(move || { - if let Err(error) = stopped.release_pressure() { + Box::new(move |ring| { + if let Err(error) = stopped.release_pressure(ring) { error!("failed to release pressure-stopped producers: {error:#}"); } }) diff --git a/crates/memtrack/src/ebpf/mod.rs b/crates/memtrack/src/ebpf/mod.rs index 4e7cf730b..76c726222 100644 --- a/crates/memtrack/src/ebpf/mod.rs +++ b/crates/memtrack/src/ebpf/mod.rs @@ -4,9 +4,9 @@ mod memtrack; pub(crate) mod pause; pub(crate) mod poller; mod proc_fs; -mod ring_stats; mod spawn; mod stacks; +pub mod stats; mod tracker; pub use memtrack::{ diff --git a/crates/memtrack/src/ebpf/pause.rs b/crates/memtrack/src/ebpf/pause.rs index 5f003d678..56a6a23b7 100644 --- a/crates/memtrack/src/ebpf/pause.rs +++ b/crates/memtrack/src/ebpf/pause.rs @@ -1,3 +1,4 @@ +use crate::ebpf::stats; use crate::prelude::*; use libbpf_rs::{MapCore, MapFlags, MapHandle}; use std::sync::atomic::{AtomicU64, Ordering::Relaxed}; @@ -40,13 +41,15 @@ impl StoppedProcesses { } /// Resume every pressure-stopped producer; call once a ring is flushed. - pub(crate) fn release_pressure(&self) -> Result<()> { + pub(crate) fn release_pressure(&self, ring: &str) -> Result<()> { // Deleting while iterating restarts hash iteration, so snapshot the keys first. let keys: Vec> = self.pressure_stopped.keys().collect(); if keys.is_empty() { return Ok(()); } self.pressure_stops.fetch_add(keys.len() as u64, Relaxed); + let record_stats = stats::enabled(); + let mut stopped_at = Vec::new(); for key in keys { let pid = u32::from_le_bytes( key.as_slice() @@ -54,8 +57,23 @@ impl StoppedProcesses { .context("Invalid pressure_stopped key size")?, ); debug!("Releasing pressure stop of pid {pid}"); + // Read the stop time before release deletes the entry; a missing + // entry means the pid already exited. + if record_stats + && let Some(value) = self.pressure_stopped.lookup(&key, MapFlags::ANY)? + && let Ok(bytes) = <[u8; 8]>::try_from(value.as_slice()) + { + stopped_at.push((pid, u64::from_le_bytes(bytes))); + } Self::release(pid, &self.pressure_stopped, &self.attach_stopped)?; } + if record_stats { + stats::emit(&stats::Record::Pressure { + ring, + t: stats::now_ns(), + pids: &stopped_at, + }); + } Ok(()) } diff --git a/crates/memtrack/src/ebpf/poller.rs b/crates/memtrack/src/ebpf/poller.rs index d34f86fe1..f383e63c6 100644 --- a/crates/memtrack/src/ebpf/poller.rs +++ b/crates/memtrack/src/ebpf/poller.rs @@ -1,4 +1,4 @@ -use crate::ebpf::ring_stats::RingStats; +use crate::ebpf::stats::RingSampler; use anyhow::{Context, Result}; use libbpf_rs::{AsRawLibbpf, MapCore, RingBuffer, RingBufferBuilder, libbpf_sys}; use parking_lot::Mutex; @@ -10,6 +10,9 @@ use std::time::Duration; /// Ring-buffer poll interval shared by every poller. pub(crate) const POLL_INTERVAL_MS: u64 = 1; +/// Called with the ring's map name each time the poller finds the ring empty. +pub(crate) type OnDrained = Box; + /// Items buffered before a channel send. `std::sync::mpsc` allocates a block /// every 31 messages, so sending one item at a time makes that allocation /// dominate the pipeline; batching amortizes it over a whole batch. @@ -91,7 +94,7 @@ impl RingBufferPoller { parse: F, tx: Sender>, poll_interval_ms: u64, - on_drained: Option>, + on_drained: Option, ) -> Result where M: MapCore, @@ -127,10 +130,10 @@ impl RingBufferPoller { // poll tick, and disconnection is the shutdown signal. let (ctl, ctl_rx) = mpsc::channel::>(); let poll_thread = std::thread::spawn(move || { - let mut stats = RingStats::enabled().then(|| RingStats::new(name, ring_of(&ringbuf))); + let mut sampler = RingSampler::new(name.clone(), ring_of(&ringbuf)); loop { let control = ctl_rx.recv_timeout(Duration::from_millis(poll_interval_ms)); - let tick = stats.as_ref().map(RingStats::begin); + let tick = sampler.as_ref().map(RingSampler::begin); let running = poll_iteration( control, || consume_all(&ringbuf, ring_of(&ringbuf)), @@ -140,8 +143,8 @@ impl RingBufferPoller { &batch, &tx, ); - if let (Some(stats), Some(tick)) = (&mut stats, tick) { - stats.end(tick); + if let (Some(sampler), Some(tick)) = (&mut sampler, tick) { + sampler.end(tick); } if !running { break; @@ -149,14 +152,11 @@ impl RingBufferPoller { if let Some(on_drained) = &on_drained && unsafe { libbpf_sys::ring__avail_data_size(ring_of(&ringbuf)) } == 0 { - on_drained(); + on_drained(&name); } } if let Some(on_drained) = &on_drained { - on_drained(); - } - if let Some(stats) = &stats { - stats.report_run(); + on_drained(&name); } }); @@ -209,7 +209,7 @@ impl ThreadedRingBufferPoller { resolve: R, tx: Sender>, poll_interval_ms: u64, - on_drained: Option>, + on_drained: Option, ) -> Result where M: MapCore, diff --git a/crates/memtrack/src/ebpf/ring_stats.rs b/crates/memtrack/src/ebpf/ring_stats.rs deleted file mode 100644 index a97e7ae69..000000000 --- a/crates/memtrack/src/ebpf/ring_stats.rs +++ /dev/null @@ -1,132 +0,0 @@ -//! Ring throughput from libbpf's mmapped producer/consumer positions, enabled -//! with `CODSPEED_MEMTRACK_RING_STATS=1`. -//! -//! The producer position counts every reserved byte, including the 8-byte -//! record headers and records the BPF side later discarded, so "wrote" is ring -//! pressure rather than artifact bytes. - -use crate::prelude::*; -use libbpf_rs::libbpf_sys; -use std::time::{Duration, Instant}; - -const REPORT_INTERVAL: Duration = Duration::from_secs(1); - -pub(crate) struct RingStats { - name: String, - ring: *const libbpf_sys::ring, - run: Window, - window: Window, -} - -/// Taken right before a poll iteration: the backlog then is everything the -/// producers wrote since the previous iteration. -pub(crate) struct Tick { - started: Instant, - backlog: u64, -} - -struct Window { - started: Instant, - produced: u64, - consumed: u64, - busy: Duration, - max_tick: Duration, - peak_backlog: u64, -} - -impl RingStats { - pub(crate) fn enabled() -> bool { - std::env::var("CODSPEED_MEMTRACK_RING_STATS").is_ok_and(|v| v == "1") - } - - /// `ring` must stay valid for as long as the stats are used. - pub(crate) fn new(name: String, ring: *const libbpf_sys::ring) -> Self { - Self { - name, - ring, - run: Window::start(ring), - window: Window::start(ring), - } - } - - pub(crate) fn begin(&self) -> Tick { - Tick { - started: Instant::now(), - // SAFETY: `ring` is valid per `new`'s contract. - backlog: unsafe { libbpf_sys::ring__avail_data_size(self.ring) } as u64, - } - } - - pub(crate) fn end(&mut self, tick: Tick) { - let busy = tick.started.elapsed(); - self.run.record(tick.backlog, busy); - self.window.record(tick.backlog, busy); - if self.window.started.elapsed() < REPORT_INTERVAL { - return; - } - debug!( - "{} ring (1s): {}", - self.name, - self.window.summary(self.ring) - ); - self.window = Window::start(self.ring); - } - - pub(crate) fn report_run(&self) { - info!("{} ring (run): {}", self.name, self.run.summary(self.ring)); - } -} - -impl Window { - fn start(ring: *const libbpf_sys::ring) -> Self { - let (produced, consumed) = positions(ring); - Self { - started: Instant::now(), - produced, - consumed, - busy: Duration::ZERO, - max_tick: Duration::ZERO, - peak_backlog: 0, - } - } - - fn record(&mut self, backlog: u64, busy: Duration) { - self.busy += busy; - self.max_tick = self.max_tick.max(busy); - self.peak_backlog = self.peak_backlog.max(backlog); - } - - /// `drain` is the consume rate while the poll thread is busy, i.e. the - /// read speed the thread can sustain; `read` is averaged over wall time. - fn summary(&self, ring: *const libbpf_sys::ring) -> String { - let (produced, consumed) = positions(ring); - let wall = self.started.elapsed(); - let consumed = consumed - self.consumed; - // SAFETY: `ring` is valid per `RingStats::new`'s contract. - let size = unsafe { libbpf_sys::ring__size(ring) } as u64; - format!( - "wrote {:.1} MB/s, read {:.1} MB/s, drain {:.1} MB/s, peak backlog {:.1} MiB ({:.1}%), busy {:.1}% (max tick {:.2} ms)", - mb_per_s(produced - self.produced, wall), - mb_per_s(consumed, wall), - mb_per_s(consumed, self.busy), - self.peak_backlog as f64 / (1024.0 * 1024.0), - self.peak_backlog as f64 * 100.0 / size as f64, - self.busy.as_secs_f64() * 100.0 / wall.as_secs_f64(), - self.max_tick.as_secs_f64() * 1e3, - ) - } -} - -fn positions(ring: *const libbpf_sys::ring) -> (u64, u64) { - // SAFETY: `ring` is valid per `RingStats::new`'s contract. - unsafe { - ( - libbpf_sys::ring__producer_pos(ring), - libbpf_sys::ring__consumer_pos(ring), - ) - } -} - -fn mb_per_s(bytes: u64, over: Duration) -> f64 { - bytes as f64 / over.as_secs_f64() / 1e6 -} diff --git a/crates/memtrack/src/ebpf/stats.rs b/crates/memtrack/src/ebpf/stats.rs new file mode 100644 index 000000000..8573978ca --- /dev/null +++ b/crates/memtrack/src/ebpf/stats.rs @@ -0,0 +1,166 @@ +//! Pipeline samples as JSON lines, enabled with `CODSPEED_MEMTRACK_STATS=`. +//! +//! Only raw counters are recorded; rates and fill levels are derived offline +//! by `scripts/plot_stats.py`. Every `t*` is CLOCK_MONOTONIC ns, the clock of +//! `bpf_ktime_get_ns()` and of the artifact's event timestamps. + +use crate::prelude::*; +use libbpf_rs::libbpf_sys; +use parking_lot::Mutex; +use serde::Serialize; +use std::fs::File; +use std::io::{BufWriter, Write}; +use std::path::PathBuf; +use std::sync::OnceLock; + +/// Emptied on the first write error, so a full disk stops sampling instead of +/// logging on every tick. +static SINK: OnceLock>>> = OnceLock::new(); + +#[derive(Serialize)] +#[serde(tag = "k", rename_all = "snake_case")] +pub(crate) enum Record<'a> { + RingOpen { + t: u64, + ring: &'a str, + size: u64, + }, + Ring { + ring: &'a str, + t0: u64, + prod0: u64, + cons0: u64, + t1: u64, + prod1: u64, + cons1: u64, + }, + Pressure { + ring: &'a str, + t: u64, + pids: &'a [(u32, u64)], + }, +} + +pub fn init_from_env() -> Result<()> { + let Some(path) = std::env::var_os("CODSPEED_MEMTRACK_STATS").map(PathBuf::from) else { + return Ok(()); + }; + let file = File::create(&path) + .with_context(|| format!("Failed to create memtrack stats file {}", path.display()))?; + SINK.set(Mutex::new(Some(BufWriter::new(file)))) + .map_err(|_| anyhow!("memtrack stats already initialized"))?; + info!("Writing memtrack stats to {}", path.display()); + Ok(()) +} + +pub fn finish() -> Result<()> { + let Some(mut out) = SINK.get().and_then(|sink| sink.lock().take()) else { + return Ok(()); + }; + out.flush().context("Failed to flush memtrack stats") +} + +pub(crate) fn enabled() -> bool { + SINK.get().is_some() +} + +pub(crate) fn emit(record: &Record) { + let Some(sink) = SINK.get() else { + return; + }; + let mut sink = sink.lock(); + let Some(out) = sink.as_mut() else { + return; + }; + let written = serde_json::to_writer(&mut *out, record) + .map_err(std::io::Error::from) + .and_then(|()| out.write_all(b"\n")); + if let Err(error) = written { + error!("Stopping memtrack stats after a write error: {error}"); + *sink = None; + } +} + +pub(crate) fn now_ns() -> u64 { + let mut ts = libc::timespec { + tv_sec: 0, + tv_nsec: 0, + }; + // SAFETY: `ts` is a valid out-pointer. + unsafe { libc::clock_gettime(libc::CLOCK_MONOTONIC, &mut ts) }; + ts.tv_sec as u64 * 1_000_000_000 + ts.tv_nsec as u64 +} + +/// Positions of one ring around each poll tick, from libbpf's mmapped +/// producer/consumer pages. The producer position counts every reserved byte, +/// including 8-byte record headers and records BPF later discarded, so it is +/// ring pressure rather than artifact bytes. +pub(crate) struct RingSampler { + name: String, + ring: *const libbpf_sys::ring, + last: (u64, u64), +} + +pub(crate) struct Tick { + t: u64, + prod: u64, + cons: u64, +} + +impl RingSampler { + /// `None` unless stats are enabled. `ring` must outlive the sampler. + pub(crate) fn new(name: String, ring: *const libbpf_sys::ring) -> Option { + if !enabled() { + return None; + } + // SAFETY: `ring` is valid per this function's contract. + let size = unsafe { libbpf_sys::ring__size(ring) } as u64; + emit(&Record::RingOpen { + t: now_ns(), + ring: &name, + size, + }); + let mut sampler = Self { + name, + ring, + last: (0, 0), + }; + sampler.last = sampler.positions(); + Some(sampler) + } + + pub(crate) fn begin(&self) -> Tick { + let t = now_ns(); + let (prod, cons) = self.positions(); + Tick { t, prod, cons } + } + + /// Positions only grow, so equal end positions mean nothing was written or + /// read since the last emitted tick. + pub(crate) fn end(&mut self, tick: Tick) { + let (prod1, cons1) = self.positions(); + if (prod1, cons1) == self.last { + return; + } + self.last = (prod1, cons1); + emit(&Record::Ring { + ring: &self.name, + t0: tick.t, + prod0: tick.prod, + cons0: tick.cons, + t1: now_ns(), + prod1, + cons1, + }); + } + + fn positions(&self) -> (u64, u64) { + // SAFETY: `ring` is valid per `new`'s contract. + unsafe { + ( + libbpf_sys::ring__producer_pos(self.ring), + libbpf_sys::ring__consumer_pos(self.ring), + ) + } + } +} diff --git a/crates/memtrack/src/main.rs b/crates/memtrack/src/main.rs index 5ff58cb41..25c30934e 100644 --- a/crates/memtrack/src/main.rs +++ b/crates/memtrack/src/main.rs @@ -4,7 +4,7 @@ static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc; use clap::Parser; use ipc_channel::ipc; use memtrack::prelude::*; -use memtrack::{MemtrackIpcMessage, Tracker, handle_ipc_message}; +use memtrack::{MemtrackIpcMessage, Tracker, handle_ipc_message, stats}; use runner_shared::artifacts::{ArtifactExt, MemtrackArtifact, encode_events}; use std::path::{Path, PathBuf}; use std::process::Command; @@ -87,6 +87,7 @@ fn track_command( None }; + stats::init_from_env()?; let tracker = Arc::new(Tracker::new()?); // Spawn IPC handler thread with the now-available tracker @@ -154,6 +155,9 @@ fn track_command( // encode pipeline join below would block forever. debug!("Stopping the ring buffer poller"); drop(session); + if let Err(error) = stats::finish() { + warn!("{error:#}"); + } debug!("Waiting for the encode pipeline to finish"); let total = pipeline_thread From 0e49eba5c4f432313eb2d9a46b8ecd915eadda10 Mon Sep 17 00:00:00 2001 From: not-matthias Date: Mon, 5 Oct 2026 18:13:16 +0200 Subject: [PATCH 4/5] feat(runner-shared): report per-step encoder stats encode_events takes an on_step callback. A step fills and submits one frame, then writes every finished frame; the last step also drains the frames in flight. Each report has the time the reader waited for input, was blocked on the workers at the in-flight cap, and spent writing, plus the msgpack and zstd sizes of the frames it wrote. --- .../runner-shared/benches/memtrack_writer.rs | 2 +- .../src/artifacts/memtrack/pipeline.rs | 86 +++++++++++++++---- 2 files changed, 70 insertions(+), 18 deletions(-) diff --git a/crates/runner-shared/benches/memtrack_writer.rs b/crates/runner-shared/benches/memtrack_writer.rs index 580424a1a..1d93904b0 100644 --- a/crates/runner-shared/benches/memtrack_writer.rs +++ b/crates/runner-shared/benches/memtrack_writer.rs @@ -200,7 +200,7 @@ fn write_stack_events(bencher: Bencher) { fn encode(events: &[MemtrackEvent], n_workers: usize) -> Vec { let mut output = Vec::new(); - encode_events(events.iter().cloned(), &mut output, n_workers).unwrap(); + encode_events(events.iter().cloned(), &mut output, n_workers, |_| {}).unwrap(); output } diff --git a/crates/runner-shared/src/artifacts/memtrack/pipeline.rs b/crates/runner-shared/src/artifacts/memtrack/pipeline.rs index 8adb9c34e..a379ac314 100644 --- a/crates/runner-shared/src/artifacts/memtrack/pipeline.rs +++ b/crates/runner-shared/src/artifacts/memtrack/pipeline.rs @@ -2,6 +2,7 @@ use std::cell::RefCell; use std::collections::VecDeque; use std::io::{BufWriter, Write}; use std::sync::mpsc::{self, Receiver, TryRecvError}; +use std::time::{Duration, Instant}; use serde::Serialize; @@ -15,9 +16,29 @@ const FRAME_EVENTS: usize = 64 * 1024; /// oldest frame, bounding memory to about `cap + 1` frames. const MAX_IN_FLIGHT_PER_WORKER: usize = 2; -/// A self-contained compressed zstd frame, plus the event buffer it was built -/// from so the reader can reuse it for the next frame instead of reallocating. -type CompressedFrame = (anyhow::Result>, Vec); +/// A self-contained compressed zstd frame with its uncompressed msgpack size, +/// plus the event buffer it was built from so the reader can reuse it for the +/// next frame instead of reallocating. +type CompressedFrame = (anyhow::Result<(Vec, u64)>, Vec); + +/// Reader-thread time and output size of one encoder step, reported through +/// `encode_events`' `on_step` callback. A step pulls events until a frame is +/// full, submits it, then writes every finished frame; the last step also +/// drains the frames still in flight. Steps run back to back, so +/// `wait + encode + write` is the wall time each one covers. +#[derive(Default)] +pub struct EncodeStats { + /// Events submitted in this step. + pub events: usize, + /// Blocked on the input iterator, i.e. waiting for events. + pub wait: Duration, + /// Blocked on the workers, i.e. the in-flight queue was at the cap. + pub encode: Duration, + pub write: Duration, + /// Sizes of the frames written in this step. + pub msgpack_bytes: u64, + pub zstd_bytes: u64, +} /// Encode a stream of events into a single compressed artifact stream, /// compressing frames in parallel across a Rayon pool of `n_workers` threads. @@ -30,7 +51,12 @@ type CompressedFrame = (anyhow::Result>, Vec); /// /// Blocks the calling thread until `events` is exhausted. Returns the total /// number of events written. -pub fn encode_events(events: S, out: W, n_workers: usize) -> anyhow::Result +pub fn encode_events( + events: S, + out: W, + n_workers: usize, + mut on_step: impl FnMut(&EncodeStats), +) -> anyhow::Result where S: IntoIterator, W: Write, @@ -59,20 +85,30 @@ where rx }; let mut collect = |(encoded, mut frame): CompressedFrame, - spare: &mut Vec>| { - out.write_all(&encoded?)?; + spare: &mut Vec>, + step: &mut EncodeStats| { + let (bytes, msgpack_bytes) = encoded?; + let writing = Instant::now(); + out.write_all(&bytes)?; + step.write += writing.elapsed(); + step.msgpack_bytes += msgpack_bytes; + step.zstd_bytes += bytes.len() as u64; wrote_any = true; frame.clear(); spare.push(frame); anyhow::Ok(()) }; + let mut step = EncodeStats::default(); + let mut filling = Instant::now(); let mut frame: Vec = Vec::with_capacity(FRAME_EVENTS); for event in events { frame.push(event); if frame.len() < FRAME_EVENTS { continue; } + step.wait = filling.elapsed(); + step.events = frame.len(); total += frame.len() as u64; let next = spare .pop() @@ -87,25 +123,38 @@ where Ok(result) => result, Err(TryRecvError::Empty) if !at_cap => break, // Blocks at the cap; fails at once if the worker is gone. - Err(_) => recv_frame(rx)?, + Err(_) => { + let blocked = Instant::now(); + let result = recv_frame(rx)?; + step.encode += blocked.elapsed(); + result + } }; in_flight.pop_front(); - collect(result, &mut spare)?; + collect(result, &mut spare, &mut step)?; } + on_step(&std::mem::take(&mut step)); + filling = Instant::now(); } + step.wait = filling.elapsed(); + step.events = frame.len(); if !frame.is_empty() { total += frame.len() as u64; in_flight.push_back(submit(frame)); } for rx in in_flight.drain(..) { - collect(recv_frame(&rx)?, &mut spare)?; + let blocked = Instant::now(); + let result = recv_frame(&rx)?; + step.encode += blocked.elapsed(); + collect(result, &mut spare, &mut step)?; } + on_step(&step); // Always emit at least one (possibly empty) frame so the artifact stream is // valid and decodable even when no events were recorded. if !wrote_any { - out.write_all(&encode_frame(&[])?)?; + out.write_all(&encode_frame(&[])?.0)?; } out.flush()?; @@ -139,12 +188,13 @@ thread_local! { static FRAME_ENCODER: RefCell> = const { RefCell::new(None) }; } -/// Encode one batch as a single self-contained zstd frame. +/// Encode one batch as a single self-contained zstd frame, returned together +/// with its uncompressed msgpack size. /// /// The batch is serialized with `rmp_serde` into a reused buffer, then /// compressed in one shot with a reused zstd context. The decoded bytes are the /// same msgpack stream `MemtrackWriter` produces. -fn encode_frame(batch: &[MemtrackEvent]) -> anyhow::Result> { +fn encode_frame(batch: &[MemtrackEvent]) -> anyhow::Result<(Vec, u64)> { FRAME_ENCODER.with(|cell| { let mut slot = cell.borrow_mut(); let enc = match slot.as_mut() { @@ -167,7 +217,7 @@ fn encode_frame(batch: &[MemtrackEvent]) -> anyhow::Result> { // Trim the worst-case compression bound so in-flight frames only hold // their actual size. compressed.shrink_to_fit(); - Ok(compressed) + Ok((compressed, enc.msgpack.len() as u64)) }) } @@ -200,7 +250,7 @@ mod tests { let events = malloc_events(0..(FRAME_EVENTS as u64 * 3 + 7)); let mut out = Vec::new(); - let total = encode_events(events.clone(), &mut out, 4)?; + let total = encode_events(events.clone(), &mut out, 4, |_| {})?; assert_eq!(total, events.len() as u64); let decoded: Vec<_> = MemtrackArtifact::decode_streamed(Cursor::new(out))?.collect(); @@ -214,7 +264,7 @@ mod tests { let events = malloc_events(0..(FRAME_EVENTS as u64 * 5 + 3)); let mut out = Vec::new(); - let total = encode_events(events.clone(), &mut out, 1)?; + let total = encode_events(events.clone(), &mut out, 1, |_| {})?; assert_eq!(total, events.len() as u64); let decoded: Vec<_> = MemtrackArtifact::decode_streamed(Cursor::new(out))?.collect(); @@ -235,8 +285,10 @@ mod tests { // Encode twice to also exercise the reused per-thread buffers. for _ in 0..2 { - let frame = zstd::decode_all(Cursor::new(encode_frame(&events)?))?; + let (frame, msgpack_bytes) = encode_frame(&events)?; + let frame = zstd::decode_all(Cursor::new(frame))?; assert_eq!(frame, reference); + assert_eq!(msgpack_bytes, reference.len() as u64); } Ok(()) @@ -247,7 +299,7 @@ mod tests { let events: Vec = Vec::new(); let mut out = Vec::new(); - let total = encode_events(events, &mut out, 4)?; + let total = encode_events(events, &mut out, 4, |_| {})?; assert_eq!(total, 0); assert!(MemtrackArtifact::is_empty(Cursor::new(out))); From 991e47c2be618a7440cb6e6d5170cbf6b7507385 Mon Sep 17 00:00:00 2001 From: not-matthias Date: Thu, 24 Sep 2026 18:52:56 +0200 Subject: [PATCH 5/5] feat(memtrack): record backlog, resolver and encoder stats With CODSPEED_MEMTRACK_STATS set, memtrack also writes: - backlog rows every 10 ms: events parsed from the rings vs. taken by the encoder, plus memtrack's own RSS - one resolve row per stack-resolver batch - one encode row per encoder step plot_stats.py adds stage-busy and backlog/RSS panels, encoder in/out throughput, and a log throughput axis. --- crates/memtrack/AGENTS.md | 2 +- crates/memtrack/scripts/plot_stats.py | 155 ++++++++++++++++------- crates/memtrack/src/ebpf/memtrack/mod.rs | 15 ++- crates/memtrack/src/ebpf/poller.rs | 12 +- crates/memtrack/src/ebpf/stats.rs | 99 +++++++++++++++ crates/memtrack/src/main.rs | 15 ++- crates/memtrack/src/perf_mappings.rs | 1 + 7 files changed, 244 insertions(+), 55 deletions(-) diff --git a/crates/memtrack/AGENTS.md b/crates/memtrack/AGENTS.md index 2ee11bdd3..ffa1bdf3d 100644 --- a/crates/memtrack/AGENTS.md +++ b/crates/memtrack/AGENTS.md @@ -84,7 +84,7 @@ sudo -E cargo test --test c_tests -- --test-threads 1 - **Build toolchain:** `clang` + BTF/vmlinux headers, `libbpf-dev`, `zlib1g-dev`, `pkgconf`, `build-essential`; vendored libbpf also needs `autopoint`/`bison`/`flex`. - `vmlinux.h` is pinned to a specific git rev; `libbpf-rs` uses the `vendored` feature (dist links `libbpf-rs/static`). -Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_MEMTRACK_STATS` (absolute path; writes per-tick ring positions and pressure episodes as JSONL, plotted by `scripts/plot_stats.py`), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). +Env vars actually wired: `CODSPEED_MEMTRACK_BINARIES` (extra static-allocator binaries), `CODSPEED_MEMTRACK_TRACK_ALLOCATORS` (0/false disables), `CODSPEED_MEMTRACK_TRACK_PHYSICAL` (1 enables), `CODSPEED_MEMTRACK_CAPTURE_STACKS` (1 enables), `CODSPEED_MEMTRACK_STACK_BUDGET` (stack copy size in bytes, default 8192), `CODSPEED_MEMTRACK_STATS` (absolute path; writes per-tick ring positions, pressure episodes, pipeline backlog + RSS, stack-resolver batches and encoder steps as JSONL, plotted by `scripts/plot_stats.py`), `CODSPEED_LOG` (log filter, default `info`), `SUDO_UID`/`SUDO_GID` (privilege drop), `GITHUB_ACTIONS` (build rebuild trigger + test gate). ### Minimum kernel version diff --git a/crates/memtrack/scripts/plot_stats.py b/crates/memtrack/scripts/plot_stats.py index f71730ffd..5270a4f5d 100644 --- a/crates/memtrack/scripts/plot_stats.py +++ b/crates/memtrack/scripts/plot_stats.py @@ -16,37 +16,43 @@ import polars as pl MB = 1e6 +MIB = 1024 * 1024 BIN_NS = 100_000_000 - - -def load(path: Path) -> tuple[pl.DataFrame, pl.DataFrame, dict[str, int]]: - rows = [json.loads(line) for line in path.read_text().splitlines() if line.strip()] - if not rows: - sys.exit(f"{path}: no records") - sizes = {r["ring"]: r["size"] for r in rows if r["k"] == "ring_open"} - ring_rows = [r for r in rows if r["k"] == "ring"] - if not ring_rows: +TIME_COLS = ("t", "t0", "t1", "stopped_at") + + +def load(path: Path) -> dict[str, pl.DataFrame]: + rows = [] + for n, line in enumerate(path.read_text().splitlines(), 1): + try: + rows.append(json.loads(line)) + except json.JSONDecodeError: + # A killed run can leave a cut-off last line. + print(f"{path}:{n}: skipping unparsable line", file=sys.stderr) + if not any(r["k"] == "ring" for r in rows): sys.exit(f"{path}: no ring records") - pressure_rows = [ + t_min = min(r.get("t", r.get("t0", 0)) for r in rows) + frames = {} + for kind in ("ring_open", "ring", "backlog", "resolve", "encode"): + kind_rows = [r for r in rows if r["k"] == kind] + frames[kind] = pl.DataFrame(kind_rows).drop("k") if kind_rows else pl.DataFrame() + pressure = [ {"ring": r["ring"], "t": r["t"], "pid": pid, "stopped_at": at} for r in rows if r["k"] == "pressure" for pid, at in r["pids"] ] - t_min = min(r.get("t", r.get("t0", 0)) for r in rows) - rings = pl.DataFrame(ring_rows).drop("k") - rings = rings.with_columns(pl.col("t0", "t1") - t_min) schema = {"ring": pl.Utf8, "t": pl.Int64, "pid": pl.Int64, "stopped_at": pl.Int64} - pressure = pl.DataFrame(pressure_rows, schema=schema) - pressure = pressure.with_columns(pl.col("t", "stopped_at") - t_min) - return rings, pressure, sizes + frames["pressure"] = pl.DataFrame(pressure, schema=schema) + for kind, df in frames.items(): + if kind != "ring_open" and len(df): + frames[kind] = df.with_columns(pl.col(c) - t_min for c in TIME_COLS if c in df.columns) + return frames def episodes(pressure: pl.DataFrame) -> pl.DataFrame: """One row per released pid; episode bounds are shared by all pids released together.""" - return pressure.with_columns( - start=pl.col("stopped_at").min().over("ring", "t"), end=pl.col("t") - ) + return pressure.with_columns(start=pl.col("stopped_at").min().over("ring", "t"), end=pl.col("t")) def write_rate(r: pl.DataFrame) -> pl.DataFrame: @@ -67,45 +73,89 @@ def drain_rate(r: pl.DataFrame) -> pl.DataFrame: ) -def plot(rings: pl.DataFrame, eps: pl.DataFrame, sizes: dict[str, int], out: Path) -> None: - n = 3 if len(eps) else 2 - fig, axes = plt.subplots(n, 1, sharex=True, figsize=(14, 3.5 * n), squeeze=False) - ax_fill, ax_tp = axes[0][0], axes[1][0] - for name, r in rings.sort("t0").group_by("ring", maintain_order=True): - name, size = name[0], sizes.get(name[0]) - if size: +def busy_pct(intervals: pl.DataFrame) -> pl.DataFrame: + """Share of each bin a thread spent inside short (t0, t1) work intervals.""" + return ( + intervals.select(t=pl.col("t1") // BIN_NS * BIN_NS, busy=pl.col("t1") - pl.col("t0")) + .group_by("t").agg(pl.col("busy").sum()).sort("t") + .select("t", pct=100 * pl.col("busy") / BIN_NS) + ) + + +def encoder_step_ns() -> pl.Expr: + # Steps run back to back, so wait + encode + write is the wall time each one covers. + return pl.col("wait_ns") + pl.col("encode_ns") + pl.col("write_ns") + + +def plot(f: dict[str, pl.DataFrame], eps: pl.DataFrame, out: Path) -> None: + sizes = dict(zip(f["ring_open"]["ring"], f["ring_open"]["size"])) if len(f["ring_open"]) else {} + rings, enc, backlog, resolve = f["ring"], f["encode"], f["backlog"], f["resolve"] + n = 4 + bool(len(eps)) + fig, axes = plt.subplots(n, 1, sharex=True, figsize=(14, 3.2 * n)) + ax_fill, ax_tp, ax_busy, ax_backlog = axes[:4] + + # One color per ring across every panel. + for i, name in enumerate(sorted(rings["ring"].unique())): + r, color = rings.filter(pl.col("ring") == name).sort("t0"), f"C{i}" + if size := sizes.get(name): xs = [v for a, b in zip(r["t0"], r["t1"]) for v in (a, b)] ys = [v for a, b in zip(r["prod0"] - r["cons0"], r["prod1"] - r["cons1"]) for v in (a, b)] - ax_fill.plot([x / 1e9 for x in xs], [100 * y / size for y in ys], lw=0.8, label=name) + ax_fill.plot([x / 1e9 for x in xs], [100 * y / size for y in ys], color=color, lw=0.8, label=name) w, d = write_rate(r), drain_rate(r) - ax_tp.step(w["t1"] / 1e9, w["mbps"], where="pre", lw=0.8, label=f"{name} write") - ax_tp.plot(d["t1"] / 1e9, d["mbps"], lw=0.8, ls=":", label=f"{name} drain") + ax_tp.step(w["t1"] / 1e9, w["mbps"], where="pre", color=color, lw=0.8, label=f"{name} ring write") + ax_tp.plot(d["t1"] / 1e9, d["mbps"], color=color, lw=0.8, ls=":", label=f"{name} drain (while busy)") + b = busy_pct(r) + ax_busy.step(b["t"] / 1e9, b["pct"], where="post", color=color, lw=0.8, label=f"{name} poller") ax_fill.axhline(75, ls="--", color="gray") for e in eps.unique(["ring", "t"]).iter_rows(named=True): - ax_fill.axvspan(e["start"] / 1e9, e["end"] / 1e9, color="red", alpha=0.15) - ax_fill.set_ylabel("fill %") - ax_tp.set_ylabel("MB/s") + for ax in axes: + ax.axvspan(e["start"] / 1e9, e["end"] / 1e9, color="red", alpha=0.08, lw=0) + ax_fill.set_ylabel("ring fill %") + + if len(enc): + # A step can span seconds, so draw each one across the time it covers. + e = enc.with_columns(step=encoder_step_ns()) + start, end = (e["t"] - e["step"]) / 1e9, e["t"] / 1e9 + for col, label, color in (("msgpack_bytes", "encoder in (msgpack)", "C6"), ("zstd_bytes", "encoder out (zstd, disk)", "C7")): + ax_tp.hlines(e[col] / MB / (e["step"] / 1e9), start, end, color=color, lw=2, label=label) + busy = 100 * (e["encode_ns"] + e["write_ns"]) / e["step"] + ax_busy.hlines(busy, start, end, color="C6", lw=2, label="encoder (blocked on workers + write)") + if len(resolve): + b = busy_pct(resolve) + ax_busy.step(b["t"] / 1e9, b["pct"], where="post", color="C5", lw=0.8, label="stack resolver") + ax_tp.set_yscale("log") + ax_tp.set_ylabel("MB/s (log)") + ax_busy.set_ylabel("stage busy %") + ax_busy.set_ylim(0, 105) + + if len(backlog): + depth = (backlog["sent"] - backlog["received"]) / 1e6 + ax_backlog.plot(backlog["t"] / 1e9, depth, color="C2", lw=1, label="events in flight") + ax_backlog.set_ylabel("M events in flight", color="C2") + ax_rss = ax_backlog.twinx() + ax_rss.plot(backlog["t"] / 1e9, backlog["rss"] / MIB, color="gray", lw=1, ls=":") + ax_rss.set_ylabel("memtrack RSS (MiB)", color="gray") + if len(eps): - ax = axes[2][0] + ax = axes[4] pids = sorted(eps["pid"].unique()) - colors = {ring: f"C{i}" for i, ring in enumerate(sorted(eps["ring"].unique()))} for e in eps.iter_rows(named=True): - ax.barh(pids.index(e["pid"]), (e["end"] - e["stopped_at"]) / 1e9, - left=e["stopped_at"] / 1e9, color=colors[e["ring"]]) + ax.barh(pids.index(e["pid"]), (e["end"] - e["stopped_at"]) / 1e9, left=e["stopped_at"] / 1e9, color="C4") ax.set_yticks(range(len(pids)), [str(p) for p in pids]) ax.set_ylabel("paused pid") - for row in axes: - if row[0].get_legend_handles_labels()[0]: - row[0].legend(loc="upper right", fontsize="small") - axes[-1][0].set_xlabel("seconds") + for ax in axes: + if ax.get_legend_handles_labels()[0]: + ax.legend(loc="upper right", fontsize="small") + axes[-1].set_xlabel("seconds") fig.tight_layout() fig.savefig(out, dpi=120) -def summary(rings: pl.DataFrame, eps: pl.DataFrame, sizes: dict[str, int]) -> None: +def summary(f: dict[str, pl.DataFrame], eps: pl.DataFrame) -> None: + sizes = dict(zip(f["ring_open"]["ring"], f["ring_open"]["size"])) if len(f["ring_open"]) else {} print(f"{'ring':<16} {'MB':>9} {'wr avg':>8} {'wr peak':>8} {'dr avg':>8} {'dr peak':>8}" f" {'fill%':>6} {'busy%':>6} {'eps':>4} {'paused ms':>10} {'max ms':>8}") - for name, r in rings.group_by("ring", maintain_order=True): + for name, r in f["ring"].group_by("ring", maintain_order=True): name = name[0] span = max(r["t1"].max() - r["t0"].min(), 1) written = r["prod1"].max() - r["prod0"].min() @@ -119,16 +169,29 @@ def summary(rings: pl.DataFrame, eps: pl.DataFrame, sizes: dict[str, int]) -> No f" {d['mbps'].mean() or 0:>8.1f} {d['mbps'].max() or 0:>8.1f} {fill:>6.1f} {busy:>6.1f}" f" {e.unique('t').height:>4} {paused.sum():>10.1f} {paused.max() or 0:>8.1f}") + if len(enc := f["encode"]): + wall = enc.select(encoder_step_ns().sum()).item() + msgpack, zstd = enc["msgpack_bytes"].sum(), enc["zstd_bytes"].sum() + print(f"encoder: {len(enc)} steps, {enc['events'].sum()} events, {msgpack / MB:.1f} MB msgpack ->" + f" {zstd / MB:.1f} MB zstd ({msgpack / max(zstd, 1):.1f}x), {msgpack / MB / (wall / 1e9):.1f} MB/s in," + f" busy {100 * (enc['encode_ns'].sum() + enc['write_ns'].sum()) / wall:.1f}%") + if len(res := f["resolve"]): + span = max(res["t1"].max() - res["t0"].min(), 1) + print(f"resolver: {len(res)} batches, {res['n'].sum()} stacks, busy {100 * (res['t1'] - res['t0']).sum() / span:.1f}%") + if len(bl := f["backlog"]): + depth = bl["sent"] - bl["received"] + print(f"backlog: peak {depth.max()} events in flight, peak RSS {bl['rss'].max() / MIB:.1f} MiB") + def main() -> None: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("stats", type=Path) parser.add_argument("-o", "--output", type=Path, default=Path("stats.png")) args = parser.parse_args() - rings, pressure, sizes = load(args.stats) - eps = episodes(pressure) - plot(rings, eps, sizes, args.output) - summary(rings, eps, sizes) + frames = load(args.stats) + eps = episodes(frames["pressure"]) + plot(frames, eps, args.output) + summary(frames, eps) print(f"wrote {args.output}") diff --git a/crates/memtrack/src/ebpf/memtrack/mod.rs b/crates/memtrack/src/ebpf/memtrack/mod.rs index dbe4680ad..b2ea068ed 100644 --- a/crates/memtrack/src/ebpf/memtrack/mod.rs +++ b/crates/memtrack/src/ebpf/memtrack/mod.rs @@ -243,9 +243,14 @@ impl MemtrackBpf { poll_interval_ms: u64, tx: std::sync::mpsc::Sender>, ) -> Result { + let parse = |data: &[u8]| { + let event = crate::ebpf::events::parse_event(data)?; + crate::ebpf::stats::add_sent(1); + Some(event) + }; with_skel!(self, skel => RingBufferPoller::new( &skel.maps.events, - crate::ebpf::events::parse_event, + parse, tx, poll_interval_ms, None, @@ -276,9 +281,15 @@ impl MemtrackBpf { event }; + let parse = |data: &[u8]| { + let stack = events::parse_stack(data)?; + crate::ebpf::stats::add_sent(1); + Some(stack) + }; + with_skel!(self, skel => ThreadedRingBufferPoller::new( &skel.maps.stacks, - events::parse_stack, + parse, resolve, tx, poll_interval_ms, diff --git a/crates/memtrack/src/ebpf/poller.rs b/crates/memtrack/src/ebpf/poller.rs index f383e63c6..8439ffa97 100644 --- a/crates/memtrack/src/ebpf/poller.rs +++ b/crates/memtrack/src/ebpf/poller.rs @@ -1,4 +1,4 @@ -use crate::ebpf::stats::RingSampler; +use crate::ebpf::stats::{self, RingSampler}; use anyhow::{Context, Result}; use libbpf_rs::{AsRawLibbpf, MapCore, RingBuffer, RingBufferBuilder, libbpf_sys}; use parking_lot::Mutex; @@ -221,8 +221,18 @@ impl ThreadedRingBufferPoller { let (parsed_tx, parsed_rx) = mpsc::channel::>(); let ring = RingBufferPoller::new(rb_map, parse, parsed_tx, poll_interval_ms, on_drained)?; let resolver = std::thread::spawn(move || { + let record_stats = stats::enabled(); for batch in parsed_rx { + let t0 = record_stats.then(stats::now_ns); + let n = batch.len(); let resolved = batch.into_iter().map(&resolve).collect(); + if let Some(t0) = t0 { + stats::emit(&stats::Record::Resolve { + t0, + t1: stats::now_ns(), + n, + }); + } let _ = tx.send(resolved); } }); diff --git a/crates/memtrack/src/ebpf/stats.rs b/crates/memtrack/src/ebpf/stats.rs index 8573978ca..dba8bb43d 100644 --- a/crates/memtrack/src/ebpf/stats.rs +++ b/crates/memtrack/src/ebpf/stats.rs @@ -7,16 +7,31 @@ use crate::prelude::*; use libbpf_rs::libbpf_sys; use parking_lot::Mutex; +use runner_shared::artifacts::EncodeStats; use serde::Serialize; use std::fs::File; use std::io::{BufWriter, Write}; use std::path::PathBuf; use std::sync::OnceLock; +use std::sync::atomic::{AtomicU64, Ordering::Relaxed}; +use std::sync::mpsc::{self, RecvTimeoutError, Sender}; +use std::thread::JoinHandle; +use std::time::Duration; /// Emptied on the first write error, so a full disk stops sampling instead of /// logging on every tick. static SINK: OnceLock>>> = OnceLock::new(); +/// Events parsed from the rings vs. taken by the encoder. The difference is +/// everything in flight between them: poll batches, the resolver queue and +/// the unbounded encoder channel. +static SENT: AtomicU64 = AtomicU64::new(0); +static RECEIVED: AtomicU64 = AtomicU64::new(0); + +/// Poll-independent so the backlog stays visible under slow poll intervals. +const BACKLOG_INTERVAL: Duration = Duration::from_millis(10); +static BACKLOG_THREAD: Mutex, JoinHandle<()>)>> = Mutex::new(None); + #[derive(Serialize)] #[serde(tag = "k", rename_all = "snake_case")] pub(crate) enum Record<'a> { @@ -39,6 +54,30 @@ pub(crate) enum Record<'a> { t: u64, pids: &'a [(u32, u64)], }, + /// `rss` is memtrack's own resident set in bytes. + Backlog { + t: u64, + sent: u64, + received: u64, + rss: u64, + }, + /// One stack-resolver batch of `n` records. + Resolve { + t0: u64, + t1: u64, + n: usize, + }, + /// One encoder step, emitted after its frames are written at `t`. + /// `encode_ns` is the time the reader was blocked on the frame workers. + Encode { + t: u64, + events: usize, + wait_ns: u64, + encode_ns: u64, + write_ns: u64, + msgpack_bytes: u64, + zstd_bytes: u64, + }, } pub fn init_from_env() -> Result<()> { @@ -50,16 +89,76 @@ pub fn init_from_env() -> Result<()> { SINK.set(Mutex::new(Some(BufWriter::new(file)))) .map_err(|_| anyhow!("memtrack stats already initialized"))?; info!("Writing memtrack stats to {}", path.display()); + + let (stop, stopped) = mpsc::channel::<()>(); + let thread = std::thread::spawn(move || { + while stopped.recv_timeout(BACKLOG_INTERVAL) == Err(RecvTimeoutError::Timeout) { + emit_backlog(); + } + emit_backlog(); + }); + *BACKLOG_THREAD.lock() = Some((stop, thread)); Ok(()) } pub fn finish() -> Result<()> { + if let Some((stop, thread)) = BACKLOG_THREAD.lock().take() { + drop(stop); + let _ = thread.join(); + } let Some(mut out) = SINK.get().and_then(|sink| sink.lock().take()) else { return Ok(()); }; out.flush().context("Failed to flush memtrack stats") } +pub(crate) fn add_sent(events: usize) { + if enabled() { + SENT.fetch_add(events as u64, Relaxed); + } +} + +pub fn add_received(events: usize) { + if enabled() { + RECEIVED.fetch_add(events as u64, Relaxed); + } +} + +pub fn encoder_step(step: &EncodeStats) { + emit(&Record::Encode { + t: now_ns(), + events: step.events, + wait_ns: step.wait.as_nanos() as u64, + encode_ns: step.encode.as_nanos() as u64, + write_ns: step.write.as_nanos() as u64, + msgpack_bytes: step.msgpack_bytes, + zstd_bytes: step.zstd_bytes, + }); +} + +fn emit_backlog() { + emit(&Record::Backlog { + t: now_ns(), + sent: SENT.load(Relaxed), + received: RECEIVED.load(Relaxed), + rss: self_rss_bytes(), + }); +} + +/// 0 if `/proc/self/statm` is unreadable; the sample is diagnostic only. +fn self_rss_bytes() -> u64 { + let Ok(statm) = std::fs::read_to_string("/proc/self/statm") else { + return 0; + }; + let pages: u64 = statm + .split_whitespace() + .nth(1) + .and_then(|v| v.parse().ok()) + .unwrap_or(0); + // SAFETY: sysconf has no preconditions. + pages * unsafe { libc::sysconf(libc::_SC_PAGESIZE) } as u64 +} + pub(crate) fn enabled() -> bool { SINK.get().is_some() } diff --git a/crates/memtrack/src/main.rs b/crates/memtrack/src/main.rs index 25c30934e..96d3945ea 100644 --- a/crates/memtrack/src/main.rs +++ b/crates/memtrack/src/main.rs @@ -134,8 +134,13 @@ fn track_command( .map(|n| n.get().saturating_sub(2).max(1)) .unwrap_or(4); - let pipeline_thread = - thread::spawn(move || encode_events(event_rx.into_iter().flatten(), out_file, n_workers)); + let pipeline_thread = thread::spawn(move || { + let events = event_rx + .into_iter() + .inspect(|batch| stats::add_received(batch.len())) + .flatten(); + encode_events(events, out_file, n_workers, stats::encoder_step) + }); // A worker failure must not skip disabling tracking, draining, joining the // encoder, or detaching probes. Keep the wait result until teardown is done. @@ -155,9 +160,6 @@ fn track_command( // encode pipeline join below would block forever. debug!("Stopping the ring buffer poller"); drop(session); - if let Err(error) = stats::finish() { - warn!("{error:#}"); - } debug!("Waiting for the encode pipeline to finish"); let total = pipeline_thread @@ -168,6 +170,9 @@ fn track_command( if let Ok(total) = &total { info!("Wrote {total} memtrack events to disk"); } + if let Err(error) = stats::finish() { + warn!("{error:#}"); + } // Stop background workers after the ring pipeline has drained. Fatal // worker errors mean the capture is incomplete. diff --git a/crates/memtrack/src/perf_mappings.rs b/crates/memtrack/src/perf_mappings.rs index 33c2122a9..6dd6cd3ab 100644 --- a/crates/memtrack/src/perf_mappings.rs +++ b/crates/memtrack/src/perf_mappings.rs @@ -66,6 +66,7 @@ impl PerfMappingPoller { } mappings.sort_unstable_by_key(|event| (event.pid, event.timestamp)); if !mappings.is_empty() { + crate::ebpf::stats::add_sent(mappings.len()); let _ = tx.send(mappings); } });