Files
windmill/benchmarks/sim/pg_latency_poller.ts
pyranotaandClaude Opus 4.7 74b662d8de feat(benchmarks): k8s sim mode + util-group dashboard + reliability fixes
Stand up a minikube-backed simulation subsystem for benching Windmill under
realistic multi-node load, with a per-bench measurement pipeline and a
dashboard renderer that consolidates throughput, queue depth, per-node CPU,
PG latency/conns, OOM events, and per-node CPU-util-vs-oversaturation into
one SVG report.

Sim infrastructure (sim/):
- k8s_provisioner: minikube up + heterogeneous node sizing from topology JSON
- helm_deploy: helm install Windmill with smoke.yaml + local.yaml overlays
- image_cache: pre-load required images so bench bringup is offline-safe
- toxiproxy_k8s: per-node toxiproxy DaemonSet for cross-node latency injection
- cpu_sampler_k8s: privileged DS reading per-cgroup cpu.stat at 10Hz, dual-
  writes to stdout AND a host-mounted log file (/var/log/wm-sim-cpu-sampler/
  sampler.tsv) so heavy benches no longer lose early samples to kubelet log
  rotation
- pg_logging: ALTER SYSTEM + SIGHUP to enable verbose PG logging without restart
- pgbadger: post-bench PG log analysis HTML report
- readiness: pre-bench cluster health check (samplers stable ≥30s, workers
  ready, PG responsive, queue empty, **deploy.status rollout-complete**) —
  the rollout-complete check catches mid-rolling-update fires that previously
  starved m04's sampler under cgroup_mutex contention

Per-bench JSONL pollers, started/finalized alongside the bench loop:
- pod_timeline: 1Hz workers-per-node Ready counts (used for the workers panel)
- oom_poller: live OOM event capture (kernel + kubelet evictions + cgroup)
- pg_latency_poller: 4Hz psql \\timing on SELECT 1 vs kubectl-exec roundtrip
- pg_conn_poller: 1Hz pg_stat_activity by state (active/idle/idle_in_xact)
- node_load_poller: 2Hz /proc/loadavg + /proc/stat procs_running per node

Dashboard renderer (sim/render_report.ts + graph.ts):
- Util group: one panel per node with translucent orange oversaturation area
  BEHIND solid blue CPU-util area, 100% reference line, phase-boundary verticals.
  cols:2 grid wraps after 2 panels per row.
- PG node tinted with [PG] flag in legend across the dashboard.
- Phase-boundary verticals + push-window shaded zones layered consistently.
- All x-axes switched from wall-clock HH:MM to relative seconds-from-bench-
  start. Shared origin sourced from meta.json's bench_start_ms so 0s on every
  panel = the same wall-clock moment (previously each chart picked its own
  earliest sample as origin, causing drift between panels).

Oversaturation metric, with explicit fallback:
- Primary: (procs_running - ncpu) / ncpu × 100 — true CPU run-queue pressure.
- Fallback to load1 when procs_running is missing (older reports).
- load1 overcounted previously because it includes uninterruptible D-state
  procs (PG backends in disk I/O, cgroup_mutex waits), inflating "saturation"
  by 5-10x under load.
- Pure helper extracted to sim/util_metrics.ts; 8 unit tests cover the
  procs_running > load1 preference, the clamp-at-zero, invalid-ncpu cases.

Sampler reliability:
- HostPath log file in addition to stdout so the bench's scp-based collector
  bypasses kubelet log rotation entirely.
- main.ts truncates the host log file on every node before pushers start
  (parallel ssh, best-effort) so it doesn't grow unbounded across runs.
- Collector falls back to kubectl-logs when scp fails for any node.

Workloads (workloads/):
- io_4phase: four-phase IO step (idle → 2.5s → 500ms → 150ms jobs)
- io_150ms_flood / io_300ms_flood / io_1s_flood / io_2s_flood: single-phase
  flood configs to isolate the worker-host CFS context-switch storm vs PG
  contention regime
- burst, ops_day, cpu_*, etc. for other scenarios

Tests:
- sim/util_metrics_test.ts — 8 cases for computeOversatPct
- sim/util_panel_snapshot_test.ts — 5 assertions guarding util-panel SVG
  invariants (orange behind blue, 100% ref line, relative-time ticks NOT
  wall-clock, phase-boundary verticals, shared-origin override)

Helm values:
- sim/values/smoke.yaml — bench-tuned: workers w/ no CPU limit & low mem
  request, PG w/ 3-core request + wm-critical priorityClass + oomImmune +
  maxConnections, app w/ wm-critical + oomImmune + no resource limits.
- sim/values/local.example.yaml — template for the gitignored local.yaml
  that carries the EE license key.
