From 58ce2e0f63e782cfd81a2f269c061ae2472e1f13 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 24 Aug 2026 02:44:29 +0000 Subject: [PATCH 1/2] chore(perf): add tokio worker balance profiler Turns Hotpath's already-present metrics server into the two numbers that distinguish a scheduling problem from a placement one: per-worker busy share, and microseconds per poll. A healthy async poll is single-digit microseconds; milliseconds mean a task ran synchronous work inside one poll, which no pool width can rebalance because a poll is not preemptible until it returns. Also attributes in-poll time to labelled futures. Hotpath's total_poll_duration_ns excludes every .await by construction, so a label with a high core count is running synchronous work on the request runtime. Reads only the metrics HTTP routes, so it needs no product code change and no restart. Threshold flags let it gate a before/after comparison. Arithmetic is covered by --self-test over fixtures, with no daemon. Co-Authored-By: Claude Opus 5 (1M context) --- scripts/lib/tokio_worker_balance.py | 373 ++++++++++++++++++++++++ scripts/profile-tokio-worker-balance.sh | 163 +++++++++++ 2 files changed, 536 insertions(+) create mode 100644 scripts/lib/tokio_worker_balance.py create mode 100755 scripts/profile-tokio-worker-balance.sh diff --git a/scripts/lib/tokio_worker_balance.py b/scripts/lib/tokio_worker_balance.py new file mode 100644 index 0000000000..56c78d8701 --- /dev/null +++ b/scripts/lib/tokio_worker_balance.py @@ -0,0 +1,373 @@ +#!/usr/bin/env python3 +"""Tokio async-worker balance analysis from Hotpath's metrics server. + +Answers one question with falsifiable numbers: is the daemon's Tokio worker +pool skewed because a few tasks execute long *synchronous* slices inside +`poll()`, or is it merely idle because the async side has little to do? + +Two independent signals are combined, both taken from the already-present +Hotpath instrumentation (no product code changes are required): + +* ``GET /tokio_runtime`` -- Tokio's own per-worker counters. ``busy_duration`` + minus ``park`` time is what the worker was *not* parked for; + ``poll_count`` divides it into per-poll occupancy. Microseconds per poll is + the discriminator: single-digit microseconds is a healthy async poll, and + milliseconds means the worker ran synchronous work inside one poll and could + not be preempted. +* ``GET /futures`` -- Hotpath's labelled-future poll accounting. + ``total_poll_duration_ns`` is wall time strictly *inside* ``Future::poll``, + so it excludes every ``.await``. Divided by the sample window it reads as + "cores of synchronous execution this label pinned to the async pool". + +Both are deltas between two samples, never lifetime totals, so a long-lived +daemon and a fresh one are comparable. + +This module is deliberately dependency-free (stdlib only) and importable, so +``--self-test`` can exercise the arithmetic over fixtures with no daemon. +""" + +from __future__ import annotations + +import json +import urllib.error +import urllib.request +from dataclasses import dataclass, field +from typing import Any, Iterable + +__all__ = [ + "WorkerDelta", + "FutureDelta", + "BalanceReport", + "fetch_json", + "worker_deltas", + "future_deltas", + "build_report", + "render_report", +] + + +@dataclass(frozen=True) +class WorkerDelta: + """One Tokio worker's counters between two samples.""" + + index: int + busy_ms: int + polls: int + steals: int + parks: int + + @property + def us_per_poll(self) -> float: + return (self.busy_ms * 1000.0 / self.polls) if self.polls else 0.0 + + +@dataclass(frozen=True) +class FutureDelta: + """One labelled future's in-poll execution between two samples.""" + + label: str + source: str + poll_ns: int + polls: int + + @property + def us_per_poll(self) -> float: + return (self.poll_ns / 1000.0 / self.polls) if self.polls else 0.0 + + def cores(self, window_secs: float) -> float: + return (self.poll_ns / 1e9 / window_secs) if window_secs > 0 else 0.0 + + +@dataclass +class BalanceReport: + window_secs: float + workers: list[WorkerDelta] = field(default_factory=list) + futures: list[FutureDelta] = field(default_factory=list) + num_workers: int = 0 + blocking_threads: int | None = None + idle_blocking_threads: int | None = None + + @property + def total_busy_ms(self) -> int: + return sum(w.busy_ms for w in self.workers) + + @property + def busy_cores(self) -> float: + if self.window_secs <= 0: + return 0.0 + return self.total_busy_ms / 1000.0 / self.window_secs + + def busy_share(self, worker: WorkerDelta) -> float: + total = self.total_busy_ms + return (100.0 * worker.busy_ms / total) if total else 0.0 + + def top_share(self, count: int) -> float: + total = self.total_busy_ms + if not total: + return 0.0 + ranked = sorted(self.workers, key=lambda w: -w.busy_ms)[:count] + return 100.0 * sum(w.busy_ms for w in ranked) / total + + @property + def active_workers(self) -> int: + return sum(1 for w in self.workers if w.busy_ms > 0) + + @property + def total_steals(self) -> int: + return sum(w.steals for w in self.workers) + + def worst_us_per_poll(self) -> float: + """Highest per-poll occupancy among workers that did real work. + + Workers with a handful of polls are excluded: one slow poll on an + otherwise idle worker is noise, not a funnel. + """ + candidates = [w for w in self.workers if w.polls >= 100] + return max((w.us_per_poll for w in candidates), default=0.0) + + +def fetch_json(host: str, port: int, path: str, timeout: float = 5.0) -> Any: + """GET one Hotpath metrics route. Returns None when it is not served.""" + url = f"http://{host}:{port}{path}" + try: + with urllib.request.urlopen(url, timeout=timeout) as response: + return json.loads(response.read().decode()) + except (urllib.error.URLError, urllib.error.HTTPError, OSError, ValueError): + return None + + +def _rows(payload: Any) -> list[dict[str, Any]]: + if isinstance(payload, dict): + payload = payload.get("data", []) + if not isinstance(payload, list): + return [] + return [row for row in payload if isinstance(row, dict)] + + +def worker_deltas(before: Any, after: Any) -> list[WorkerDelta]: + """Per-worker counter deltas between two ``/tokio_runtime`` snapshots.""" + if not isinstance(before, dict) or not isinstance(after, dict): + return [] + first = {w["index"]: w for w in _rows(before.get("workers"))} + deltas = [] + for row in _rows(after.get("workers")): + prior = first.get(row["index"]) + if prior is None: + continue + deltas.append( + WorkerDelta( + index=row["index"], + busy_ms=row["busy_duration_ms"] - prior["busy_duration_ms"], + polls=(row.get("poll_count") or 0) - (prior.get("poll_count") or 0), + steals=(row.get("steal_count") or 0) - (prior.get("steal_count") or 0), + parks=row["park_count"] - prior["park_count"], + ) + ) + deltas.sort(key=lambda w: -w.busy_ms) + return deltas + + +def future_deltas(before: Any, after: Any) -> list[FutureDelta]: + """Per-label in-poll deltas between two ``/futures`` snapshots.""" + first = {row["id"]: row for row in _rows(before) if "id" in row} + deltas = [] + for row in _rows(after): + prior = first.get(row.get("id")) + if prior is None: + continue + deltas.append( + FutureDelta( + label=row.get("label") or row.get("source", "?"), + source=row.get("source", "?"), + poll_ns=row["total_poll_duration_ns"] - prior["total_poll_duration_ns"], + polls=row["total_polls"] - prior["total_polls"], + ) + ) + deltas.sort(key=lambda f: -f.poll_ns) + return deltas + + +def build_report( + window_secs: float, + runtime_before: Any, + runtime_after: Any, + futures_before: Any, + futures_after: Any, +) -> BalanceReport: + after = runtime_after if isinstance(runtime_after, dict) else {} + return BalanceReport( + window_secs=window_secs, + workers=worker_deltas(runtime_before, runtime_after), + futures=future_deltas(futures_before, futures_after), + num_workers=after.get("num_workers", 0), + blocking_threads=after.get("num_blocking_threads"), + idle_blocking_threads=after.get("num_idle_blocking_threads"), + ) + + +def render_report(report: BalanceReport, future_limit: int = 8) -> str: + out: list[str] = [] + out.append( + f"window {report.window_secs:.1f}s workers {report.num_workers} " + f"blocking threads {report.blocking_threads} " + f"({report.idle_blocking_threads} idle)" + ) + out.append( + f"async pool busy {report.total_busy_ms / 1000.0:.2f}s " + f"= {report.busy_cores:.2f} cores steals {report.total_steals}" + ) + out.append("") + out.append( + f"{'wrk':>4} {'busy_ms':>10} {'share%':>8} {'polls':>10} " + f"{'us/poll':>10} {'steals':>8} {'parks':>9}" + ) + for worker in report.workers: + out.append( + f"{worker.index:>4} {worker.busy_ms:>10} " + f"{report.busy_share(worker):>8.2f} {worker.polls:>10} " + f"{worker.us_per_poll:>10.1f} {worker.steals:>8} {worker.parks:>9}" + ) + out.append("") + out.append(f"top-1 busy share {report.top_share(1):.2f}%") + out.append(f"top-2 busy share {report.top_share(2):.2f}%") + out.append(f"workers with any busy time {report.active_workers}/{len(report.workers)}") + out.append(f"worst us/poll (>=100 polls) {report.worst_us_per_poll():.1f}") + out.append("") + out.append("in-poll execution by labelled future (excludes every .await):") + out.append(f"{'cores':>8} {'poll_s':>9} {'polls':>9} {'us/poll':>10} label") + for entry in report.futures[:future_limit]: + if entry.poll_ns <= 0: + continue + out.append( + f"{entry.cores(report.window_secs):>8.2f} {entry.poll_ns / 1e9:>9.2f} " + f"{entry.polls:>9} {entry.us_per_poll:>10.1f} {entry.label}" + ) + return "\n".join(out) + + +def check_thresholds( + report: BalanceReport, + max_top2_share: float | None, + max_us_per_poll: float | None, +) -> list[str]: + """Return one message per breached threshold; empty means the gate passed.""" + failures: list[str] = [] + if max_top2_share is not None: + share = report.top_share(2) + if share > max_top2_share: + failures.append( + f"top-2 worker busy share {share:.2f}% exceeds {max_top2_share:.2f}%" + ) + if max_us_per_poll is not None: + worst = report.worst_us_per_poll() + if worst > max_us_per_poll: + failures.append( + f"worst worker occupancy {worst:.1f} us/poll exceeds " + f"{max_us_per_poll:.1f} us/poll" + ) + return failures + + +# -------------------------------------------------------------------------- +# Self-test: exercise the arithmetic over fixtures, with no daemon involved. +# -------------------------------------------------------------------------- + + +def _runtime_fixture(busy: Iterable[int], polls: Iterable[int]) -> dict[str, Any]: + workers = [ + { + "index": index, + "park_count": 10 * index, + "busy_duration_ms": busy_ms, + "poll_count": poll_count, + "steal_count": 0, + } + for index, (busy_ms, poll_count) in enumerate(zip(busy, polls)) + ] + return { + "num_workers": len(workers), + "num_blocking_threads": 6, + "num_idle_blocking_threads": 6, + "workers": workers, + } + + +def self_test() -> int: + failures: list[str] = [] + + def check(name: str, condition: bool) -> None: + if not condition: + failures.append(name) + + zero = _runtime_fixture([0, 0, 0, 0], [0, 0, 0, 0]) + skewed = _runtime_fixture([9000, 900, 50, 0], [9000, 900, 50, 0]) + report = build_report(10.0, zero, skewed, [], []) + + check("four worker deltas", len(report.workers) == 4) + check("busy total", report.total_busy_ms == 9950) + check("busy cores", abs(report.busy_cores - 0.995) < 1e-6) + check("top1 share", abs(report.top_share(1) - 90.452) < 0.01) + check("top2 share", abs(report.top_share(2) - 99.497) < 0.01) + check("active workers", report.active_workers == 3) + # 9000 ms over 9000 polls is exactly 1000 us per poll. + check("us per poll", abs(report.workers[0].us_per_poll - 1000.0) < 1e-6) + + # A worker with very few polls must not set the funnel verdict. + noisy = _runtime_fixture([5, 0, 0, 0], [1, 0, 0, 0]) + noisy_report = build_report(10.0, zero, noisy, [], []) + check("low-poll worker ignored", noisy_report.worst_us_per_poll() == 0.0) + + before = [ + {"id": 1, "label": "a", "source": "s1", "total_polls": 0, "total_poll_duration_ns": 0}, + {"id": 2, "label": "b", "source": "s2", "total_polls": 0, "total_poll_duration_ns": 0}, + ] + after = [ + { + "id": 1, + "label": "a", + "source": "s1", + "total_polls": 1000, + "total_poll_duration_ns": 2_000_000_000, + }, + { + "id": 2, + "label": "b", + "source": "s2", + "total_polls": 1000, + "total_poll_duration_ns": 1_000_000, + }, + ] + futures = future_deltas(before, after) + check("futures ranked by poll time", [f.label for f in futures] == ["a", "b"]) + check("future cores", abs(futures[0].cores(10.0) - 0.2) < 1e-9) + check("future us per poll", abs(futures[0].us_per_poll - 2000.0) < 1e-6) + + # Ids absent from the earlier snapshot cannot produce a delta. + check("unmatched ids dropped", future_deltas([], after) == []) + # Envelope form ({"data": [...]}) must parse the same as a bare list. + check("envelope form", len(future_deltas({"data": before}, {"data": after})) == 2) + + check( + "threshold gate fails on skew", + check_thresholds(report, max_top2_share=60.0, max_us_per_poll=None) != [], + ) + check( + "threshold gate passes when balanced", + check_thresholds(report, max_top2_share=100.0, max_us_per_poll=None) == [], + ) + check( + "threshold gate fails on slow polls", + check_thresholds(report, max_top2_share=None, max_us_per_poll=100.0) != [], + ) + + # A report with no samples must be inert rather than divide by zero. + empty = build_report(0.0, None, None, None, None) + check("empty report is inert", empty.busy_cores == 0.0 and empty.top_share(2) == 0.0) + check("empty renders", "window 0.0s" in render_report(empty)) + + if failures: + for name in failures: + print(f"FAIL {name}") + return 1 + print("tokio_worker_balance self-test: ok") + return 0 diff --git a/scripts/profile-tokio-worker-balance.sh b/scripts/profile-tokio-worker-balance.sh new file mode 100755 index 0000000000..83166f59db --- /dev/null +++ b/scripts/profile-tokio-worker-balance.sh @@ -0,0 +1,163 @@ +#!/usr/bin/env bash +# Measure how a running TraceDecay process spreads work across its Tokio +# async workers, and which labelled futures pin those workers. +# +# Reads only Hotpath's metrics server, so it needs no product code change and +# no restart: build the binary with `hotpath` (see +# .claude/skills/using-hotpath) and point this at the live port. +# +# The two numbers this exists to produce: +# +# * per-worker busy share -- how lopsided the pool is; +# * microseconds per poll -- WHY. A healthy async poll is single-digit +# microseconds. Milliseconds per poll means a task ran synchronous work +# inside one `poll()`, which no amount of work-stealing can rebalance +# because the worker is not preemptible until poll returns. +# +# Both are deltas across the sample window, so a long-lived daemon and a fresh +# one are directly comparable. `--futures` attributes the in-poll time to +# labelled futures; `total_poll_duration_ns` excludes every `.await`, so a +# label with a high core count is running synchronous work on the async pool. +# +# Sample a daemon for 60s: +# scripts/profile-tokio-worker-balance.sh --seconds 60 +# +# Gate a before/after comparison (non-zero exit when breached): +# scripts/profile-tokio-worker-balance.sh --seconds 60 \ +# --max-top2-share 60 --max-us-per-poll 200 +# +# Hermetic arithmetic tests (no daemon, no cargo): +# scripts/profile-tokio-worker-balance.sh --self-test +set -euo pipefail + +SCRIPT_DIR=$(CDPATH= cd -P -- "$(dirname -- "$0")" && pwd) + +usage() { + sed -n '2,30p' "$0" >&2 + exit 2 +} + +HOST=127.0.0.1 +PORT=6770 +SECONDS_TO_SAMPLE=60 +JSON_OUT= +MAX_TOP2= +MAX_US_PER_POLL= +SELF_TEST=0 + +while [ $# -gt 0 ]; do + case "$1" in + --host) HOST=${2:?--host needs a value}; shift 2 ;; + --port) PORT=${2:?--port needs a value}; shift 2 ;; + --seconds) SECONDS_TO_SAMPLE=${2:?--seconds needs a value}; shift 2 ;; + --json) JSON_OUT=${2:?--json needs a path}; shift 2 ;; + --max-top2-share) MAX_TOP2=${2:?--max-top2-share needs a value}; shift 2 ;; + --max-us-per-poll) MAX_US_PER_POLL=${2:?--max-us-per-poll needs a value}; shift 2 ;; + --self-test) SELF_TEST=1; shift ;; + -h|--help) usage ;; + *) echo "unknown argument: $1" >&2; usage ;; + esac +done + +PYTHON=${PYTHON:-python3} + +if [ "$SELF_TEST" = 1 ]; then + exec "$PYTHON" -c " +import sys +sys.path.insert(0, '$SCRIPT_DIR/lib') +import tokio_worker_balance as twb +raise SystemExit(twb.self_test()) +" +fi + +exec "$PYTHON" - "$HOST" "$PORT" "$SECONDS_TO_SAMPLE" "$JSON_OUT" \ + "$MAX_TOP2" "$MAX_US_PER_POLL" < 0 + ], + }, + handle, + indent=2, + ) + print(f"\nwrote {json_out}") + +failures = twb.check_thresholds( + report, + float(max_top2) if max_top2 else None, + float(max_us) if max_us else None, +) +if failures: + print("") + for failure in failures: + print(f"THRESHOLD BREACHED: {failure}") + raise SystemExit(1) +PY From 05c815261811a198e8a8de6f2d60d30380c9952a Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 24 Aug 2026 02:44:53 +0000 Subject: [PATCH 2/2] docs(perf): record the async-pool in-poll breach Principle 2 claims long CPU slices never run on the request runtime's workers. Measurement contradicts it for the historical half of the session-refresh sweep: over a 60s window one labelled future, session_temporal_refresh.history, accounts for 89.68s of the async pool's 91.92s of busy time -- 97.6% of everything the workers did was synchronous execution, at 2.5ms per poll and 8.9ms on the worst worker. Records why the resulting 5-of-16 worker skew must not be answered by resizing the pool. Runnable async concurrency is 2 (one project and one profile ingestor), so there is nothing for work-stealing to take, and the eleven idle workers are precisely why the breach has not surfaced as a latency regression: an arriving request never queues behind a 9ms poll. Narrowing the pool would delete that headroom. Also records the instrument that cannot answer the remaining question: Tokio names async workers and blocking threads with one thread_name_fn, so perf --comms cannot split the two pools, and naming the exact in-poll leaf needs instrumentation inside the pass instead. Co-Authored-By: Claude Opus 5 (1M context) --- docs/SERVING-PATH-PERFORMANCE.md | 85 ++++++++++++++++++++++++++++++++ 1 file changed, 85 insertions(+) diff --git a/docs/SERVING-PATH-PERFORMANCE.md b/docs/SERVING-PATH-PERFORMANCE.md index 8b659f777c..2c1ca23841 100644 --- a/docs/SERVING-PATH-PERFORMANCE.md +++ b/docs/SERVING-PATH-PERFORMANCE.md @@ -199,6 +199,10 @@ have no finish line, so they stay paced: them defer or shrink a slice while agents are actively querying. Results are identical; only pacing changes. +The historical half of the session-refresh sweep does none of this today; see +"Open breach: historical transcript ingest violates Principle 2" below for the +measurement and for why the answer is not a wider or narrower worker pool. + ### 3. Hash where data is born, never where it is served Content digests (payload references, canonical record digests) are computed @@ -340,6 +344,87 @@ cost, not indexing interference. Peak daemon RSS on this host is ~13.6GB with no indexing at all, well over the 6GB gate budget — a live, separate breach of Principle 5 that this measurement did not introduce. +## Open breach: historical transcript ingest violates Principle 2 + +Principle 2 says long CPU slices "run in `spawn_blocking` chunks, never on the +request runtime's workers for unbounded stretches", and names projection +refresh as an open-ended sweep that must stay paced. Measurement contradicts +both for the *historical* half of that sweep. + +Reproduce with `scripts/profile-tokio-worker-balance.sh`, which reads Hotpath's +metrics server and needs no product code change. The numbers below are one +60-second window on a 96-core host, `--profile perf`, features +`production,hotpath,hotpath-mcp`, `--cfg tokio_unstable`, against an isolated +profile root under `/tmp` with one registered 3,114-file corpus: + +``` +window 60.0s workers 16 blocking threads 8 (7 idle) +async pool busy 91.92s = 1.53 cores steals 14 + + wrk busy_ms share% polls us/poll steals parks + 14 61824 67.26 6926 8926.4 2 13084 + 12 13398 14.58 10979 1220.3 2 21509 + 6 8814 9.59 7343 1200.3 2 15314 + 15 5892 6.41 9293 634.0 4 17940 + 10 1990 2.16 3452 576.5 4 6424 + (workers 0-5, 7-9, 11, 13 ran zero polls) + +top-2 busy share 81.84% +workers with any busy time 5/16 + +in-poll execution by labelled future (excludes every .await): + cores poll_s polls us/poll label + 1.49 89.68 35648 2515.7 session_temporal_refresh.history +``` + +`total_poll_duration_ns` is wall time strictly inside `Future::poll`, so an +`.await` cannot contribute to it. One label therefore accounts for 89.68s of +the pool's 91.92s of busy time — **97.6% of everything the request runtime's +async workers did was synchronous execution belonging to +`session_temporal_refresh.history`**, at 2.5ms per poll and a worst-worker +occupancy of 8.9ms per poll. A 325-second run of the same scenario shows the +same shape: 5 of 16 workers carry 98.5% of busy time, the reserved indexing +pool consumes 4.22 CPU-seconds over the whole run, and cumulative steal counts +stay in single digits. + +### The skew is a symptom; do not resize the pool + +The concentration is not a scheduling defect and it is not fixed by changing +`async_worker_threads()`. `session_history_refresh` runs as exactly two +long-lived tasks — one project-scoped ingestor and one profile-scoped one — so +runnable async concurrency is 2, not 16. Work-stealing has nothing to steal +(14 steals in 60s), and no pool width redistributes a slice that a worker is +already inside: a poll is not preemptible until it returns. + +The width is in fact why the breach has not surfaced as a latency regression. +Eleven of sixteen workers ran zero polls in the sampled minute, so an arriving +request always finds a parked worker rather than queueing behind a 9ms poll. +Narrowing the pool to "match" the observed concurrency would delete exactly +that headroom and convert a benign skew into a real tail. The measurement to +watch is microseconds per poll, not busy share. + +### Where the fix has to land + +The cost is CPU placement, not balance: this work belongs off the request +runtime, in the same sense the indexing pool already is. Fixing it means +restructuring the ingest pass, not adding a yield — the pass is an async +pipeline whose synchronous slices sit *between* store awaits, so a yield point +redistributes the slices across workers without moving a single cycle off the +serving pool, and `spawn_blocking` cannot wrap the pass as a whole because it +is a future, not a closure. The work has to be batched into blocking chunks at +a boundary that owns its data. + +Leaf attribution is not yet pinned down, and the obvious instrument cannot pin +it: Tokio names async worker threads and blocking-pool threads with the same +`thread_name_fn`, so `perf --comms` cannot separate the two pools. Under that +combined thread-name family the largest single product symbol during ingest is +`privacy::rules::contains_ignore_ascii_case` (5.16% of whole-process cycles, +~7.4% for the privacy-detector family), which is a plausible in-poll leaf +because redaction is pure CPU over transcript text — but the SQLite parser +symbols in the same family belong to `ReadSnapshot::query`, which already uses +`spawn_blocking` and is correctly placed. Naming the exact leaf needs a +`#[hotpath::measure]` inside the pass or `HOTPATH_FOCUS`, not a thread filter. + ## Open breach: the vector-generation store violates Principle 5 `DatabaseVectorGenerationStoreV1` persists its entire state as one JSON blob in