diff --git a/snmalloc-rs/README.md b/snmalloc-rs/README.md index 876eac028..bdcf41f05 100644 --- a/snmalloc-rs/README.md +++ b/snmalloc-rs/README.md @@ -93,6 +93,32 @@ Then render to SVG: inferno-flamegraph < heap.folded > heap.svg ``` +### When to use snapshot vs streaming + +The two profiling modes answer different questions and have different +biases. Pick the one that matches your workload: + +| | `SnMalloc::snapshot()` | `ProfilingSession::start` (streaming) | +| - | - | - | +| **What it captures** | Sampled allocations *currently live* in the process at the time of the call. | *Every* sampled event (alloc / dealloc / resize) as it happens. | +| **Best for** | "What is holding memory *right now*?" — heap-state audits, leak triage, before/after diffs across a steady-state. | "Which call site is the highest-rate allocator?" — hot-path optimisation, rate-based attribution, transient-churn analysis. | +| **Bias** | Biased toward long-lived allocations — short-lived churn (allocate-and-free inside a request, scratch buffers in a tight loop) is freed before the snapshot and vanishes from view. | None on the event stream itself, but the consumer pays for storage / aggregation of every event. | +| **Output** | In-process `HeapProfile`; serialise via `write_pprof` / `write_flamegraph`. | Live callback; the application chooses how to persist events (commonly a JSON-Lines log file). | +| **Tooling** | `snmalloc-tools profile-top` for top-N live sites. | `snmalloc-tools rate-report` for per-site alloc/dealloc rate + peak-live-bytes. | + +**Rule of thumb.** If the question is "where is my live heap?" use a +snapshot. If the question is "which call site is hottest and how +churny is it?" use streaming. A snapshot will systematically +under-count a hot allocate-and-free site; a streaming log captures +that churn but requires you to keep an event log around. + +The `snmalloc-tools` CLI ships dedicated subcommands for each mode: +`profile-top` walks a snapshot, and `rate-report` stream-parses a +streaming event log file without loading the whole log into memory +(safe for multi-million-event traces). See +[`snmalloc-tools/README.md`](../snmalloc-tools/README.md) for the +streaming-log on-disk schema. + ### Streaming mode For long-running services, `ProfilingSession::start` registers a diff --git a/snmalloc-tools/README.md b/snmalloc-tools/README.md index 170bf897d..0aacfc599 100644 --- a/snmalloc-tools/README.md +++ b/snmalloc-tools/README.md @@ -27,10 +27,56 @@ snmalloc-tools branch-misses --perf-script --hints [- Parse `perf script` output and cross-reference with the Phase 10.2 branch-hint inventory. High-miss-rate inverted hints are candidates for `LIKELY` <-> `UNLIKELY` swap. + +snmalloc-tools rate-report --input [--top N] [--pretty] + Stream-parse a snmalloc streaming event log (JSON Lines) and + emit a per-site row: alloc/dealloc counts, peak live bytes, + alloc-rate per second. Output is CSV by default; `--pretty` + emits a fixed-width table. Stream-based — 6M-event logs use + O(distinct sites) memory, not O(events). +``` + +All subcommands except `rate-report` accept `--json` for structured +output; the default is a plain-text table. `rate-report` emits CSV +by default (the friendliest format for downstream awk/jq/spreadsheet +pipelines) and a fixed-width table under `--pretty`. + +## Streaming event-log schema (`rate-report`) + +`rate-report` consumes **JSON Lines** (UTF-8, one event object per +line). The producer is typically an application using +[`snmalloc_rs::ProfilingSession`](../snmalloc-rs/src/streaming.rs) +that serialises each callback to a file. Schema: + +```jsonl +{"ts_ns": 1000000, "kind": "alloc", "site": "0x55a0c0001000", "size": 4096} +{"ts_ns": 1001000, "kind": "dealloc", "site": "0x55a0c0001000", "size": 4096} ``` -All subcommands accept `--json` for structured output; the default is -a plain-text table. +Fields: + +- `ts_ns` (u64, optional) — monotonic-clock timestamp in nanoseconds. + Used to compute the alloc-rate denominator; when missing across all + records the rate column is reported as `0.0`. +- `kind` (string, required) — one of `"alloc"`, `"dealloc"`, + `"resize"`. Unknown values are skipped (forward-compat). +- `site` (string, required) — the allocation site key. Typically the + leaf-frame address as `0x` + 16 hex digits, matching the + `site_leaf` field emitted by the other subcommands. +- `size` (u64, optional) — bytes attributable to this event. + +Malformed lines are skipped silently — the reader is resilient to +truncated tails and the occasional blank line. See +`tests/fixtures/streaming_log_sample.jsonl` for a worked example. + +## Snapshot vs streaming + +`profile-top` walks a `HeapProfile::snapshot()` (currently-live +sampled allocations) and is biased toward long-lived state; +`rate-report` walks a streaming log and captures transient churn. +See the "When to use snapshot vs streaming" section in +[`../snmalloc-rs/README.md`](../snmalloc-rs/README.md) for a fuller +treatment of the tradeoff. ## Live-process limitation (important) @@ -70,6 +116,9 @@ the branch-hint inventory is a static sidecar. - `perf_c2c_sample.txt` — two contended cache lines with detail rows. - `branch_hints_sample.json` — three hint sites matching the schema in `scripts/dump_branch_hints.py`. +- `streaming_log_sample.jsonl` — eight events across two sites, + exercising alloc, dealloc, resize, and the peak-then-drop pattern + that `rate-report` is built to surface. The integration tests in `tests/integration.rs` exercise each parser/joiner against these fixtures. diff --git a/snmalloc-tools/src/lib.rs b/snmalloc-tools/src/lib.rs index 45fb462f5..289f0184e 100644 --- a/snmalloc-tools/src/lib.rs +++ b/snmalloc-tools/src/lib.rs @@ -7,3 +7,4 @@ pub mod branch_hints; pub mod joiner; pub mod perf_c2c; pub mod perf_script; +pub mod rate_report; diff --git a/snmalloc-tools/src/main.rs b/snmalloc-tools/src/main.rs index 3c7f6739a..44ab5537b 100644 --- a/snmalloc-tools/src/main.rs +++ b/snmalloc-tools/src/main.rs @@ -32,6 +32,7 @@ use snmalloc_tools::branch_hints::{BranchHintIndex, HintKind}; use snmalloc_tools::joiner; use snmalloc_tools::perf_c2c::{self, C2cLine}; use snmalloc_tools::perf_script; +use snmalloc_tools::rate_report; /// snmalloc-tools — CLI for joining perf PMU output with snmalloc's /// in-tree allocation-site lookup and branch-hint inventory. @@ -57,6 +58,9 @@ enum Cmd { /// Cross-reference `perf script` branch-miss samples with the /// Phase 10.2 branch-hint inventory. BranchMisses(BranchMissesArgs), + /// Stream-parse a snmalloc streaming event log and emit a per-site + /// rate report (alloc/dealloc counts, peak live bytes, alloc rate). + RateReport(RateReportArgs), } #[derive(Args, Debug)] @@ -121,6 +125,25 @@ struct C2cArgs { json: bool, } +#[derive(Args, Debug)] +struct RateReportArgs { + /// Path to the streaming event log to read. Must be JSON-Lines + /// (one event object per line). See + /// `snmalloc_tools::rate_report` module docs for the schema. The + /// file is stream-parsed -- 6M-event logs are fine. + #[arg(long)] + input: PathBuf, + /// Limit to the top-N highest-alloc-count sites. `0` means "no + /// limit"; the report still arrives sorted by alloc-count desc. + #[arg(long, default_value_t = 0)] + top: usize, + /// Render as a fixed-width pretty table instead of CSV. The + /// default (CSV) is the friendliest format for downstream + /// awk/jq/spreadsheet pipelines. + #[arg(long)] + pretty: bool, +} + #[derive(Args, Debug)] struct BranchMissesArgs { /// Path to the `perf script` output to parse. @@ -146,6 +169,7 @@ fn main() -> Result<()> { PmuJoinKind::C2c(c) => run_c2c(c), }, Cmd::BranchMisses(a) => run_branch_misses(a), + Cmd::RateReport(a) => run_rate_report(a), } } @@ -375,3 +399,23 @@ fn run_branch_misses(args: BranchMissesArgs) -> Result<()> { Ok(()) } +// -- rate-report ---------------------------------------------------------- + +fn run_rate_report(args: RateReportArgs) -> Result<()> { + let mut rows = rate_report::read_path(&args.input)?; + if args.top > 0 && rows.len() > args.top { + rows.truncate(args.top); + } + // Writing to a locked stdout once per run is materially faster + // than repeated `println!` for large reports, and matters when a + // user pipes the output through downstream tools. + let stdout = std::io::stdout(); + let mut out = stdout.lock(); + if args.pretty { + rate_report::write_pretty(&rows, &mut out)?; + } else { + rate_report::write_csv(&rows, &mut out)?; + } + Ok(()) +} + diff --git a/snmalloc-tools/src/rate_report.rs b/snmalloc-tools/src/rate_report.rs new file mode 100644 index 000000000..2a34b4f39 --- /dev/null +++ b/snmalloc-tools/src/rate_report.rs @@ -0,0 +1,405 @@ +//! Streaming-mode rate reporter. +//! +//! Reads a line-oriented streaming event log emitted by an application +//! using [`snmalloc_rs::ProfilingSession`] (Phase 5.1 streaming-mode +//! API) and produces a per-site rate report: how many alloc / dealloc +//! events landed at each site, the peak live-bytes high-watermark +//! attributable to that site, and the alloc rate (events per second). +//! +//! ## Why "streaming" vs "snapshot" +//! +//! `SnMalloc::snapshot()` (in `snmalloc-rs`) materialises an in-memory +//! view of allocations that are **currently** sampled-and-live in the +//! process. That answers "what's holding memory right now?" but +//! systematically under-counts call sites whose allocations are +//! short-lived churn (allocate-and-free inside a request, scratch +//! buffers in a tight loop) -- those allocations are freed before +//! `snapshot()` is called, so they vanish from the live set. +//! +//! Streaming mode records **every** sampled event as it happens +//! (alloc, dealloc, resize), so it captures that transient churn. +//! Feeding the resulting log into this reporter answers a different +//! question: "which call site is the highest-rate allocator?" -- which +//! is what you actually want when optimising a hot path. +//! +//! ## On-disk format +//! +//! The expected log is **JSON Lines (JSONL)**: one JSON object per +//! line, UTF-8. The reporter accepts a permissive schema (extra +//! fields are ignored) and only the minimum fields are load-bearing: +//! +//! ```jsonl +//! {"ts_ns": 1000000, "kind": "alloc", "site": "0x55a0c0001000", "size": 4096} +//! {"ts_ns": 1001000, "kind": "alloc", "site": "0x55a0c0002000", "size": 256} +//! {"ts_ns": 1002000, "kind": "dealloc", "site": "0x55a0c0001000", "size": 4096} +//! ``` +//! +//! Field semantics: +//! +//! - `ts_ns` (u64, optional) -- monotonic-clock timestamp in +//! nanoseconds. Used to compute the alloc-rate denominator. If +//! any event in the log lacks `ts_ns`, rates fall back to events +//! divided by 1 second. +//! - `kind` (string, required) -- one of `"alloc"`, `"dealloc"`, +//! `"resize"`. Unknown values are skipped (forward-compat). +//! - `site` (string, required) -- the allocation site key. Typically +//! the leaf-frame address rendered as `0x` + 16 hex digits (matches +//! the `site_leaf` field emitted by other snmalloc-tools +//! subcommands), but any stable string works. +//! - `size` (u64, optional) -- bytes attributable to this event. For +//! alloc, the bytes added to the live set; for dealloc, the bytes +//! removed. Missing/zero size is treated as 0 (the row is still +//! counted in alloc/dealloc tallies but doesn't move peak-live). +//! +//! ## Streaming guarantees +//! +//! The reader is strictly stream-based: events are read one line at a +//! time through a buffered reader, and only per-site aggregates are +//! retained in memory. A 6M-event log uses memory proportional to +//! the number of distinct sites, not the number of events. + +use std::collections::HashMap; +use std::fs::File; +use std::io::{self, BufRead, BufReader, Read}; +use std::path::Path; + +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize}; + +/// One emitted row of the rate report. +/// +/// `site` is the raw site key from the log (typically a leaf-frame +/// address rendered as `0x...`); `peak_live_bytes` is the maximum +/// running-sum of `alloc_size - dealloc_size` for that site over the +/// log window. `alloc_rate_per_sec` is computed as +/// `alloc_count / (last_ts_ns - first_ts_ns) * 1e9`; if the log has +/// fewer than two timestamps or the span is zero, the field is +/// reported as `0.0` rather than NaN/inf. +#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] +pub struct RateRow { + /// Allocation-site key (leaf-frame hex, or whatever the producer + /// emitted). + pub site: String, + /// Number of `kind == "alloc"` events for this site. + pub alloc_count: u64, + /// Number of `kind == "dealloc"` events for this site. + pub dealloc_count: u64, + /// Peak live-bytes high-watermark: the maximum of the running sum + /// of alloc bytes minus dealloc bytes observed across the log. + pub peak_live_bytes: u64, + /// Alloc events per second, derived from the timestamp span of + /// the log. `0.0` if the log lacks usable timestamps. + pub alloc_rate_per_sec: f64, +} + +/// Permissive on-disk record. Every field is optional except `kind` +/// and `site` (verified during reduction). Extra fields are ignored. +#[derive(Debug, Deserialize)] +struct RawEvent { + #[serde(default)] + ts_ns: Option, + kind: Option, + site: Option, + #[serde(default)] + size: Option, +} + +/// Per-site running accumulator. Owned by the reducer; never escapes +/// the function. Keeping this off the `RateRow` keeps the public +/// output type narrow and serialisation-friendly. +#[derive(Default)] +struct SiteAcc { + alloc_count: u64, + dealloc_count: u64, + /// Running `alloc_bytes - dealloc_bytes` for this site. Saturates + /// at zero on underflow so a log that emits a dealloc before its + /// matching alloc (e.g. wraparound after a restart) doesn't panic. + live_bytes: u64, + /// Watermark of `live_bytes`. + peak_live_bytes: u64, +} + +/// Read a streaming event log file and emit per-site rate rows. +/// +/// Streams the file one line at a time -- never loads the whole log +/// into memory. Per-site aggregates are O(distinct sites), so 6M +/// events touching 1k sites consume O(1k) entries' worth of memory. +/// +/// Lines that fail to parse are skipped silently; the reader is +/// resilient to truncated tail records and to the occasional extra +/// blank line. Returns rows sorted by `alloc_count` descending, then +/// by `site` ascending for deterministic output. +pub fn read_path>(path: P) -> Result> { + let p = path.as_ref(); + let f = File::open(p) + .with_context(|| format!("opening streaming event log {}", p.display()))?; + read_reader(BufReader::new(f)) +} + +/// Same as [`read_path`] but takes any [`Read`] (used by tests with +/// in-memory fixtures and by callers piping from stdin). +pub fn read_reader(reader: R) -> Result> { + let buf = BufReader::new(reader); + reduce_lines(buf.lines()) +} + +/// Core stream-reducer: walks an iterator of `io::Result` and +/// folds per-site state. Pulled out as a standalone function so tests +/// can drive it with in-memory iterators without round-tripping +/// through a `Read`. +fn reduce_lines(lines: I) -> Result> +where + I: IntoIterator>, +{ + let mut sites: HashMap = HashMap::new(); + let mut first_ts: Option = None; + let mut last_ts: Option = None; + let mut any_ts = false; + + for line in lines { + let line = match line { + Ok(l) => l, + // I/O errors during streaming read aren't fatal -- a + // truncated tail is the common case. Stop reading; emit + // what we have. + Err(_) => break, + }; + let trimmed = line.trim(); + if trimmed.is_empty() { + continue; + } + + let raw: RawEvent = match serde_json::from_str(trimmed) { + Ok(r) => r, + // Malformed line: skip and keep going. + Err(_) => continue, + }; + + let site = match raw.site { + Some(s) if !s.is_empty() => s, + _ => continue, + }; + let kind = match raw.kind.as_deref() { + Some(k) => k, + None => continue, + }; + let size = raw.size.unwrap_or(0); + + if let Some(ts) = raw.ts_ns { + any_ts = true; + first_ts = Some(first_ts.map(|f| f.min(ts)).unwrap_or(ts)); + last_ts = Some(last_ts.map(|l| l.max(ts)).unwrap_or(ts)); + } + + let acc = sites.entry(site).or_default(); + match kind { + "alloc" => { + acc.alloc_count += 1; + acc.live_bytes = acc.live_bytes.saturating_add(size); + if acc.live_bytes > acc.peak_live_bytes { + acc.peak_live_bytes = acc.live_bytes; + } + } + "dealloc" => { + acc.dealloc_count += 1; + acc.live_bytes = acc.live_bytes.saturating_sub(size); + } + // Forward-compat: unknown kinds (incl. "resize") are + // counted as size-neutral churn -- recorded only as + // alloc-rate input via timestamp, not as a per-site delta. + // We deliberately do not bump alloc_count for resize so + // the rate denominator stays the "fresh-alloc rate", not + // "fresh-alloc + churn". + _ => {} + } + } + + // Derive seconds spanned by the log. When the span is zero or + // unknown, the rate is reported as 0.0 -- producers without + // timestamps get a clear "no rate" signal rather than infinities. + let span_sec: f64 = match (first_ts, last_ts, any_ts) { + (Some(f), Some(l), true) if l > f => (l - f) as f64 / 1_000_000_000.0, + _ => 0.0, + }; + + let mut rows: Vec = sites + .into_iter() + .map(|(site, acc)| { + let rate = if span_sec > 0.0 { + acc.alloc_count as f64 / span_sec + } else { + 0.0 + }; + RateRow { + site, + alloc_count: acc.alloc_count, + dealloc_count: acc.dealloc_count, + peak_live_bytes: acc.peak_live_bytes, + alloc_rate_per_sec: rate, + } + }) + .collect(); + + // Deterministic order: alloc_count desc, then site asc. Stable + // ordering matters for CSV/table snapshot tests and for diffing + // two reports across runs. + rows.sort_by(|a, b| { + b.alloc_count + .cmp(&a.alloc_count) + .then_with(|| a.site.cmp(&b.site)) + }); + + Ok(rows) +} + +/// Write rows in CSV format with a header line. Numeric columns are +/// rendered without thousands separators; the rate column is rendered +/// with six decimal digits (enough resolution for sub-Hz rates). +pub fn write_csv(rows: &[RateRow], w: &mut W) -> io::Result<()> { + writeln!( + w, + "site,alloc_count,dealloc_count,peak_live_bytes,alloc_rate_per_sec" + )?; + for r in rows { + writeln!( + w, + "{},{},{},{},{:.6}", + r.site, r.alloc_count, r.dealloc_count, r.peak_live_bytes, r.alloc_rate_per_sec + )?; + } + Ok(()) +} + +/// Write rows in a fixed-width pretty table (no external crate +/// dependency). Column widths are constants -- the output is +/// readable in 120-column terminals and stable across runs. +pub fn write_pretty(rows: &[RateRow], w: &mut W) -> io::Result<()> { + writeln!( + w, + "{:<20} {:>11} {:>13} {:>16} {:>20}", + "site", "alloc_count", "dealloc_count", "peak_live_bytes", "alloc_rate_per_sec" + )?; + for r in rows { + writeln!( + w, + "{:<20} {:>11} {:>13} {:>16} {:>20.6}", + r.site, r.alloc_count, r.dealloc_count, r.peak_live_bytes, r.alloc_rate_per_sec + )?; + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn lines_iter(s: &str) -> Vec> { + s.lines().map(|l| Ok(l.to_string())).collect() + } + + #[test] + fn empty_input_yields_no_rows() { + let rows = reduce_lines(lines_iter("")).unwrap(); + assert!(rows.is_empty()); + } + + #[test] + fn skips_malformed_and_blank_lines() { + let log = "\n\ + not-json\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":10}\n\ + {garbled}\n\ + \n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].site, "0xA"); + assert_eq!(rows[0].alloc_count, 1); + } + + #[test] + fn peak_live_tracks_running_max_not_final() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":100,\"ts_ns\":0}\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":50,\"ts_ns\":1}\n\ + {\"kind\":\"dealloc\",\"site\":\"0xA\",\"size\":120,\"ts_ns\":2}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows.len(), 1); + // peak is 150 (after the two allocs), even though final live=30 + assert_eq!(rows[0].peak_live_bytes, 150); + assert_eq!(rows[0].alloc_count, 2); + assert_eq!(rows[0].dealloc_count, 1); + } + + #[test] + fn rate_uses_timestamp_span() { + // Two allocs 1 second apart -> rate 2 allocs / 1s = 2.0 + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1,\"ts_ns\":0}\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1,\"ts_ns\":1000000000}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows.len(), 1); + assert!((rows[0].alloc_rate_per_sec - 2.0).abs() < 1e-9); + } + + #[test] + fn rate_is_zero_when_no_timestamps() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1}\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].alloc_rate_per_sec, 0.0); + } + + #[test] + fn sort_is_alloc_count_desc_then_site_asc() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xB\",\"size\":1}\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1}\n\ + {\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows[0].site, "0xA"); // 2 allocs wins + assert_eq!(rows[1].site, "0xB"); + } + + #[test] + fn unknown_kind_is_ignored_for_counts() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":10}\n\ + {\"kind\":\"resize\",\"site\":\"0xA\",\"size\":99}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows[0].alloc_count, 1); + // resize does not bump alloc/dealloc tallies and does not + // affect peak (we only track alloc/dealloc deltas). + assert_eq!(rows[0].peak_live_bytes, 10); + } + + #[test] + fn dealloc_underflow_saturates_at_zero() { + // Dealloc with no prior alloc -- should not panic. + let log = "{\"kind\":\"dealloc\",\"site\":\"0xA\",\"size\":999}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + assert_eq!(rows[0].peak_live_bytes, 0); + assert_eq!(rows[0].dealloc_count, 1); + } + + #[test] + fn write_csv_has_header_and_one_row_per_site() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1,\"ts_ns\":0}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + let mut out: Vec = Vec::new(); + write_csv(&rows, &mut out).unwrap(); + let s = String::from_utf8(out).unwrap(); + assert!(s.starts_with("site,alloc_count,")); + assert!(s.contains("0xA,1,0,1,")); + } + + #[test] + fn write_pretty_emits_aligned_columns() { + let log = "{\"kind\":\"alloc\",\"site\":\"0xA\",\"size\":1}\n"; + let rows = reduce_lines(lines_iter(log)).unwrap(); + let mut out: Vec = Vec::new(); + write_pretty(&rows, &mut out).unwrap(); + let s = String::from_utf8(out).unwrap(); + // Header columns separated by whitespace; site column is left- + // aligned, so the "site" header appears first. + assert!(s.lines().next().unwrap().starts_with("site")); + // Data row begins with the site key. + assert!(s.lines().nth(1).unwrap().starts_with("0xA")); + } +} diff --git a/snmalloc-tools/tests/fixtures/streaming_log_sample.jsonl b/snmalloc-tools/tests/fixtures/streaming_log_sample.jsonl new file mode 100644 index 000000000..87a176293 --- /dev/null +++ b/snmalloc-tools/tests/fixtures/streaming_log_sample.jsonl @@ -0,0 +1,8 @@ +{"ts_ns": 1000000000, "kind": "alloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 1000100000, "kind": "alloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 1000200000, "kind": "alloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 1000300000, "kind": "alloc", "site": "0x0000bbbb00000002", "size": 4096} +{"ts_ns": 1000400000, "kind": "dealloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 1500000000, "kind": "alloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 1900000000, "kind": "dealloc", "site": "0x0000aaaa00000001", "size": 1024} +{"ts_ns": 2000000000, "kind": "resize", "site": "0x0000bbbb00000002", "size": 8192} diff --git a/snmalloc-tools/tests/rate_report.rs b/snmalloc-tools/tests/rate_report.rs new file mode 100644 index 000000000..57dd404f0 --- /dev/null +++ b/snmalloc-tools/tests/rate_report.rs @@ -0,0 +1,166 @@ +//! Integration tests for the `rate-report` subcommand and its +//! library counterpart in `snmalloc_tools::rate_report`. +//! +//! The fixture under `tests/fixtures/streaming_log_sample.jsonl` +//! exercises the full event matrix: multiple sites, multiple +//! alloc/dealloc events per site, a peak-then-drop pattern, a +//! resize event (size-neutral churn), and explicit timestamps. + +use std::path::PathBuf; +use std::process::Command; + +use snmalloc_tools::rate_report::{self, RateRow}; + +fn fixture(name: &str) -> PathBuf { + let mut p = PathBuf::from(env!("CARGO_MANIFEST_DIR")); + p.push("tests"); + p.push("fixtures"); + p.push(name); + p +} + +fn row_for<'a>(rows: &'a [RateRow], site: &str) -> &'a RateRow { + rows.iter() + .find(|r| r.site == site) + .unwrap_or_else(|| panic!("no row for site {}", site)) +} + +#[test] +fn streaming_log_fixture_aggregates_per_site() { + let rows = rate_report::read_path(fixture("streaming_log_sample.jsonl")) + .expect("rate-report must read fixture"); + assert_eq!(rows.len(), 2, "fixture has two distinct sites"); + + let a = row_for(&rows, "0x0000aaaa00000001"); + // Four allocs, two deallocs in the fixture for site A. + assert_eq!(a.alloc_count, 4); + assert_eq!(a.dealloc_count, 2); + // Peak: three back-to-back 1024-byte allocs before the first + // dealloc = 3072 bytes. The fourth alloc happens after a dealloc + // has freed 1024, so it brings live to 3072 again (not higher). + assert_eq!(a.peak_live_bytes, 3072); + + let b = row_for(&rows, "0x0000bbbb00000002"); + assert_eq!(b.alloc_count, 1); + assert_eq!(b.dealloc_count, 0); + assert_eq!(b.peak_live_bytes, 4096); +} + +#[test] +fn rate_is_computed_from_timestamp_span() { + let rows = rate_report::read_path(fixture("streaming_log_sample.jsonl")).unwrap(); + // Fixture spans 1_000_000_000ns (start) to 2_000_000_000ns (end) = + // 1 second. Site A has 4 allocs in 1 second -> rate 4.0/s. + let a = row_for(&rows, "0x0000aaaa00000001"); + assert!( + (a.alloc_rate_per_sec - 4.0).abs() < 1e-9, + "expected rate ~4.0, got {}", + a.alloc_rate_per_sec + ); +} + +#[test] +fn rows_are_sorted_by_alloc_count_desc() { + let rows = rate_report::read_path(fixture("streaming_log_sample.jsonl")).unwrap(); + // Site A has 4 allocs, site B has 1: A must come first. + assert_eq!(rows[0].site, "0x0000aaaa00000001"); + assert_eq!(rows[1].site, "0x0000bbbb00000002"); +} + +#[test] +fn write_csv_round_trips_via_serde() { + let rows = rate_report::read_path(fixture("streaming_log_sample.jsonl")).unwrap(); + let mut out: Vec = Vec::new(); + rate_report::write_csv(&rows, &mut out).unwrap(); + let s = String::from_utf8(out).unwrap(); + + // Header present, two data rows present (site A then site B). + let lines: Vec<&str> = s.lines().collect(); + assert_eq!(lines.len(), 3); + assert!(lines[0].starts_with("site,alloc_count,")); + assert!(lines[1].starts_with("0x0000aaaa00000001,4,2,3072,")); + assert!(lines[2].starts_with("0x0000bbbb00000002,1,0,4096,")); +} + +#[test] +fn cli_rate_report_help_lists_subcommand() { + // The clap top-level `--help` must mention `rate-report`. Use the + // cargo-injected `CARGO_BIN_EXE_snmalloc-tools` path so this works + // regardless of the workspace layout. + let exe = env!("CARGO_BIN_EXE_snmalloc-tools"); + let out = Command::new(exe) + .arg("rate-report") + .arg("--help") + .output() + .expect("failed to spawn snmalloc-tools"); + assert!(out.status.success(), "rate-report --help should succeed"); + let stdout = String::from_utf8_lossy(&out.stdout); + assert!( + stdout.contains("--input") && stdout.contains("--pretty"), + "rate-report --help should mention --input and --pretty (got: {})", + stdout + ); +} + +#[test] +fn cli_rate_report_emits_csv_by_default() { + let exe = env!("CARGO_BIN_EXE_snmalloc-tools"); + let fixture_path = fixture("streaming_log_sample.jsonl"); + let out = Command::new(exe) + .arg("rate-report") + .arg("--input") + .arg(&fixture_path) + .output() + .expect("failed to spawn snmalloc-tools"); + assert!(out.status.success(), "rate-report should exit 0"); + let stdout = String::from_utf8_lossy(&out.stdout); + // CSV header line is the first thing emitted. + assert!( + stdout.starts_with("site,alloc_count,dealloc_count,peak_live_bytes,alloc_rate_per_sec"), + "expected CSV header, got: {}", + stdout + ); + // Both sites appear. + assert!(stdout.contains("0x0000aaaa00000001")); + assert!(stdout.contains("0x0000bbbb00000002")); +} + +#[test] +fn cli_rate_report_pretty_flag_switches_format() { + let exe = env!("CARGO_BIN_EXE_snmalloc-tools"); + let fixture_path = fixture("streaming_log_sample.jsonl"); + let out = Command::new(exe) + .arg("rate-report") + .arg("--input") + .arg(&fixture_path) + .arg("--pretty") + .output() + .expect("failed to spawn snmalloc-tools"); + assert!(out.status.success()); + let stdout = String::from_utf8_lossy(&out.stdout); + // Pretty header is whitespace-separated, not comma-separated. + assert!(!stdout.contains(','), "pretty output should not contain commas: {}", stdout); + assert!(stdout.contains("site")); + assert!(stdout.contains("alloc_count")); +} + +#[test] +fn cli_rate_report_top_truncates() { + let exe = env!("CARGO_BIN_EXE_snmalloc-tools"); + let fixture_path = fixture("streaming_log_sample.jsonl"); + let out = Command::new(exe) + .arg("rate-report") + .arg("--input") + .arg(&fixture_path) + .arg("--top") + .arg("1") + .output() + .expect("failed to spawn snmalloc-tools"); + assert!(out.status.success()); + let stdout = String::from_utf8_lossy(&out.stdout); + // Header + 1 data row = 2 lines total. + let lines: Vec<&str> = stdout.lines().collect(); + assert_eq!(lines.len(), 2, "expected --top 1 to limit to 1 row, got: {}", stdout); + // Top site by alloc-count is A. + assert!(lines[1].starts_with("0x0000aaaa00000001")); +}