- Depends on the wm-critical PriorityClass + oomImmune + maxConnections
  knobs landing in windmill-helm-charts (separate PR).

graph.ts additions:
- areaFills param: ordered list of per-kind translucent area fills drawn
  before lines, used by the util panel for orange-behind-blue layering
- lineColorOverrides: pin per-kind line colors so oversaturation reliably
  renders orange regardless of d3 ordinal-color insertion order
- highlightKindToken: substring-match flag for the PG-node tint in Node CPU
- xRelativeOriginMs: shared bench-start origin for the relative-time x-axis
- DataPointMulti is now exported for downstream tests

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-06-08 11:43:47 +02:00

106 lines
4.2 KiB
TypeScript

// PG response-time poller. Every interval (default 1s), runs a trivial
// `SELECT 1` against the bundled PG via `kubectl exec ... psql` and records
// how long it took. The point is to make PG's responsiveness visible OVER
// TIME — when the DB is healthy this is a flat line near 0; when PG gets
// contended (heavy bench load, connection storm, autovacuum stalling, etc.)
// the latency rises and you can correlate it with the throughput dip on the
// same x-axis.
//
// Output: JSONL, one line per poll
// {"ts": <ms>, "latency_ms": <number>, "ok": true|false, "err"?: "..."}
//
// Why `SELECT 1` (vs a real queue query): we want PG's *infrastructure*
// latency, not query cost. A trivial query exposes connection-pool /
// backend-fork / network round-trip cost, which is exactly what bench
// pressure inflates. Switching to a heavier query muddies "is PG fast" with
// "is this query fast".
import { MinikubeProvisioner } from "./k8s_provisioner.ts";
export type PgLatencyPoller = {
cont: { value: boolean };
done: Promise<void>;
};
export function startPgLatencyPoller(
prov: MinikubeProvisioner,
outPath: string,
opts: { intervalMs?: number; namespace?: string; pgPodSelector?: string } = {},
): PgLatencyPoller {
const intervalMs = opts.intervalMs ?? 250;
const namespace = opts.namespace ?? "default";
const selector = opts.pgPodSelector ?? "app=windmill-postgresql-demo-app";
const cont = { value: true };
const f = Deno.openSync(outPath, { write: true, create: true, truncate: true });
const enc = new TextEncoder();
// Resolve the PG pod name once at startup. If PG restarts mid-bench the
// name shouldn't change (StatefulSet), so caching is safe.
const done = (async () => {
let podName = "";
try {
const r = await prov.kubectl([
"-n", namespace, "get", "pods", "-l", selector,
"-o", "jsonpath={.items[0].metadata.name}",
]);
if (r.code === 0) podName = r.stdout.trim();
} catch (_e) { /* fall through; loop will retry */ }
while (cont.value) {
const startMs = Date.now();
let ok = false;
let err: string | undefined;
let pgQueryMs: number | undefined;
try {
if (!podName) {
// Retry pod lookup if the initial resolve failed.
const r = await prov.kubectl([
"-n", namespace, "get", "pods", "-l", selector,
"-o", "jsonpath={.items[0].metadata.name}",
]);
if (r.code === 0) podName = r.stdout.trim();
}
if (podName) {
// Run with `\timing on` so we get TWO signals:
// 1. Wall-clock around the whole kubectl exec → "kubectl
// roundtrip" (~200ms baseline = API hop + containerd exec +
// fork psql + open new pg conn). Useful as a cluster-control-
// plane health signal, not as a PG signal.
// 2. The `Time: X.XXX ms` line psql emits → pure PG query
// latency (sub-millisecond when PG is happy; rises only
// under contention/lock waits).
// Both are emitted as separate `kind` series so the chart shows
// both lines on the same axis.
const r = await prov.kubectl([
"-n", namespace, "exec", podName, "--",
"psql", "-U", "postgres", "-d", "windmill",
"-c", "\\timing on",
"-c", "SELECT 1",
]);
ok = r.code === 0;
if (!ok) err = (r.stderr || r.stdout).slice(0, 200);
// Parse psql's "Time: 0.420 ms" line for the pure PG query time.
const m = r.stdout.match(/Time:\s*([\d.]+)\s*ms/);
if (m) {
pgQueryMs = parseFloat(m[1]);
}
} else {
err = "no PG pod";
}
} catch (e) {
err = (e as Error).message;
}
const latency_ms = Date.now() - startMs;
const row = JSON.stringify({ ts: startMs, latency_ms, pg_query_ms: pgQueryMs, ok, err });
f.writeSync(enc.encode(row + "\n"));
const elapsed = Date.now() - startMs;
if (cont.value && elapsed < intervalMs) {
await new Promise((r) => setTimeout(r, intervalMs - elapsed));
}
}
f.close();
})();
return { cont, done };
}