Snapshot of verae-nats-cluster (optimal NATS config study)
This commit is contained in:
commit
8639ca27ee
184 changed files with 10626 additions and 0 deletions
192
scripts/bench-report.py
Executable file
192
scripts/bench-report.py
Executable file
|
|
@ -0,0 +1,192 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Turn nats bench text logs + latency JSON into markdown."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
def fmt_int(s: str | None) -> str:
|
||||
if not s:
|
||||
return "—"
|
||||
return f"{int(s):,}"
|
||||
|
||||
|
||||
def parse_bench(text: str) -> dict[str, str]:
|
||||
out: dict[str, str] = {"kind": "throughput"}
|
||||
def _rate(pattern: str, msgs_key: str, mb_key: str) -> None:
|
||||
m = re.search(pattern, text)
|
||||
if not m:
|
||||
return
|
||||
out[msgs_key] = m.group(1).replace(",", "")
|
||||
n = float(m.group(2))
|
||||
unit = m.group(3).upper()
|
||||
if unit == "KB":
|
||||
n = n / 1024.0
|
||||
elif unit == "GB":
|
||||
n = n * 1024.0
|
||||
out[mb_key] = f"{n:.2f}"
|
||||
|
||||
_rate(
|
||||
r"(?m)^\s*Pub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
|
||||
"pub_msgs",
|
||||
"pub_mb",
|
||||
)
|
||||
_rate(
|
||||
r"(?m)^\s*Sub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
|
||||
"sub_msgs",
|
||||
"sub_mb",
|
||||
)
|
||||
_rate(
|
||||
r"NATS Pub/Sub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
|
||||
"agg_msgs",
|
||||
"agg_mb",
|
||||
)
|
||||
if "JetStream" in text or "--js" in text or "js-" in text:
|
||||
out["mode"] = "jetstream r=3 file"
|
||||
else:
|
||||
out["mode"] = "core pub/sub"
|
||||
# nats 0.1.6 prints min/avg/max as msgs/sec across publishers, not µs delay
|
||||
m = re.search(
|
||||
r"min\s+([0-9,]+)\s*\|\s*avg\s+([0-9,]+)\s*\|\s*max\s+([0-9,]+)\s*\|\s*stddev\s+([0-9,]+)\s*msgs",
|
||||
text,
|
||||
)
|
||||
if m:
|
||||
out["pub_spread"] = f"{m.group(1)}–{m.group(3)} (avg {m.group(2)})"
|
||||
return out
|
||||
|
||||
|
||||
def parse_lat(text: str) -> dict[str, str] | None:
|
||||
for line in text.splitlines():
|
||||
line = line.strip()
|
||||
if line.startswith("{") and "p99_us" in line:
|
||||
d = json.loads(line)
|
||||
mode = d.get("mode") or ""
|
||||
return {
|
||||
"kind": "latency",
|
||||
"count": str(d.get("count", "")),
|
||||
"pubs": str(d.get("pubs", "")),
|
||||
"size": str(d.get("size", "")),
|
||||
"mode": str(mode),
|
||||
"min": d.get("min", ""),
|
||||
"avg": d.get("avg", ""),
|
||||
"p50": d.get("p50", ""),
|
||||
"p90": d.get("p90", ""),
|
||||
"p99": d.get("p99", ""),
|
||||
"max": d.get("max", ""),
|
||||
}
|
||||
return None
|
||||
|
||||
|
||||
def lat_mode(run: str, recorded: str) -> str:
|
||||
if recorded in ("ping", "flood"):
|
||||
return recorded
|
||||
if "ping" in run:
|
||||
return "ping"
|
||||
return "flood"
|
||||
|
||||
|
||||
def thru_sort(p: dict[str, str]) -> tuple:
|
||||
return (0 if p.get("mode", "").startswith("core") else 1, p.get("run", ""))
|
||||
|
||||
|
||||
def lat_sort(p: dict[str, str]) -> tuple:
|
||||
mode = lat_mode(p.get("run", ""), p.get("mode", ""))
|
||||
return (0 if mode == "ping" else 1, int(p.get("count") or 0), p.get("run", ""))
|
||||
|
||||
|
||||
def main() -> int:
|
||||
folder = Path(sys.argv[1] if len(sys.argv) > 1 else ".")
|
||||
thru: list[dict[str, str]] = []
|
||||
lats: list[dict[str, str]] = []
|
||||
for f in sorted(folder.glob("*.txt")):
|
||||
text = f.read_text(encoding="utf-8", errors="replace")
|
||||
lat = parse_lat(text)
|
||||
if lat:
|
||||
lat["run"] = f.stem
|
||||
lats.append(lat)
|
||||
continue
|
||||
p = parse_bench(text)
|
||||
if p.get("pub_msgs") or p.get("agg_msgs"):
|
||||
p["run"] = f.stem
|
||||
thru.append(p)
|
||||
|
||||
stamp = folder.name if re.fullmatch(r"\d{8}T\d{6}Z", folder.name) else ""
|
||||
|
||||
print("# NATS cluster message speed")
|
||||
print()
|
||||
if stamp:
|
||||
print(f"Run **`{stamp}`** (UTC). ", end="")
|
||||
print(
|
||||
"Client: LXC **510** `verae-px-worker` (`10.10.10.20`), not a nats-* server. "
|
||||
"Servers: `nats-a/b/c` on `10.10.10.21–23` (`vmbr1` only)."
|
||||
)
|
||||
print()
|
||||
print("Client URL:")
|
||||
print()
|
||||
print("```text")
|
||||
print("nats://10.10.10.21:4222,nats://10.10.10.22:4222,nats://10.10.10.23:4222")
|
||||
print("```")
|
||||
print()
|
||||
print("## Method")
|
||||
print()
|
||||
print("- **Core NATS** is fire-and-forget pub/sub (`nats bench`). No disk, no replica ack.")
|
||||
print(
|
||||
"- **JetStream** uses **file** storage and **replicas=3** (same as product streams). "
|
||||
"The unique stream `benchstream` is deleted between JS loads."
|
||||
)
|
||||
print("- Throughput is **msgs/sec** from nats CLI **0.1.6** (`--no-progress --csv`). Its min/avg/max are publisher **rate spread**, not delay.")
|
||||
print(
|
||||
"- **Ping** delay: one publisher, sequential publish-then-wait. This is one-message round-trip through the cluster."
|
||||
)
|
||||
print(
|
||||
"- **Flood** delay: N publishers dump the whole batch, then the subscriber drains. "
|
||||
"This is **queueing under burst**, not wire RTT."
|
||||
)
|
||||
print("- Probe: `scripts/latency.mjs` (two connections, header timestamp).")
|
||||
print()
|
||||
print("## Throughput")
|
||||
print()
|
||||
print("| Run | Mode | Aggregate msgs/s | Pub msgs/s | Pub MB/s | Sub msgs/s | Sub MB/s |")
|
||||
print("|-----|------|------------------|------------|----------|------------|----------|")
|
||||
for p in sorted(thru, key=thru_sort):
|
||||
print(
|
||||
f"| `{p['run']}` | {p.get('mode', '')} | {fmt_int(p.get('agg_msgs'))} | "
|
||||
f"{fmt_int(p.get('pub_msgs'))} | {p.get('pub_mb') or '—'} | "
|
||||
f"{fmt_int(p.get('sub_msgs'))} | {p.get('sub_mb') or '—'} |"
|
||||
)
|
||||
print()
|
||||
print("## Round-trip delay")
|
||||
print()
|
||||
print("| Run | Kind | Count | Pubs | Size | min | avg | p50 | p90 | p99 | max |")
|
||||
print("|-----|------|-------|------|------|-----|-----|-----|-----|-----|-----|")
|
||||
for p in sorted(lats, key=lat_sort):
|
||||
kind = lat_mode(p.get("run", ""), p.get("mode", ""))
|
||||
label = "ping (sequential RTT)" if kind == "ping" else "flood (burst queueing)"
|
||||
print(
|
||||
f"| `{p['run']}` | {label} | {p.get('count', '')} | {p.get('pubs', '')} | "
|
||||
f"{p.get('size', '')} B | {p.get('min', '')} | {p.get('avg', '')} | "
|
||||
f"{p.get('p50', '')} | {p.get('p90', '')} | {p.get('p99', '')} | {p.get('max', '')} |"
|
||||
)
|
||||
print()
|
||||
print("## What the numbers mean")
|
||||
print()
|
||||
print(
|
||||
"Product job/event/archive traffic is **JetStream r=3 file**. On this three-LXC stand that is about "
|
||||
"**16k durable 128 B pubs/s** (about **13k** at 1 KiB). Core NATS is an upper bound for "
|
||||
"non-durable fan-out: about **0.7–2.0M msgs/s** aggregate at 128 B, or **~630k msgs/s (~616 MB/s)** at 1 KiB with 4 publishers."
|
||||
)
|
||||
print()
|
||||
print(
|
||||
"A quiet request-reply is **~0.3 ms** average, **p99 < 1 ms**. Flood rows in the **150–500 ms** band "
|
||||
"are the subscriber catching up after a burst, which is what a job-events mailbox sees if publishers outrun consumers."
|
||||
)
|
||||
print()
|
||||
print("Re-run on NS1: `bash scripts/bench.sh`. Raw logs/CSVs are under `results/<utc>/`.")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
135
scripts/bench.sh
Executable file
135
scripts/bench.sh
Executable file
|
|
@ -0,0 +1,135 @@
|
|||
#!/usr/bin/env bash
|
||||
# Message throughput and delay ladder against the 3-node vmbr1 cluster.
|
||||
# Prefers a client that is not a nats-* server (px-worker LXC 510).
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/client.env"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:/usr/local/bin:$PATH"
|
||||
STAMP="$(date -u +%Y%m%dT%H%M%SZ)"
|
||||
OUT="${BENCH_OUT:-$ROOT/results/$STAMP}"
|
||||
CLIENT_VMID="${CLIENT_VMID:-510}"
|
||||
mkdir -p "$OUT"
|
||||
|
||||
ensure_nats_cli() {
|
||||
local vmid="$1"
|
||||
sudo pct exec "$vmid" -- bash -lc '
|
||||
set -e
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
if [[ ! -x /usr/local/bin/nats ]]; then
|
||||
apt-get install -y --no-install-recommends unzip curl ca-certificates >/dev/null
|
||||
curl -fsSL https://github.com/nats-io/natscli/releases/download/v0.1.6/nats-0.1.6-linux-amd64.zip -o /tmp/natscli.zip
|
||||
rm -rf /tmp/natscli && mkdir -p /tmp/natscli
|
||||
unzip -o /tmp/natscli.zip -d /tmp/natscli >/dev/null
|
||||
BIN=$(find /tmp/natscli -type f -name nats | head -1)
|
||||
install -m 0755 "$BIN" /usr/local/bin/nats
|
||||
fi
|
||||
nats --version
|
||||
'
|
||||
}
|
||||
|
||||
run_one() {
|
||||
local name="$1"
|
||||
shift
|
||||
echo "=== $name ===" | tee "$OUT/$name.txt"
|
||||
# nats bench writes csv itself when --csv is a path inside the guest
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -lc "
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
export NATS_URL='$NATS_URL'
|
||||
nats bench --no-progress --csv=/tmp/bench.csv $*
|
||||
" | tee -a "$OUT/$name.txt"
|
||||
sudo pct exec "$CLIENT_VMID" -- cat /tmp/bench.csv >"$OUT/$name.csv" || true
|
||||
}
|
||||
|
||||
echo "client LXC $CLIENT_VMID NATS_URL=$NATS_URL out=$OUT"
|
||||
ensure_nats_cli "$CLIENT_VMID"
|
||||
|
||||
# Core NATS pub/sub — increasing publishers (same 128 B payload)
|
||||
run_one core-1p1s-50k-128 bench.core.a --pub 1 --sub 1 --msgs 50000 --size 128
|
||||
run_one core-4p4s-100k-128 bench.core.b --pub 4 --sub 4 --msgs 100000 --size 128
|
||||
run_one core-8p8s-200k-128 bench.core.c --pub 8 --sub 8 --msgs 200000 --size 128
|
||||
run_one core-4p4s-50k-1k bench.core.d --pub 4 --sub 4 --msgs 50000 --size 1024
|
||||
|
||||
js_rm() {
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -lc "
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
export NATS_URL='$NATS_URL'
|
||||
nats stream rm benchstream --force >/dev/null 2>&1 || true
|
||||
"
|
||||
}
|
||||
# JetStream file store, replicas=3 (matches product streams)
|
||||
js_rm
|
||||
run_one js-1p-20k-128-r3 bench.js.a --js --purge --pub 1 --msgs 20000 --size 128 --replicas 3 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-4p-50k-128-r3 bench.js.b --js --purge --pub 4 --msgs 50000 --size 128 --replicas 3 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-4p-20k-1k-r3 bench.js.c --js --purge --pub 4 --msgs 20000 --size 1024 --replicas 3 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-2p2s-20k-128-r3 bench.js.d --js --purge --pub 2 --sub 2 --msgs 20000 --size 128 --replicas 3 --storage file --maxbytes=512MB --pull --stream=benchstream
|
||||
js_rm
|
||||
|
||||
# Optional native memory-store ladder (same replica count). Used by maximize-ns1-study.sh.
|
||||
if [[ "${JS_EXTRA_MEMORY:-0}" == "1" ]]; then
|
||||
js_rm
|
||||
run_one js-mem-1p-20k-128-r3 bench.js.m1 --js --purge --pub 1 --msgs 20000 --size 128 --replicas 3 --storage memory --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-mem-4p-50k-128-r3 bench.js.m2 --js --purge --pub 4 --msgs 50000 --size 128 --replicas 3 --storage memory --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-mem-4p-20k-1k-r3 bench.js.m3 --js --purge --pub 4 --msgs 20000 --size 1024 --replicas 3 --storage memory --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
fi
|
||||
|
||||
# Factorial extras: replicas=1 vs 3, file vs memory (does not touch product streams).
|
||||
if [[ "${EXHAUSTIVE:-0}" == "1" ]]; then
|
||||
js_rm
|
||||
run_one js-file-1p-20k-128-r1 bench.js.e1 --js --purge --pub 1 --msgs 20000 --size 128 --replicas 1 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-file-4p-50k-128-r1 bench.js.e2 --js --purge --pub 4 --msgs 50000 --size 128 --replicas 1 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-mem-1p-20k-128-r1 bench.js.e3 --js --purge --pub 1 --msgs 20000 --size 128 --replicas 1 --storage memory --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-mem-4p-50k-128-r1 bench.js.e4 --js --purge --pub 4 --msgs 50000 --size 128 --replicas 1 --storage memory --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
run_one js-file-1p-20k-4k-r3 bench.js.e5 --js --purge --pub 1 --msgs 20000 --size 4096 --replicas 3 --storage file --maxbytes=512MB --stream=benchstream
|
||||
js_rm
|
||||
fi
|
||||
|
||||
# Round-trip delay (two connections, through the cluster) at several loads
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -lc "
|
||||
set -e
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
export NATS_URL='$NATS_URL'
|
||||
mkdir -p /tmp/nats-lat
|
||||
cd /tmp/nats-lat
|
||||
if [[ ! -d node_modules/nats ]]; then
|
||||
npm init -y >/dev/null
|
||||
npm install --no-audit --no-fund nats@2 >/dev/null
|
||||
fi
|
||||
" >/dev/null
|
||||
sudo pct push "$CLIENT_VMID" "$ROOT/scripts/latency.mjs" /tmp/nats-lat/latency.mjs || \
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -c 'cat > /tmp/nats-lat/latency.mjs' < "$ROOT/scripts/latency.mjs"
|
||||
lat() {
|
||||
local name="$1" n="$2" sz="$3" p="$4" mode="${5:-flood}"
|
||||
echo "=== $name ===" | tee "$OUT/$name.txt"
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -lc "
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
export NATS_URL='$NATS_URL'
|
||||
cd /tmp/nats-lat
|
||||
node latency.mjs $n $sz $p $mode
|
||||
" | tee -a "$OUT/$name.txt"
|
||||
}
|
||||
# copy latest probe
|
||||
sudo pct exec "$CLIENT_VMID" -- bash -c 'cat > /tmp/nats-lat/latency.mjs' < "$ROOT/scripts/latency.mjs"
|
||||
lat lat-ping-1k-128 1000 128 1 ping
|
||||
if [[ "${EXHAUSTIVE:-0}" == "1" ]]; then
|
||||
lat lat-reconnect-200-128 200 128 1 reconnect
|
||||
fi
|
||||
lat lat-1p-5k-128 5000 128 1 flood
|
||||
lat lat-4p-10k-128 10000 128 4 flood
|
||||
lat lat-8p-20k-128 20000 128 8 flood
|
||||
lat lat-4p-5k-1k 5000 1024 4 flood
|
||||
|
||||
python3 "$ROOT/scripts/bench-report.py" "$OUT" >"$OUT/BENCH.md"
|
||||
cp "$OUT/BENCH.md" "$ROOT/BENCH.md"
|
||||
echo "wrote $OUT/BENCH.md and $ROOT/BENCH.md"
|
||||
571
scripts/build-ns1-study-report.py
Executable file
571
scripts/build-ns1-study-report.py
Executable file
|
|
@ -0,0 +1,571 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Build the NS1-host study report (charts + markdown + HTML + PDF) from a results dir.
|
||||
|
||||
Must be able to run entirely on NS1.GEORGELAMBERT.ORG with python3, matplotlib,
|
||||
pandoc, and weasyprint. Parses nats bench logs; does not hard-code rates.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
import importlib.util
|
||||
|
||||
_spec = importlib.util.spec_from_file_location(
|
||||
"bench_report", Path(__file__).resolve().parent / "bench-report.py"
|
||||
)
|
||||
_br = importlib.util.module_from_spec(_spec)
|
||||
assert _spec.loader is not None
|
||||
_spec.loader.exec_module(_br)
|
||||
fmt_int = _br.fmt_int
|
||||
lat_mode = _br.lat_mode
|
||||
lat_sort = _br.lat_sort
|
||||
parse_bench = _br.parse_bench
|
||||
parse_lat = _br.parse_lat
|
||||
thru_sort = _br.thru_sort
|
||||
|
||||
try:
|
||||
import matplotlib
|
||||
|
||||
matplotlib.use("Agg")
|
||||
import matplotlib.pyplot as plt
|
||||
from matplotlib.ticker import FuncFormatter
|
||||
except ImportError as e:
|
||||
raise SystemExit(f"matplotlib required on NS1: {e}") from e
|
||||
|
||||
INDIGO = "#4f46e5"
|
||||
DEEP = "#312e81"
|
||||
TEAL = "#047857"
|
||||
AMBER = "#b45309"
|
||||
LILAC = "#7c74f0"
|
||||
INK = "#171a26"
|
||||
MUTED = "#5b6178"
|
||||
GRID = "#d9dce8"
|
||||
|
||||
CORE_LABELS = {
|
||||
"core-1p1s-50k-128": "1p1s\n50k×128 B",
|
||||
"core-4p4s-100k-128": "4p4s\n100k×128 B",
|
||||
"core-8p8s-200k-128": "8p8s\n200k×128 B",
|
||||
"core-4p4s-50k-1k": "4p4s\n50k×1 KiB",
|
||||
}
|
||||
JS_LABELS = {
|
||||
"js-1p-20k-128-r3": "1p 20k×128 B",
|
||||
"js-4p-50k-128-r3": "4p 50k×128 B",
|
||||
"js-4p-20k-1k-r3": "4p 20k×1 KiB",
|
||||
"js-2p2s-20k-128-r3": "2p2s pull 20k×128 B",
|
||||
"js-mem-1p-20k-128-r3": "mem 1p 128 B",
|
||||
"js-mem-4p-50k-128-r3": "mem 4p 128 B",
|
||||
"js-mem-4p-20k-1k-r3": "mem 4p 1 KiB",
|
||||
}
|
||||
LAT_LABELS = {
|
||||
"lat-ping-1k-128": "Ping\n1k×128 B",
|
||||
"lat-1p-5k-128": "Flood 1p\n5k×128 B",
|
||||
"lat-4p-5k-1k": "Flood 4p\n5k×1 KiB",
|
||||
"lat-4p-10k-128": "Flood 4p\n10k×128 B",
|
||||
"lat-8p-20k-128": "Flood 8p\n20k×128 B",
|
||||
}
|
||||
|
||||
|
||||
def ms(s: str) -> float:
|
||||
return float(s.replace("ms", "").replace(",", "").strip())
|
||||
|
||||
|
||||
def k_fmt(x: float, _pos: int | None = None) -> str:
|
||||
if x >= 1_000_000:
|
||||
return f"{x / 1_000_000:.2f}M"
|
||||
if x >= 1000:
|
||||
return f"{x / 1000:.0f}k"
|
||||
return f"{x:.0f}"
|
||||
|
||||
|
||||
def style() -> None:
|
||||
plt.rcParams.update(
|
||||
{
|
||||
"font.family": "sans-serif",
|
||||
"font.size": 10,
|
||||
"axes.titlesize": 12,
|
||||
"axes.titleweight": "semibold",
|
||||
"axes.edgecolor": GRID,
|
||||
"axes.labelcolor": INK,
|
||||
"text.color": INK,
|
||||
"xtick.color": MUTED,
|
||||
"ytick.color": MUTED,
|
||||
"figure.facecolor": "white",
|
||||
"axes.facecolor": "white",
|
||||
"axes.grid": True,
|
||||
"grid.color": GRID,
|
||||
"grid.linewidth": 0.8,
|
||||
"legend.frameon": False,
|
||||
"savefig.bbox": "tight",
|
||||
"savefig.dpi": 160,
|
||||
"savefig.facecolor": "white",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def save(fig: plt.Figure, path: Path) -> None:
|
||||
fig.savefig(path, dpi=160)
|
||||
plt.close(fig)
|
||||
|
||||
|
||||
def load_runs(folder: Path) -> tuple[list[dict[str, str]], list[dict[str, str]]]:
|
||||
thru: list[dict[str, str]] = []
|
||||
lats: list[dict[str, str]] = []
|
||||
for f in sorted(folder.glob("*.txt")):
|
||||
if f.name.startswith("host-"):
|
||||
continue
|
||||
text = f.read_text(encoding="utf-8", errors="replace")
|
||||
lat = parse_lat(text)
|
||||
if lat:
|
||||
lat["run"] = f.stem
|
||||
lats.append(lat)
|
||||
continue
|
||||
p = parse_bench(text)
|
||||
if p.get("pub_msgs") or p.get("agg_msgs"):
|
||||
p["run"] = f.stem
|
||||
thru.append(p)
|
||||
return sorted(thru, key=thru_sort), sorted(lats, key=lat_sort)
|
||||
|
||||
|
||||
def kv_file(path: Path) -> dict[str, str]:
|
||||
out: dict[str, str] = {}
|
||||
if not path.exists():
|
||||
return out
|
||||
for line in path.read_text(encoding="utf-8", errors="replace").splitlines():
|
||||
if "=" in line and not line.startswith("---"):
|
||||
k, _, v = line.partition("=")
|
||||
if k.strip() in out:
|
||||
continue
|
||||
out[k.strip()] = v.strip()
|
||||
return out
|
||||
|
||||
|
||||
def thru_table(thru: list[dict[str, str]]) -> str:
|
||||
lines = [
|
||||
"| Run | Mode | Aggregate msgs/s | Pub msgs/s | Pub MB/s | Sub msgs/s | Sub MB/s |",
|
||||
"|-----|------|------------------|------------|----------|------------|----------|",
|
||||
]
|
||||
for p in thru:
|
||||
lines.append(
|
||||
f"| `{p['run']}` | {p.get('mode', '')} | {fmt_int(p.get('agg_msgs'))} | "
|
||||
f"{fmt_int(p.get('pub_msgs'))} | {p.get('pub_mb') or '—'} | "
|
||||
f"{fmt_int(p.get('sub_msgs'))} | {p.get('sub_mb') or '—'} |"
|
||||
)
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def delay_table(lats: list[dict[str, str]]) -> str:
|
||||
lines = [
|
||||
"| Run | Kind | Count | Pubs | Size | min | avg | p50 | p90 | p99 | max |",
|
||||
"|-----|------|-------|------|------|-----|-----|-----|-----|-----|-----|",
|
||||
]
|
||||
for p in lats:
|
||||
kind = lat_mode(p.get("run", ""), p.get("mode", ""))
|
||||
label = "ping (sequential RTT)" if kind == "ping" else "flood (burst queueing)"
|
||||
lines.append(
|
||||
f"| `{p['run']}` | {label} | {p.get('count', '')} | {p.get('pubs', '')} | "
|
||||
f"{p.get('size', '')} B | {p.get('min', '')} | {p.get('avg', '')} | "
|
||||
f"{p.get('p50', '')} | {p.get('p90', '')} | {p.get('p99', '')} | {p.get('max', '')} |"
|
||||
)
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def varz_table(path: Path) -> str:
|
||||
if not path.exists():
|
||||
return "_varz snapshot not captured._"
|
||||
rows = json.loads(path.read_text(encoding="utf-8"))
|
||||
lines = [
|
||||
"| Node | VMID | connections | in_msgs | out_msgs | cpu | cores | mem (B) | jetstream |",
|
||||
"|------|------|-------------|---------|----------|-----|-------|---------|-----------|",
|
||||
]
|
||||
for r in rows:
|
||||
if r.get("error"):
|
||||
lines.append(f"| {r.get('name')} | {r.get('vmid')} | error: {r['error']} | | | | | | |")
|
||||
continue
|
||||
lines.append(
|
||||
f"| {r.get('name')} | {r.get('vmid')} | {r.get('connections')} | "
|
||||
f"{r.get('in_msgs'):,} | {r.get('out_msgs'):,} | {r.get('cpu')} | "
|
||||
f"{r.get('cores')} | {r.get('mem'):,} | {r.get('jetstream')} |"
|
||||
)
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def charts(thru: list[dict[str, str]], lats: list[dict[str, str]], dest: Path) -> None:
|
||||
dest.mkdir(parents=True, exist_ok=True)
|
||||
style()
|
||||
by = {p["run"]: p for p in thru}
|
||||
core_keys = [k for k in CORE_LABELS if k in by]
|
||||
if core_keys:
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.4))
|
||||
x = list(range(len(core_keys)))
|
||||
w = 0.25
|
||||
agg = [int(by[k].get("agg_msgs") or 0) for k in core_keys]
|
||||
pub = [int(by[k].get("pub_msgs") or 0) for k in core_keys]
|
||||
sub = [int(by[k].get("sub_msgs") or 0) for k in core_keys]
|
||||
ax.bar([i - w for i in x], agg, w, label="Aggregate", color=DEEP)
|
||||
ax.bar(x, pub, w, label="Publish", color=INDIGO)
|
||||
ax.bar([i + w for i in x], sub, w, label="Subscribe", color=TEAL)
|
||||
ax.set_xticks(x, [CORE_LABELS[k] for k in core_keys])
|
||||
ax.set_ylabel("messages / second")
|
||||
ax.set_title("Core NATS throughput (fire-and-forget) — NS1 host run")
|
||||
ax.yaxis.set_major_formatter(FuncFormatter(k_fmt))
|
||||
ax.legend(loc="upper left")
|
||||
ax.set_axisbelow(True)
|
||||
save(fig, dest / "core-throughput.png")
|
||||
|
||||
js_keys = [k for k in JS_LABELS if k in by]
|
||||
if js_keys:
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.4))
|
||||
pubs = [int(by[k].get("pub_msgs") or 0) for k in js_keys]
|
||||
colors = [INDIGO, INDIGO, AMBER, LILAC][: len(js_keys)]
|
||||
ax.bar([JS_LABELS[k] for k in js_keys], pubs, color=colors)
|
||||
ax.set_ylabel("durable publish messages / second")
|
||||
ax.set_title("JetStream file store, replicas=3 — NS1 host run")
|
||||
ax.yaxis.set_major_formatter(FuncFormatter(k_fmt))
|
||||
ax.set_axisbelow(True)
|
||||
for i, v in enumerate(pubs):
|
||||
ax.text(i, v * 1.02, f"{v:,}", ha="center", va="bottom", fontsize=9, color=MUTED)
|
||||
save(fig, dest / "js-throughput.png")
|
||||
|
||||
pair = [("core-1p1s-50k-128", "js-1p-20k-128-r3"), ("core-4p4s-100k-128", "js-4p-50k-128-r3"), ("core-4p4s-50k-1k", "js-4p-20k-1k-r3")]
|
||||
if all(c in by and j in by for c, j in pair):
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.4))
|
||||
labels = ["1 publisher\n128 B", "4 publishers\n128 B", "4 publishers\n1 KiB"]
|
||||
core_pub = [int(by[c]["pub_msgs"]) for c, _ in pair]
|
||||
js_pub = [int(by[j]["pub_msgs"]) for _, j in pair]
|
||||
x = list(range(3))
|
||||
w = 0.35
|
||||
ax.bar([i - w / 2 for i in x], core_pub, w, label="Core NATS (no disk)", color=INDIGO)
|
||||
ax.bar([i + w / 2 for i in x], js_pub, w, label="JetStream r=3 file", color=AMBER)
|
||||
ax.set_xticks(x, labels)
|
||||
ax.set_yscale("log")
|
||||
ax.set_ylabel("publish messages / second (log)")
|
||||
ax.set_title("Core vs JetStream — NS1 host run")
|
||||
ax.legend(loc="upper right")
|
||||
ax.set_axisbelow(True)
|
||||
save(fig, dest / "core-vs-js.png")
|
||||
|
||||
if "core-4p4s-100k-128" in by and "core-4p4s-50k-1k" in by:
|
||||
fig, axes = plt.subplots(1, 2, figsize=(9.2, 4.2))
|
||||
labels = ["128 B\n4p4s", "1 KiB\n4p4s"]
|
||||
msgs = [int(by["core-4p4s-100k-128"].get("agg_msgs") or 0), int(by["core-4p4s-50k-1k"].get("agg_msgs") or 0)]
|
||||
mb = [float(by["core-4p4s-100k-128"].get("agg_mb") or 0), float(by["core-4p4s-50k-1k"].get("agg_mb") or 0)]
|
||||
axes[0].bar(labels, msgs, color=[INDIGO, AMBER])
|
||||
axes[0].set_title("Aggregate messages / second")
|
||||
axes[0].yaxis.set_major_formatter(FuncFormatter(k_fmt))
|
||||
axes[1].bar(labels, mb, color=[INDIGO, AMBER])
|
||||
axes[1].set_title("Aggregate MB / second")
|
||||
fig.suptitle("Core NATS payload effect — NS1 host run", fontsize=12, fontweight="semibold")
|
||||
fig.tight_layout()
|
||||
save(fig, dest / "payload-size.png")
|
||||
|
||||
if lats:
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.6))
|
||||
ordered = [p for p in lats]
|
||||
labels = [LAT_LABELS.get(p["run"], p["run"]) for p in ordered]
|
||||
x = list(range(len(ordered)))
|
||||
w = 0.25
|
||||
p50 = [ms(p["p50"]) for p in ordered]
|
||||
p90 = [ms(p["p90"]) for p in ordered]
|
||||
p99 = [ms(p["p99"]) for p in ordered]
|
||||
ax.bar([i - w for i in x], p50, w, label="p50", color=TEAL)
|
||||
ax.bar(x, p90, w, label="p90", color=INDIGO)
|
||||
ax.bar([i + w for i in x], p99, w, label="p99", color=AMBER)
|
||||
ax.set_xticks(x, labels)
|
||||
ax.set_yscale("log")
|
||||
ax.set_ylabel("milliseconds (log)")
|
||||
ax.set_title("Round-trip delay — NS1 host run")
|
||||
ax.axhline(1.0, color=GRID, linestyle="--", linewidth=1)
|
||||
ax.legend(loc="upper left")
|
||||
ax.set_axisbelow(True)
|
||||
save(fig, dest / "delay-percentiles.png")
|
||||
|
||||
|
||||
def figure(name: str, caption: str) -> str:
|
||||
return f"\n\n*{caption}*"
|
||||
|
||||
|
||||
def ratio(new: str | None, old: str | None) -> str:
|
||||
if not new or not old:
|
||||
return "—"
|
||||
a, b = float(new), float(old)
|
||||
if b == 0:
|
||||
return "—"
|
||||
return f"{a / b:.2f}×"
|
||||
|
||||
|
||||
def delay_ms_val(p: dict[str, str] | None, key: str) -> str | None:
|
||||
if not p or not p.get(key):
|
||||
return None
|
||||
return str(ms(p[key]))
|
||||
|
||||
|
||||
def delta_table(
|
||||
thru: list[dict[str, str]],
|
||||
lats: list[dict[str, str]],
|
||||
base_thru: list[dict[str, str]],
|
||||
base_lats: list[dict[str, str]],
|
||||
base_stamp: str,
|
||||
) -> str:
|
||||
bt = {p["run"]: p for p in base_thru}
|
||||
nt = {p["run"]: p for p in thru}
|
||||
bl = {p["run"]: p for p in base_lats}
|
||||
nl = {p["run"]: p for p in lats}
|
||||
keys = [
|
||||
("core-1p1s-50k-128", "pub", "Core 1p1s 128 B pub msgs/s"),
|
||||
("core-8p8s-200k-128", "agg", "Core 8p8s 128 B aggregate msgs/s"),
|
||||
("js-1p-20k-128-r3", "pub", "JS file r=3 1p 128 B pub msgs/s"),
|
||||
("js-4p-50k-128-r3", "pub", "JS file r=3 4p 128 B pub msgs/s"),
|
||||
("js-4p-20k-1k-r3", "pub", "JS file r=3 4p 1 KiB pub msgs/s"),
|
||||
("js-mem-1p-20k-128-r3", "pub", "JS memory r=3 1p 128 B pub msgs/s"),
|
||||
("js-mem-4p-50k-128-r3", "pub", "JS memory r=3 4p 128 B pub msgs/s"),
|
||||
]
|
||||
lines = [
|
||||
f"| Metric | Baseline `{base_stamp}` | This run | Ratio |",
|
||||
"|--------|-------------------------|----------|-------|",
|
||||
]
|
||||
for run, kind, label in keys:
|
||||
old, new = bt.get(run), nt.get(run)
|
||||
ok = "pub_msgs" if kind == "pub" else "agg_msgs"
|
||||
ov = old.get(ok) if old else None
|
||||
nv = new.get(ok) if new else None
|
||||
lines.append(f"| {label} | {fmt_int(ov)} | {fmt_int(nv)} | {ratio(nv, ov)} |")
|
||||
old_p, new_p = bl.get("lat-ping-1k-128"), nl.get("lat-ping-1k-128")
|
||||
if old_p or new_p:
|
||||
ov = delay_ms_val(old_p, "p99")
|
||||
nv = delay_ms_val(new_p, "p99")
|
||||
# smaller delay is better — invert ratio label
|
||||
r = "—"
|
||||
if ov and nv and float(nv) != 0:
|
||||
r = f"{float(ov) / float(nv):.2f}× faster" if float(nv) < float(ov) else f"{float(nv) / float(ov):.2f}× slower"
|
||||
lines.append(
|
||||
f"| Ping p99 (ms) | {old_p.get('p99') if old_p else '—'} | {new_p.get('p99') if new_p else '—'} | {r} |"
|
||||
)
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def chart_delta(
|
||||
thru: list[dict[str, str]],
|
||||
base_thru: list[dict[str, str]],
|
||||
dest: Path,
|
||||
) -> None:
|
||||
bt = {p["run"]: p for p in base_thru}
|
||||
nt = {p["run"]: p for p in thru}
|
||||
labels = ["Core 1p\n128 B pub", "JS file 1p\n128 B", "JS file 4p\n128 B", "JS mem 1p\n128 B"]
|
||||
keys = ["core-1p1s-50k-128", "js-1p-20k-128-r3", "js-4p-50k-128-r3", "js-mem-1p-20k-128-r3"]
|
||||
old = [int(bt[k]["pub_msgs"]) if k in bt and bt[k].get("pub_msgs") else 0 for k in keys]
|
||||
new = [int(nt[k]["pub_msgs"]) if k in nt and nt[k].get("pub_msgs") else 0 for k in keys]
|
||||
if not any(new):
|
||||
return
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.4))
|
||||
x = list(range(len(labels)))
|
||||
w = 0.35
|
||||
ax.bar([i - w / 2 for i in x], old, w, label="Baseline 1c/1G/ZFS", color=MUTED)
|
||||
ax.bar([i + w / 2 for i in x], new, w, label="8c/16G/tmpfs (+ mem rows)", color=INDIGO)
|
||||
ax.set_xticks(x, labels)
|
||||
ax.set_yscale("log")
|
||||
ax.set_ylabel("publish messages / second (log)")
|
||||
ax.set_title("Measured delta vs 20260912T051237Z")
|
||||
ax.legend(loc="upper right")
|
||||
ax.set_axisbelow(True)
|
||||
save(fig, dest / "delta-vs-baseline.png")
|
||||
|
||||
|
||||
def write_markdown(
|
||||
folder: Path,
|
||||
thru: list[dict[str, str]],
|
||||
lats: list[dict[str, str]],
|
||||
compare: Path | None = None,
|
||||
) -> str:
|
||||
before = kv_file(folder / "host-before.txt")
|
||||
after = kv_file(folder / "host-after.txt")
|
||||
stamp = folder.name
|
||||
method = (Path(__file__).resolve().parent / "ns1-study-methodology.md").read_text(encoding="utf-8")
|
||||
delta_md = ""
|
||||
base_thru: list[dict[str, str]] = []
|
||||
base_lats: list[dict[str, str]] = []
|
||||
if compare and compare.is_dir():
|
||||
base_thru, base_lats = load_runs(compare)
|
||||
delta_md = (
|
||||
f"## Measured delta vs `{compare.name}`\n\n"
|
||||
"Baseline: 1 core / 1 GiB / JetStream on ZFS. This run: 8 cores / 16 GiB / "
|
||||
"JetStream **tmpfs** (file r=3) plus extra **memory** store rows. veth/10G unchanged.\n\n"
|
||||
+ delta_table(thru, lats, base_thru, base_lats, compare.name)
|
||||
+ "\n"
|
||||
)
|
||||
if (folder / "charts" / "delta-vs-baseline.png").exists():
|
||||
delta_md += "\n" + figure("delta-vs-baseline.png", "Baseline vs maximized publish rates (log)")
|
||||
delta_md += "\n"
|
||||
figs = []
|
||||
charts_dir = folder / "charts"
|
||||
if (charts_dir / "core-throughput.png").exists():
|
||||
figs.append("### Core NATS\n\n" + figure("core-throughput.png", "Core NATS throughput at four loads (NS1 host run)"))
|
||||
if (charts_dir / "payload-size.png").exists():
|
||||
figs.append("### Payload size (core)\n\n" + figure("payload-size.png", "Core NATS 128 B vs 1 KiB (NS1 host run)"))
|
||||
if (charts_dir / "js-throughput.png").exists():
|
||||
figs.append("### JetStream r=3 file\n\n" + figure("js-throughput.png", "JetStream durable publish rate (NS1 host run)"))
|
||||
if (charts_dir / "core-vs-js.png").exists():
|
||||
figs.append("### Core vs JetStream\n\n" + figure("core-vs-js.png", "Core vs JetStream publish rate, log scale (NS1 host run)"))
|
||||
if (charts_dir / "delay-percentiles.png").exists():
|
||||
figs.append("### Delay\n\n" + figure("delay-percentiles.png", "Ping vs flood delay percentiles, log scale (NS1 host run)"))
|
||||
|
||||
ping = next((p for p in lats if "ping" in p.get("run", "")), None)
|
||||
js1 = next((p for p in thru if p["run"] == "js-1p-20k-128-r3"), None)
|
||||
core1 = next((p for p in thru if p["run"] == "core-1p1s-50k-128"), None)
|
||||
|
||||
md = f"""**Progress report (maximized NS1 study)** · run `{stamp}` (UTC)
|
||||
|
||||
> **Execution provenance.** Every process for this study ran on **NS1.GEORGELAMBERT.ORG** (`70.88.205.138`): `maximize-ns1-study.sh` (cores/RAM/`max_mem`/tmpfs), then `study-on-ns1.sh`, `nats bench`, `latency.mjs` (LXC 510), matplotlib, pandoc, weasyprint. Traffic stayed on `vmbr1`. veth/10G was **not** changed. After the ladder, JetStream was put back on ZFS and product streams were re-created; **8 cores / 16 GiB / max_mem 8G stay**.
|
||||
|
||||
{delta_md}
|
||||
|
||||
---
|
||||
|
||||
## 1. Executive summary
|
||||
|
||||
| Item | This NS1-host run |
|
||||
|------|-------------------|
|
||||
| Control plane | NS1.GEORGELAMBERT.ORG (`70.88.205.138`), user `{before.get("whoami", "marchon")}` |
|
||||
| Bench client | LXC {before.get("client_vmid", "510")} `verae-px-worker` |
|
||||
| Brokers | LXC 511/512/513 `nats-a/b/c` on `10.10.10.21–23` |
|
||||
| Client URL | `{before.get("nats_url", "")}` |
|
||||
| Host load before | `{before.get("loadavg", "n/a")}` |
|
||||
| Host load after | `{after.get("loadavg", "n/a")}` |
|
||||
| Core 1p1s 128 B pub | {fmt_int(core1.get("pub_msgs") if core1 else None)} msgs/s |
|
||||
| JetStream 1p 128 B r=3 | {fmt_int(js1.get("pub_msgs") if js1 else None)} durable pubs/s |
|
||||
| Ping p50 / p99 | {ping.get("p50") if ping else "—"} / {ping.get("p99") if ping else "—"} |
|
||||
|
||||
Product traffic is the JetStream row. Ping is one-message delay. Flood is mailbox catch-up after a burst.
|
||||
|
||||
---
|
||||
|
||||
## 2. Where it ran (and where it did not)
|
||||
|
||||
```text
|
||||
Operator laptop ──ssh──► NS1.GEORGELAMBERT.ORG 70.88.205.138
|
||||
study-on-ns1.sh
|
||||
python3 build-ns1-study-report.py
|
||||
sudo pct exec 510 ──► nats bench / latency.mjs
|
||||
│
|
||||
▼ vmbr1
|
||||
10.10.10.21-23 :4222
|
||||
```
|
||||
|
||||
- **Did run on 138:** bash, python3, matplotlib, pandoc, weasyprint, `pct`, nats-server (in LXC), nats CLI and Node (in LXC 510).
|
||||
- **Did not run on the laptop:** no local `nats bench`, no local charting, no local WeasyPrint for this file.
|
||||
|
||||
---
|
||||
|
||||
## 3. Results (this run)
|
||||
|
||||
### Host and brokers
|
||||
|
||||
**Before**
|
||||
|
||||
{varz_table(folder / "varz-before.json")}
|
||||
|
||||
**After**
|
||||
|
||||
{varz_table(folder / "varz-after.json")}
|
||||
|
||||
nproc={before.get("nproc", "?")} · uname=`{before.get("uname", "")}`
|
||||
|
||||
### Throughput
|
||||
|
||||
{thru_table(thru)}
|
||||
|
||||
### Round-trip delay
|
||||
|
||||
{delay_table(lats)}
|
||||
|
||||
{chr(10).join(figs)}
|
||||
|
||||
---
|
||||
|
||||
{method}
|
||||
|
||||
---
|
||||
|
||||
## 6. Reproducing this study
|
||||
|
||||
On **NS1 only**:
|
||||
|
||||
```bash
|
||||
cd ~/verae-src/verae-nats-cluster
|
||||
bash scripts/study-on-ns1.sh
|
||||
```
|
||||
|
||||
The script exits if `hostname` is not NS1. Outputs land in `results/<utc>/` including `nats-cluster-bench-ns1.{{md,html,pdf}}` and `charts/`. Copy those into `zapier-decisions/reports/` for the progress repo and catalog.
|
||||
|
||||
Raw logs for this run: `results/{stamp}/`.
|
||||
"""
|
||||
return md
|
||||
|
||||
|
||||
def render(md_path: Path, html_path: Path, pdf_path: Path) -> None:
|
||||
css = Path(__file__).resolve().parent / "docs-print.css"
|
||||
header = html_path.with_suffix(".hdr.html")
|
||||
banner = html_path.with_suffix(".ban.html")
|
||||
css_text = css.read_text(encoding="utf-8") if css.exists() else ""
|
||||
header.write_text(f"<style>{css_text}</style>\n", encoding="utf-8")
|
||||
banner.write_text(
|
||||
'<div class="doc-banner">'
|
||||
'<nav class="site"><a href="/">zapier.georgelambert.org</a>'
|
||||
' · <a href="/index-md.html">Markdown indexes</a></nav>'
|
||||
'<div class="kicker">Verae Time × Zapier · progress report · maximized NS1 study</div>'
|
||||
"<h1>NATS cluster message speed — maximized (RAM disk + 8 cores)</h1>"
|
||||
'<div class="source-path">packages/zapier-decisions/reports/nats-cluster-bench-ns1.md</div>'
|
||||
"</div>\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
r = subprocess.run(
|
||||
[
|
||||
"pandoc",
|
||||
str(md_path),
|
||||
"-o",
|
||||
str(html_path),
|
||||
"--standalone",
|
||||
f"--resource-path={md_path.parent}",
|
||||
"--highlight-style=breezedark",
|
||||
"--metadata=title=NATS cluster message speed — NS1 host study",
|
||||
f"--include-in-header={header}",
|
||||
f"--include-before-body={banner}",
|
||||
],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
header.unlink(missing_ok=True)
|
||||
banner.unlink(missing_ok=True)
|
||||
if r.returncode != 0:
|
||||
raise SystemExit(f"pandoc failed: {r.stderr[-800:]}")
|
||||
w = subprocess.run(["weasyprint", str(html_path), str(pdf_path)], capture_output=True, text=True)
|
||||
if w.returncode != 0:
|
||||
raise SystemExit(f"weasyprint failed: {w.stderr[-800:]}")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
folder = Path(sys.argv[1] if len(sys.argv) > 1 else ".")
|
||||
compare = Path(sys.argv[2]) if len(sys.argv) > 2 and sys.argv[2] else None
|
||||
thru, lats = load_runs(folder)
|
||||
charts(thru, lats, folder / "charts")
|
||||
if compare and compare.is_dir():
|
||||
base_thru, _base_lats = load_runs(compare)
|
||||
chart_delta(thru, base_thru, folder / "charts")
|
||||
md = write_markdown(folder, thru, lats, compare if compare and compare.is_dir() else None)
|
||||
md_path = folder / "nats-cluster-bench-ns1.md"
|
||||
md_path.write_text(md, encoding="utf-8")
|
||||
html_path = folder / "nats-cluster-bench-ns1.html"
|
||||
pdf_path = folder / "nats-cluster-bench-ns1.pdf"
|
||||
render(md_path, html_path, pdf_path)
|
||||
print(f"wrote {md_path}")
|
||||
print(f"wrote {html_path}")
|
||||
print(f"wrote {pdf_path}")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
400
scripts/build-optimal-report.py
Normal file
400
scripts/build-optimal-report.py
Normal file
|
|
@ -0,0 +1,400 @@
|
|||
#!/usr/bin/env python3
|
||||
"""One large comparison report from all NS1 result folders + extras (UDP/MQTT/reconnect)."""
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
import json
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
import matplotlib
|
||||
|
||||
matplotlib.use("Agg")
|
||||
import matplotlib.pyplot as plt
|
||||
from matplotlib.ticker import FuncFormatter
|
||||
|
||||
_spec = importlib.util.spec_from_file_location(
|
||||
"bench_report", Path(__file__).resolve().parent / "bench-report.py"
|
||||
)
|
||||
_br = importlib.util.module_from_spec(_spec)
|
||||
assert _spec.loader
|
||||
_spec.loader.exec_module(_br)
|
||||
fmt_int = _br.fmt_int
|
||||
parse_bench = _br.parse_bench
|
||||
parse_lat = _br.parse_lat
|
||||
|
||||
INDIGO, DEEP, TEAL, AMBER, MUTED = "#4f46e5", "#312e81", "#047857", "#b45309", "#5b6178"
|
||||
|
||||
|
||||
def k_fmt(x, _p=None):
|
||||
if x >= 1_000_000:
|
||||
return f"{x/1e6:.2f}M"
|
||||
if x >= 1000:
|
||||
return f"{x/1000:.0f}k"
|
||||
return f"{x:.0f}"
|
||||
|
||||
|
||||
def load_folder(folder: Path) -> dict:
|
||||
thru, lats, extra = {}, {}, {}
|
||||
if not folder.is_dir():
|
||||
return {"thru": thru, "lats": lats, "extra": extra, "stamp": folder.name}
|
||||
for f in folder.glob("*.txt"):
|
||||
text = f.read_text(encoding="utf-8", errors="replace")
|
||||
if f.name.startswith("host-"):
|
||||
continue
|
||||
if "mqtt" in f.name or "udp-ping" in f.name:
|
||||
for line in text.splitlines():
|
||||
line = line.strip()
|
||||
if line.startswith("{"):
|
||||
extra[f.stem] = json.loads(line)
|
||||
break
|
||||
continue
|
||||
lat = parse_lat(text)
|
||||
if lat:
|
||||
lats[f.stem] = lat
|
||||
continue
|
||||
p = parse_bench(text)
|
||||
if p.get("pub_msgs") or p.get("agg_msgs"):
|
||||
thru[f.stem] = p
|
||||
return {"thru": thru, "lats": lats, "extra": extra, "stamp": folder.name}
|
||||
|
||||
|
||||
def pub(d, run):
|
||||
p = d["thru"].get(run) or {}
|
||||
return p.get("pub_msgs")
|
||||
|
||||
|
||||
def latp(d, run, key="p99"):
|
||||
p = d["lats"].get(run) or {}
|
||||
return p.get(key, "")
|
||||
|
||||
|
||||
def row(*cells):
|
||||
return "| " + " | ".join(cells) + " |"
|
||||
|
||||
|
||||
def save(fig, path: Path):
|
||||
fig.savefig(path, dpi=150, bbox_inches="tight", facecolor="white")
|
||||
plt.close(fig)
|
||||
|
||||
|
||||
def charts(latest: dict, folders: list[dict], dest: Path):
|
||||
dest.mkdir(parents=True, exist_ok=True)
|
||||
plt.rcParams.update({"font.size": 9, "axes.grid": True, "grid.color": "#d9dce8"})
|
||||
# r=1 vs r=3 file/mem from latest
|
||||
labels, file_r, mem_r = [], [], []
|
||||
for lab, fr, mr in (
|
||||
("1p file", "js-file-1p-20k-128-r1", "js-1p-20k-128-r3"),
|
||||
("4p file", "js-file-4p-50k-128-r1", "js-4p-50k-128-r3"),
|
||||
("1p mem", "js-mem-1p-20k-128-r1", "js-mem-1p-20k-128-r3"),
|
||||
("4p mem", "js-mem-4p-50k-128-r1", "js-mem-4p-50k-128-r3"),
|
||||
):
|
||||
a, b = pub(latest, fr), pub(latest, mr)
|
||||
if a or b:
|
||||
labels.append(lab)
|
||||
file_r.append(int(a or 0))
|
||||
mem_r.append(int(b or 0))
|
||||
if labels:
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.3))
|
||||
x = range(len(labels))
|
||||
ax.bar([i - 0.2 for i in x], file_r, 0.4, label="replicas=1", color=TEAL)
|
||||
ax.bar([i + 0.2 for i in x], mem_r, 0.4, label="replicas=3", color=AMBER)
|
||||
ax.set_xticks(list(x), labels)
|
||||
ax.set_ylabel("pub msgs/s")
|
||||
ax.set_title("This run: replica cost (1 vs 3)")
|
||||
ax.yaxis.set_major_formatter(FuncFormatter(k_fmt))
|
||||
ax.legend()
|
||||
save(fig, dest / "replicas.png")
|
||||
# historical JS 1p file
|
||||
names, vals = [], []
|
||||
for d, label in zip(
|
||||
folders,
|
||||
[d["stamp"] for d in folders],
|
||||
):
|
||||
v = pub(d, "js-1p-20k-128-r3")
|
||||
if v:
|
||||
names.append(label[-7:] if len(label) > 8 else label)
|
||||
vals.append(int(v))
|
||||
if names:
|
||||
fig, ax = plt.subplots(figsize=(9.2, 4.0))
|
||||
ax.bar(names, vals, color=INDIGO)
|
||||
ax.set_title("JS file r=3 1p 128 B across studies")
|
||||
ax.set_ylabel("pub msgs/s")
|
||||
ax.yaxis.set_major_formatter(FuncFormatter(k_fmt))
|
||||
save(fig, dest / "history-js1p.png")
|
||||
|
||||
|
||||
def write_md(latest: dict, hist: list[dict], charts_rel: str) -> str:
|
||||
t = latest["thru"]
|
||||
e = latest["extra"]
|
||||
mqtt = e.get("mqtt-qos0-5k-128") or {}
|
||||
udp = e.get("udp-ping-1k-128") or {}
|
||||
ping = latest["lats"].get("lat-ping-1k-128") or {}
|
||||
recon = latest["lats"].get("lat-reconnect-200-128") or {}
|
||||
|
||||
def js_table():
|
||||
runs = [
|
||||
("js-file-1p-20k-128-r1", "file r=1 1p 128 B"),
|
||||
("js-file-4p-50k-128-r1", "file r=1 4p 128 B"),
|
||||
("js-1p-20k-128-r3", "file r=3 1p 128 B"),
|
||||
("js-4p-50k-128-r3", "file r=3 4p 128 B"),
|
||||
("js-4p-20k-1k-r3", "file r=3 4p 1 KiB"),
|
||||
("js-file-1p-20k-4k-r3", "file r=3 1p 4 KiB"),
|
||||
("js-mem-1p-20k-128-r1", "memory r=1 1p 128 B"),
|
||||
("js-mem-4p-50k-128-r1", "memory r=1 4p 128 B"),
|
||||
("js-mem-1p-20k-128-r3", "memory r=3 1p 128 B"),
|
||||
("js-mem-4p-50k-128-r3", "memory r=3 4p 128 B"),
|
||||
("js-mem-4p-20k-1k-r3", "memory r=3 4p 1 KiB"),
|
||||
]
|
||||
lines = [
|
||||
row("Run", "What", "Pub msgs/s", "Pub MB/s"),
|
||||
row("---", "---", "---", "---"),
|
||||
]
|
||||
for k, lab in runs:
|
||||
p = t.get(k)
|
||||
if not p:
|
||||
continue
|
||||
lines.append(row(f"`{k}`", lab, fmt_int(p.get("pub_msgs")), p.get("pub_mb") or "—"))
|
||||
return "\n".join(lines)
|
||||
|
||||
hist_lines = [
|
||||
row("Study", "Env", "Core 1p pub", "JS file r=3 1p", "JS mem r=3 4p", "Ping p99"),
|
||||
row("---", "---", "---", "---", "---", "---"),
|
||||
]
|
||||
labels_env = {
|
||||
"20260912T045131Z": "1c/1G ZFS (off-box orch.)",
|
||||
"20260912T051237Z": "1c/1G ZFS (NS1 orch.)",
|
||||
"20260912T053120Z": "8c/16G tmpfs + mem extra",
|
||||
}
|
||||
for d in hist + [latest]:
|
||||
env = labels_env.get(d["stamp"], f"8c/16G ZFS exhaustive `{d['stamp']}`")
|
||||
hist_lines.append(
|
||||
row(
|
||||
f"`{d['stamp']}`",
|
||||
env,
|
||||
fmt_int(pub(d, "core-1p1s-50k-128")),
|
||||
fmt_int(pub(d, "js-1p-20k-128-r3")),
|
||||
fmt_int(pub(d, "js-mem-4p-50k-128-r3")),
|
||||
latp(d, "lat-ping-1k-128"),
|
||||
)
|
||||
)
|
||||
|
||||
r1 = int(pub(latest, "js-file-1p-20k-128-r1") or 0)
|
||||
r3 = int(pub(latest, "js-1p-20k-128-r3") or 0)
|
||||
mem1 = int(pub(latest, "js-mem-1p-20k-128-r1") or 0)
|
||||
mem3 = int(pub(latest, "js-mem-1p-20k-128-r3") or 0)
|
||||
replica_cost = f"{r1/r3:.2f}×" if r3 else "—"
|
||||
mem_gain = f"{mem1/r3:.2f}×" if r3 and mem1 else "—"
|
||||
|
||||
mqtt_rate = mqtt.get("pubs_per_sec", "—")
|
||||
udp_p99 = udp.get("p99", "—")
|
||||
|
||||
return f"""**Progress report — optimal configuration study** · `{latest['stamp']}` (UTC) · all code on **NS1.GEORGELAMBERT.ORG** (`70.88.205.138`)
|
||||
|
||||
This document folds every ladder we have run (1-core ZFS, NS1-orchestrated, tmpfs maximize, and this exhaustive 8c/16G **ZFS** factorial) plus UDP / MQTT / reconnect probes. It recommends a lab config and a **three-box HP DL360 Gen10** projection. veth/10G was not changed.
|
||||
|
||||
---
|
||||
|
||||
## 1. Verdict (read this first)
|
||||
|
||||
**Keep NATS + JetStream.** Do not replace the fabric with MQTT, UDP, or a custom persistent-socket protocol for Verae jobs/events/archive. Those are either slower, less durable, or already what NATS is.
|
||||
|
||||
**Lab (NS1, one host, three LXC) — optimal now**
|
||||
|
||||
| Stream | Storage | Replicas | Why |
|
||||
|--------|---------|----------|-----|
|
||||
| `ZAPIER_JOBS`, `ZAPIER_WEBHOOKS`, `VERAE_ARCHIVE` | **file** (ZFS) | **3** | Survive a nats LXC death; archive must persist |
|
||||
| `ZAPIER_EVENTS` | **memory** | **3** | Waiters are latency-sensitive; events rebuild from job status |
|
||||
| `ZAPIER_USAGE` | file | 3 | Telemetry, limits + max-age |
|
||||
|
||||
Keep **8 cores / 16 GiB / `max_mem: 8G`** on 510–513 (already live). Do **not** leave JetStream on tmpfs. Do **not** drop product streams to r=1. Reuse **one NATS connection per process** (already true in middleware); never connect-per-message.
|
||||
|
||||
**Metal (3× DL360 Gen10) — optimal later**
|
||||
|
||||
Same stream table. File store on **local NVMe/M.2**, not a shared SAN. Cluster + client on **10GbE** (or 25GbE if you already have it). Dual Gold Xeon is surplus CPU for this workload; 8–16 cores dedicated to `nats-server` is enough. Expected JS file r=3: **~40–80k** 128 B pubs/s (about **3–6×** this lab’s 8c ZFS 1p, **2–4×** tmpfs 1p) — bounded by **10GbE replica RTT**, not by Xeon clocks. Core NATS will sit in the **1–3M msgs/s** band until the NIC saturates (~9 Gbit/s ≈ 8–9M × 128 B theoretical; CPU and client will hit first).
|
||||
|
||||
---
|
||||
|
||||
## 2. What we actually ran (this exhaustive pass)
|
||||
|
||||
Live cluster during this run: LXC 510–513 **8 cores / 16 GiB**, JetStream **on ZFS** (tmpfs from the maximize study was already unmounted). Extra factorial: file/memory × replicas 1/3, 4 KiB file r=3, reconnect-per-message ping, UDP echo 510→511, MQTT QoS0 against nats-a `:1883`. Product streams were not the bench target.
|
||||
|
||||
### 2.1 Cross-study history
|
||||
|
||||
{chr(10).join(hist_lines)}
|
||||
|
||||

|
||||
|
||||
### 2.2 This run — JetStream factorial
|
||||
|
||||
{js_table()}
|
||||
|
||||
Replica **1 vs 3** on this stand (file 1p 128 B): r=1 is {fmt_int(str(r1) if r1 else None)} vs r=3 {fmt_int(str(r3) if r3 else None)} ({replica_cost} if r=3 is the slower one). Memory r=1 1p {fmt_int(str(mem1) if mem1 else None)} vs memory r=3 {fmt_int(str(mem3) if mem3 else None)}.
|
||||
|
||||

|
||||
|
||||
### 2.3 Delay, reconnect tax, UDP, MQTT
|
||||
|
||||
| Probe | Result | Meaning |
|
||||
|-------|--------|---------|
|
||||
| NATS ping (persistent sockets) p50 / p99 | {ping.get('p50','—')} / {ping.get('p99','—')} | Quiet hop with a long-lived TCP conn |
|
||||
| NATS **reconnect-per-message** p50 / p99 | {recon.get('p50','—')} / {recon.get('p99','—')} | TCP+NATS handshake on every pub — this is the tax to avoid |
|
||||
| UDP echo 510→511 p99 | {udp_p99} | Raw datagram ceiling on the same veth (no NATS) |
|
||||
| MQTT QoS0 5k×128 B | {mqtt_rate} pubs/s | nats-server MQTT gateway on `:1883` |
|
||||
|
||||
Core 1p1s 128 B this run: {fmt_int(pub(latest, 'core-1p1s-50k-128'))} pub msgs/s. Flood delay is still backlog/consume_rate, not RTT.
|
||||
|
||||
---
|
||||
|
||||
## 3. Alternative transports (why we are not switching the fabric)
|
||||
|
||||
NATS already **is** persistent TCP sockets with a tiny binary protocol, automatic reconnect, and optional JetStream durability. “Reduce connection overhead” is a **client** discipline: hold the connection. The reconnect probe exists to prove that opening a socket per job would dominate ping RTT.
|
||||
|
||||
| Idea | Fit for Verae jobs/events/archive | Throughput vs NATS core | Durability |
|
||||
|------|-----------------------------------|-------------------------|------------|
|
||||
| **NATS core pub/sub** | Fan-out, request-reply (`verae.billing.*`) | Highest we measured (~0.5–2M msgs/s) | None |
|
||||
| **NATS JetStream file r=3** | Jobs, webhooks, archive | ~8–23k on this lab; see metal projection | Disk + 1-node loss |
|
||||
| **NATS JetStream memory r=3** | Events mailbox | ~22–36k on this lab | RAM + 1-node loss; **empty on full restart** |
|
||||
| **MQTT** (NATS gateway or Mosquitto) | IoT endpoints that already speak MQTT | This probe: {mqtt_rate} pubs/s QoS0 — typically **well below** NATS core; QoS1 ≈ JetStream-ish with more chatter | QoS1/2 session state; not our WORM model |
|
||||
| **UDP** | Telemetry that may drop | RTT {udp_p99} p99 — fastest hop, **no** reliability, no cluster, no auth | None |
|
||||
| **Custom persistent sockets / HTTP long-poll** | Worse NATS | You would re-implement reconnect, flow control, and fan-out | DIY |
|
||||
| **WebSocket** | Browsers only | Extra framing; NATS already has WS for UIs, not for middleware | Same as core/JS behind it |
|
||||
| **QUIC / WebTransport** | Lossy WAN / browsers | NATS QUIC is not the lab path; 10GbE LAN does not need it | Same |
|
||||
| **Kafka / Redis streams** | Heavy log replay | Higher ops cost; not on `vmbr1` today | Yes, heavier |
|
||||
|
||||
**MQTT:** NATS documents MQTT as an *enabling* gateway for existing IoT, and prefers NATS end-to-end for greenfield. Zapier cloud never talks NATS or MQTT; it talks HTTPS. Putting MQTT in the middle of timestamp jobs adds protocol translation and QoS timers without helping `jobId → events`. Use MQTT only if a device already cannot speak NATS.
|
||||
|
||||
**UDP:** Fine as a *measurement* of veth RTT. Unusable as the job fabric (no ack, no replica, no flow control). NATS ping is already within a small multiple of UDP on this bridge.
|
||||
|
||||
**Persistence sockets:** Middleware and keep already keep `NATS_URL` connections open. Optimal: one connection (or a small pool) per process, `max_reconnect`, jitter, no `connect()` in the per-job path. The reconnect ladder is the anti-pattern.
|
||||
|
||||
---
|
||||
|
||||
## 4. Optimal configurations
|
||||
|
||||
### 4.1 NS1 lab (now)
|
||||
|
||||
1. **Leave 8 cores / 16 GiB** on nats-a/b/c and the worker. Host has 40 cores / 377 GiB; this is cheap.
|
||||
2. **`max_mem: 8G`** stays. Required for memory streams.
|
||||
3. **File r=3 on ZFS** for jobs/webhooks/archive. tmpfs doubled JS 1p (7.4k→17k) but **loses the stream on reboot** — unacceptable for archive.
|
||||
4. **Memory r=3 for `ZAPIER_EVENTS`** if we accept “all three nats CTs reboot ⇒ in-flight waiters fall back to HTTP poll.” That matches the designed wait path (`GET /api/status/{{jobId}}`).
|
||||
5. **r=1 only for throwaway benches**, never product streams. Replica=3 is the point of three guests.
|
||||
6. **veth on vmbr1, no fake 10G NICs.** Already 10000Mb/s; JS does not fill it.
|
||||
7. **Pin cpusets** later if keep/fleet steal; not required to beat these numbers.
|
||||
8. Clients: persistent NATS connections; pull consumers with bounded `max_ack_pending` for webhooks.
|
||||
|
||||
### 4.2 Three HP DL360 Gen10 (projection — not measured)
|
||||
|
||||
Assumed bill of materials (state it in the buy):
|
||||
|
||||
| Piece | Assumption |
|
||||
|-------|------------|
|
||||
| Chassis | 3× DL360 Gen10 1U |
|
||||
| CPU | Dual 2nd-gen Xeon **Gold** (e.g. 6226R 16c or 6248 20c — **32–40 cores/box**) |
|
||||
| Memory | DDR4-2933, **192–384 GiB**/box (6–12×32 GiB); NATS will not use most of it |
|
||||
| Storage | **NVMe M.2 or U.2** for `/var/lib/nats/jetstream` (XFS or ext4, **not** shared ZFS over the network). RAID1 of two NVMe if you want disk HA *inside* a box |
|
||||
| Network | **10GbE** (FlexibleLOM or PCIe); dedicated VLAN for `:4222`+`:6222`. Do not share with public `vmbr0` traffic |
|
||||
| OS | Debian/Ubuntu bare metal, `nats-server` systemd, same `nats.conf` as lab (bind private IP only) |
|
||||
|
||||
**What changes vs NS1 LXC**
|
||||
|
||||
| Factor | NS1 today | 3× DL360 | Effect on JS file r=3 |
|
||||
|--------|-----------|----------|------------------------|
|
||||
| Failure domain | 1 Proxmox host | 3 chassis, 3 NVMe, 3 NICs | r=3 **means** something |
|
||||
| Disk | Shared ZFS SSD2 | Local NVMe fsync ~50–150 µs | Big win vs ZFS; similar to tmpfs for sequential 128 B |
|
||||
| Replica path | veth/bridge (~µs–tens of µs) | 10GbE RTT typically **50–200 µs** | **Slower than same-host tmpfs**, faster than a bad SAN |
|
||||
| CPU | 8 of 40 shared | 32–40 dedicated Gold cores | Headroom for many clients, not 10× JS |
|
||||
| NIC | software 10G veth, already ~5 Gbit/s core | real 10GbE ~9 Gbit/s TCP | Core NATS can grow; JS r=3 stays replica-bound |
|
||||
|
||||
**Projected bands** (128 B, 3-node cluster, dedicated 10GbE, local NVMe, 8+ cores pinned to nats-server):
|
||||
|
||||
| Workload | NS1 measured (best) | DL360 projection | Confidence |
|
||||
|----------|---------------------|------------------|------------|
|
||||
| Core pub/sub 1p | 0.5–0.8M | **0.8–2M** | Medium — NIC + syscall, plenty of CPU |
|
||||
| Core 4p4s 1 KiB | ~0.6–0.7M (~0.6 GB/s) | **~1M msgs/s / ~1 GB/s** approaching 10GbE | Medium |
|
||||
| JS file r=1 | this run r=1 | **80–200k** pubs/s | Medium — NVMe + no replica wait |
|
||||
| JS file r=3 | 7–23k (ZFS/tmpfs) | **40–80k** pubs/s | Medium-low — replica RTT dominates; 3 NVMe still help vs shared ZFS |
|
||||
| JS memory r=3 | 22–36k | **50–100k** | Medium-low — RAM + 10GbE ack |
|
||||
| Ping p99 | 0.7–1.4 ms | **0.2–0.6 ms** | Medium — real NIC but no Proxmox tax |
|
||||
|
||||
These are **not** DL360 measurements. Scale from: (a) our replica-1 vs replica-3 ratio once this run’s r=1 numbers exist, (b) tmpfs vs ZFS ratio (2.35× on 1p), (c) Synadia/nats bench async file r=1 ~100–400k on NVMe loopback, derated for 10GbE RTT.
|
||||
|
||||
**Buy notes:** M.2 via Dual uFF / enablement kit; put JetStream on NVMe **directly**, not behind a RAID controller write-through unless you measure. 1GbE onboard is a trap — use 10GbE for `:6222`. Dual Gold is for isolation (nats vs worm/tree vs OS), not because JS needs 56 cores.
|
||||
|
||||
---
|
||||
|
||||
## 5. What we are not doing
|
||||
|
||||
- MQTT as the Zapier or middleware transport.
|
||||
- UDP for jobs.
|
||||
- Emulated 10G fiber NICs on LXC.
|
||||
- tmpfs as the production store.
|
||||
- r=1 for product streams.
|
||||
- Connect-per-job.
|
||||
|
||||
Re-run exhaustive: `bash scripts/exhaustive-ns1-study.sh` on NS1.
|
||||
"""
|
||||
|
||||
|
||||
def render(md_path: Path, html_path: Path, pdf_path: Path) -> None:
|
||||
css = Path(__file__).resolve().parent / "docs-print.css"
|
||||
header = html_path.with_suffix(".hdr.html")
|
||||
banner = html_path.with_suffix(".ban.html")
|
||||
header.write_text(f"<style>{css.read_text() if css.exists() else ''}</style>\n", encoding="utf-8")
|
||||
banner.write_text(
|
||||
'<div class="doc-banner">'
|
||||
'<nav class="site"><a href="/">zapier.georgelambert.org</a></nav>'
|
||||
'<div class="kicker">Verae Time × Zapier · progress report</div>'
|
||||
"<h1>NATS optimal configuration study</h1>"
|
||||
'<div class="source-path">packages/zapier-decisions/reports/optimal-config/REPORT.md</div>'
|
||||
"</div>\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
r = subprocess.run(
|
||||
[
|
||||
"pandoc",
|
||||
str(md_path),
|
||||
"-o",
|
||||
str(html_path),
|
||||
"--standalone",
|
||||
f"--resource-path={md_path.parent}",
|
||||
"--highlight-style=breezedark",
|
||||
"--metadata=title=NATS optimal configuration study",
|
||||
f"--include-in-header={header}",
|
||||
f"--include-before-body={banner}",
|
||||
],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
header.unlink(missing_ok=True)
|
||||
banner.unlink(missing_ok=True)
|
||||
if r.returncode != 0:
|
||||
raise SystemExit(f"pandoc failed: {r.stderr[-600:]}")
|
||||
w = subprocess.run(["weasyprint", str(html_path), str(pdf_path)], capture_output=True, text=True)
|
||||
if w.returncode != 0:
|
||||
raise SystemExit(f"weasyprint failed: {w.stderr[-600:]}")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
latest = Path(sys.argv[1])
|
||||
hist_dirs = [Path(p) for p in sys.argv[2:] if p and Path(p).is_dir()]
|
||||
data_latest = load_folder(latest)
|
||||
hist = [load_folder(p) for p in hist_dirs]
|
||||
charts_dir = latest / "charts-optimal"
|
||||
charts(data_latest, hist + [data_latest], charts_dir)
|
||||
md = write_md(data_latest, hist, "charts-optimal")
|
||||
md_path = latest / "optimal-config.md"
|
||||
md_path.write_text(md, encoding="utf-8")
|
||||
html_path = latest / "optimal-config.html"
|
||||
pdf_path = latest / "optimal-config.pdf"
|
||||
render(md_path, html_path, pdf_path)
|
||||
print(f"wrote {md_path}")
|
||||
print(f"wrote {html_path}", file=sys.stderr)
|
||||
print(f"wrote {pdf_path}", file=sys.stderr)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
69
scripts/create-cluster.sh
Executable file
69
scripts/create-cluster.sh
Executable file
|
|
@ -0,0 +1,69 @@
|
|||
#!/usr/bin/env bash
|
||||
# Create three distinct Proxmox LXC guests and start a JetStream cluster on vmbr1.
|
||||
# Does not touch host loopback NATS (127.0.0.1:4222) or vmbr0.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/cluster.env"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/scripts/lib-ct.sh"
|
||||
ct_require_proxmox
|
||||
|
||||
mapfile -t rows < <(printf '%s\n' "$NODES" | awk 'NF==3 {print}')
|
||||
[[ ${#rows[@]} -eq 3 ]] || { echo "need exactly 3 nodes in cluster.env" >&2; exit 1; }
|
||||
|
||||
declare -a VMIDS NAMES IPS
|
||||
for row in "${rows[@]}"; do
|
||||
# shellcheck disable=SC2086
|
||||
set -- $row
|
||||
VMIDS+=("$1"); NAMES+=("$2"); IPS+=("$3")
|
||||
done
|
||||
|
||||
i=0
|
||||
for i in 0 1 2; do
|
||||
ct_ensure "${VMIDS[$i]}" "${NAMES[$i]}" "${IPS[$i]}"
|
||||
ct_bootstrap_user "${VMIDS[$i]}"
|
||||
done
|
||||
|
||||
# Install nats-server + conf + systemd on each guest
|
||||
for i in 0 1 2; do
|
||||
routes=""
|
||||
for j in 0 1 2; do
|
||||
[[ $i -eq $j ]] && continue
|
||||
routes="${routes} nats-route://${IPS[$j]}:6222"$'\n'
|
||||
done
|
||||
tmpconf="$(mktemp)"
|
||||
NAME="${NAMES[$i]}" IP="${IPS[$i]}" CLUSTER="$CLUSTER_NAME" ROUTES="$routes" \
|
||||
python3 - "$ROOT/conf/nats.conf.tmpl" "$tmpconf" <<'PY'
|
||||
import os, pathlib, sys
|
||||
t = pathlib.Path(sys.argv[1]).read_text()
|
||||
out = t.replace("{{NAME}}", os.environ["NAME"]).replace("{{IP}}", os.environ["IP"]).replace("{{CLUSTER}}", os.environ["CLUSTER"]).replace("{{ROUTES}}", os.environ["ROUTES"])
|
||||
pathlib.Path(sys.argv[2]).write_text(out)
|
||||
PY
|
||||
sudo pct exec "${VMIDS[$i]}" -- bash -c 'cat > /tmp/nats.conf' < "$tmpconf"
|
||||
sudo pct exec "${VMIDS[$i]}" -- bash -c 'cat > /tmp/nats-server.service' < "$ROOT/systemd/nats-server.service"
|
||||
rm -f "$tmpconf"
|
||||
sudo pct exec "${VMIDS[$i]}" -- bash -lc "
|
||||
set -e
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
id nats >/dev/null 2>&1 || useradd -r -s /usr/sbin/nologin nats
|
||||
install -d -m 755 -o nats -g nats /var/lib/nats/jetstream /etc/nats
|
||||
mv /tmp/nats.conf /etc/nats/nats.conf
|
||||
chown root:root /etc/nats/nats.conf
|
||||
chmod 644 /etc/nats/nats.conf
|
||||
if [[ ! -x /usr/local/bin/nats-server ]]; then
|
||||
curl -fsSL https://github.com/nats-io/nats-server/releases/download/v${NATS_VER}/nats-server-v${NATS_VER}-linux-amd64.tar.gz -o /tmp/nats.tgz
|
||||
tar -xzf /tmp/nats.tgz -C /tmp
|
||||
install -m 0755 /tmp/nats-server-v${NATS_VER}-linux-amd64/nats-server /usr/local/bin/nats-server
|
||||
rm -rf /tmp/nats.tgz /tmp/nats-server-v${NATS_VER}-linux-amd64
|
||||
fi
|
||||
install -m 644 /tmp/nats-server.service /etc/systemd/system/nats-server.service
|
||||
systemctl daemon-reload
|
||||
systemctl enable --now nats-server
|
||||
"
|
||||
echo "nats-server ${NAMES[$i]} ${IPS[$i]}:4222 cluster ${IPS[$i]}:6222"
|
||||
done
|
||||
|
||||
echo "cluster client URL: nats://${IPS[0]}:4222,nats://${IPS[1]}:4222,nats://${IPS[2]}:4222"
|
||||
echo "lab loopback NATS on the host is unchanged (127.0.0.1:4222)"
|
||||
echo "next: bash $ROOT/scripts/test.sh"
|
||||
10
scripts/cutover-ns1.sh
Executable file
10
scripts/cutover-ns1.sh
Executable file
|
|
@ -0,0 +1,10 @@
|
|||
#!/usr/bin/env bash
|
||||
# Point NS1 test modules at the 3-node vmbr1 cluster. Does not change MOCK_VERAE or Zapier.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/client.env"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:$PATH"
|
||||
bash "$ROOT/scripts/ensure-streams.sh"
|
||||
echo "NATS_URL=$NATS_URL"
|
||||
echo "streams ensured. restart keep + fleet on the host after copying overlay/service JSON."
|
||||
122
scripts/docs-print.css
Normal file
122
scripts/docs-print.css
Normal file
|
|
@ -0,0 +1,122 @@
|
|||
/* Colored print + screen stylesheet for zapier.georgelambert.org */
|
||||
:root {
|
||||
--ink: #171a26;
|
||||
--muted: #5b6178;
|
||||
--line: #d9dce8;
|
||||
--bg: #f4f5fb;
|
||||
--paper: #ffffff;
|
||||
--accent: #4f46e5;
|
||||
--accent-deep: #312e81;
|
||||
--accent-soft: #eef0fe;
|
||||
--ok: #047857;
|
||||
--warn: #8a5a00;
|
||||
--code-bg: #1b1f33;
|
||||
--code-fg: #e8ecff;
|
||||
}
|
||||
html { background: var(--bg); }
|
||||
body {
|
||||
margin: 0 auto;
|
||||
padding: 1.5rem 1.25rem 3rem;
|
||||
max-width: 48rem;
|
||||
font: 15px/1.55 -apple-system, "Segoe UI", Georgia, serif;
|
||||
color: var(--ink);
|
||||
background: var(--paper);
|
||||
}
|
||||
.doc-banner {
|
||||
background: linear-gradient(160deg, #312e81 0%, #4f46e5 60%, #7c74f0 100%);
|
||||
color: #eef0fe;
|
||||
margin: -1.5rem -1.25rem 1.5rem;
|
||||
padding: 1.1rem 1.25rem 1rem;
|
||||
}
|
||||
.doc-banner a { color: #fff; }
|
||||
.doc-banner .kicker {
|
||||
letter-spacing: 0.12em;
|
||||
text-transform: uppercase;
|
||||
font: 700 10px system-ui, sans-serif;
|
||||
opacity: 0.8;
|
||||
}
|
||||
.doc-banner h1 { margin: 0.25rem 0 0; font-size: 1.45rem; color: #fff; }
|
||||
h1, h2, h3, h4 { color: var(--accent-deep); page-break-after: avoid; }
|
||||
h1 { font-size: 1.7rem; }
|
||||
h2 {
|
||||
font-size: 1.2rem;
|
||||
border-bottom: 2px solid var(--accent);
|
||||
padding-bottom: 0.2rem;
|
||||
margin-top: 1.6rem;
|
||||
}
|
||||
h3 { font-size: 1.05rem; color: var(--accent); }
|
||||
a { color: var(--accent); }
|
||||
p, li { orphans: 3; widows: 3; }
|
||||
code {
|
||||
font-family: ui-monospace, Menlo, Consolas, monospace;
|
||||
font-size: 0.86em;
|
||||
background: var(--accent-soft);
|
||||
color: var(--accent-deep);
|
||||
padding: 0.08em 0.28em;
|
||||
border-radius: 4px;
|
||||
}
|
||||
pre, div.sourceCode, div.sourceCode pre {
|
||||
background: var(--code-bg) !important;
|
||||
color: var(--code-fg) !important;
|
||||
padding: 0.85rem 1rem;
|
||||
border-radius: 10px;
|
||||
overflow: auto;
|
||||
font-size: 0.78rem;
|
||||
line-height: 1.4;
|
||||
page-break-inside: avoid;
|
||||
}
|
||||
pre code { background: transparent; color: inherit; padding: 0; }
|
||||
#title-block-header, header#title-block-header, h1.title { display: none; }
|
||||
.doc-banner + h1 { display: none; }
|
||||
table {
|
||||
border-collapse: collapse;
|
||||
width: 100%;
|
||||
margin: 0.8rem 0 1.2rem;
|
||||
font-size: 0.9rem;
|
||||
page-break-inside: avoid;
|
||||
}
|
||||
th, td { border: 1px solid var(--line); padding: 0.38rem 0.55rem; text-align: left; vertical-align: top; }
|
||||
th {
|
||||
background: var(--accent);
|
||||
color: #fff;
|
||||
font: 650 12px system-ui, sans-serif;
|
||||
}
|
||||
tr:nth-child(even) td { background: var(--accent-soft); }
|
||||
blockquote {
|
||||
margin: 1rem 0;
|
||||
padding: 0.4rem 0.9rem;
|
||||
border-left: 4px solid var(--accent);
|
||||
background: var(--accent-soft);
|
||||
color: var(--accent-deep);
|
||||
}
|
||||
img { max-width: 100%; height: auto; border-radius: 8px; page-break-inside: avoid; }
|
||||
hr { border: 0; border-top: 1px solid var(--line); }
|
||||
ul, ol { padding-left: 1.25rem; }
|
||||
nav.site { font: 13px system-ui, sans-serif; margin-bottom: 0.4rem; }
|
||||
.source-path { font: 11px ui-monospace, Menlo, monospace; color: var(--muted); }
|
||||
|
||||
@page {
|
||||
size: letter;
|
||||
margin: 0.65in 0.7in 0.8in 0.7in;
|
||||
@top-left {
|
||||
content: "Verae Time × Zapier";
|
||||
font: 700 8pt system-ui, sans-serif;
|
||||
color: #4f46e5;
|
||||
}
|
||||
@top-right {
|
||||
content: "zapier.georgelambert.org";
|
||||
font: 8pt system-ui, sans-serif;
|
||||
color: #6b7186;
|
||||
}
|
||||
@bottom-center {
|
||||
content: counter(page) " / " counter(pages);
|
||||
font: 8pt system-ui, sans-serif;
|
||||
color: #6b7186;
|
||||
}
|
||||
}
|
||||
@media print {
|
||||
html, body { background: #fff; max-width: none; padding: 0; }
|
||||
.doc-banner { margin: 0 0 1rem; border-radius: 8px; -webkit-print-color-adjust: exact; print-color-adjust: exact; }
|
||||
a { text-decoration: none; }
|
||||
th, tr:nth-child(even) td, pre, blockquote, code { -webkit-print-color-adjust: exact; print-color-adjust: exact; }
|
||||
}
|
||||
34
scripts/ensure-streams.sh
Executable file
34
scripts/ensure-streams.sh
Executable file
|
|
@ -0,0 +1,34 @@
|
|||
#!/usr/bin/env bash
|
||||
# Create product JetStream streams with replicas=3 on the Proxmox cluster.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/client.env"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:/usr/local/bin:$PATH"
|
||||
VMID="${1:-511}"
|
||||
sudo pct exec "$VMID" -- bash -lc "
|
||||
set -e
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
export NATS_URL=nats://10.10.10.21:4222
|
||||
if [[ ! -x /usr/local/bin/nats ]]; then
|
||||
apt-get install -y --no-install-recommends unzip >/dev/null
|
||||
curl -fsSL https://github.com/nats-io/natscli/releases/download/v0.1.6/nats-0.1.6-linux-amd64.zip -o /tmp/natscli.zip
|
||||
rm -rf /tmp/natscli && mkdir -p /tmp/natscli
|
||||
unzip -o /tmp/natscli.zip -d /tmp/natscli >/dev/null
|
||||
BIN=\$(find /tmp/natscli -type f -name nats | head -1)
|
||||
install -m 0755 \"\$BIN\" /usr/local/bin/nats
|
||||
fi
|
||||
add() {
|
||||
local name=\$1 subj=\$2
|
||||
nats stream info \"\$name\" >/dev/null 2>&1 && return 0
|
||||
nats stream add \"\$name\" --subjects=\"\$subj\" --replicas=3 --storage=file --retention=limits --discard=old --max-msgs=-1 --max-bytes=-1 --max-age=24h --dupe-window=2m --defaults
|
||||
}
|
||||
add ZAPIER_JOBS 'verae.zapier.jobs.watch'
|
||||
add ZAPIER_EVENTS 'verae.zapier.jobs.events'
|
||||
add ZAPIER_WEBHOOKS 'verae.zapier.webhooks.deliver'
|
||||
add ZAPIER_USAGE 'verae.zapier.usage'
|
||||
add VERAE_ARCHIVE 'verae.archive.>'
|
||||
nats stream ls
|
||||
"
|
||||
echo "streams ready on cluster (replicas=3)"
|
||||
84
scripts/exhaustive-ns1-study.sh
Executable file
84
scripts/exhaustive-ns1-study.sh
Executable file
|
|
@ -0,0 +1,84 @@
|
|||
#!/usr/bin/env bash
|
||||
# Exhaustive NS1 ladder on the *current* 8c/16G cluster with JetStream on ZFS.
|
||||
# Adds r=1 vs r=3, file vs memory, reconnect tax, UDP echo, MQTT gateway probe.
|
||||
# Does not tmpfs (product streams stay). Must run on NS1.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:/usr/local/bin:$PATH"
|
||||
HOST="$(hostname -f 2>/dev/null || hostname)"
|
||||
case "$HOST" in
|
||||
NS1.GEORGELAMBERT.ORG|NS1|ns1.georgelambert.org|ns1) ;;
|
||||
*) echo "refusing: exhaustive-ns1-study.sh must run on NS1, got '$HOST'" >&2; exit 1 ;;
|
||||
esac
|
||||
|
||||
export EXHAUSTIVE=1
|
||||
export JS_EXTRA_MEMORY=1
|
||||
export COMPARE_DIR="${COMPARE_DIR:-$ROOT/results/20260912T051237Z}"
|
||||
STAMP="$(date -u +%Y%m%dT%H%M%SZ)"
|
||||
export BENCH_OUT="$ROOT/results/$STAMP"
|
||||
mkdir -p "$BENCH_OUT"
|
||||
|
||||
# MQTT gateway on nats-a only (vmbr1). Restored after.
|
||||
MQTT_CONF=/etc/nats/nats.conf
|
||||
enable_mqtt() {
|
||||
sudo pct exec 511 -- bash -lc '
|
||||
set -e
|
||||
f=/etc/nats/nats.conf
|
||||
grep -q "^mqtt {" "$f" && exit 0
|
||||
cat >> "$f" <<EOF
|
||||
|
||||
mqtt {
|
||||
host: 10.10.10.21
|
||||
port: 1883
|
||||
}
|
||||
EOF
|
||||
systemctl kill -s HUP nats-server || systemctl restart nats-server
|
||||
'
|
||||
sleep 2
|
||||
}
|
||||
disable_mqtt() {
|
||||
sudo pct exec 511 -- bash -lc '
|
||||
f=/etc/nats/nats.conf
|
||||
python3 - "$f" <<'"'"'PY'"'"'
|
||||
from pathlib import Path
|
||||
import re, sys
|
||||
p = Path(sys.argv[1])
|
||||
t = re.sub(r"\nmqtt \{[^}]*\}\n", "\n", p.read_text(), flags=re.S)
|
||||
p.write_text(t)
|
||||
PY
|
||||
systemctl kill -s HUP nats-server || true
|
||||
' || true
|
||||
}
|
||||
|
||||
enable_mqtt
|
||||
trap disable_mqtt EXIT
|
||||
|
||||
bash "$ROOT/scripts/study-on-ns1.sh"
|
||||
OUT="$BENCH_OUT"
|
||||
echo "exhaustive extras into $OUT"
|
||||
|
||||
# UDP echo: server in 511, client in 510
|
||||
sudo pct exec 511 -- bash -lc 'pkill -f "udp-probe.mjs server" >/dev/null 2>&1 || true'
|
||||
sudo pct exec 511 -- bash -c 'cat > /tmp/udp-probe.mjs' < "$ROOT/scripts/udp-probe.mjs"
|
||||
sudo pct exec 510 -- bash -c 'cat > /tmp/nats-lat/udp-probe.mjs' < "$ROOT/scripts/udp-probe.mjs"
|
||||
sudo pct exec 511 -- bash -lc 'setsid node /tmp/udp-probe.mjs server 9999 >/tmp/udp-echo.log 2>&1 < /dev/null &'
|
||||
sleep 1
|
||||
echo "=== udp-ping-1k-128 ===" | tee "$OUT/udp-ping-1k-128.txt"
|
||||
sudo pct exec 510 -- bash -lc 'node /tmp/nats-lat/udp-probe.mjs client 10.10.10.21 1000 9999 128' | tee -a "$OUT/udp-ping-1k-128.txt"
|
||||
sudo pct exec 511 -- bash -lc 'pkill -f "udp-probe.mjs server" || true'
|
||||
|
||||
# MQTT QoS0
|
||||
sudo pct exec 510 -- bash -lc '
|
||||
set -e
|
||||
cd /tmp/nats-lat
|
||||
if [[ ! -d node_modules/mqtt ]]; then npm install --no-audit --no-fund mqtt@10 >/dev/null; fi
|
||||
'
|
||||
sudo pct exec 510 -- bash -c 'cat > /tmp/nats-lat/mqtt-probe.mjs' < "$ROOT/scripts/mqtt-probe.mjs"
|
||||
echo "=== mqtt-qos0-5k-128 ===" | tee "$OUT/mqtt-qos0-5k-128.txt"
|
||||
sudo pct exec 510 -- bash -lc 'cd /tmp/nats-lat && node mqtt-probe.mjs mqtt://10.10.10.21:1883 5000 128' | tee -a "$OUT/mqtt-qos0-5k-128.txt" || echo '{"error":"mqtt probe failed"}' | tee -a "$OUT/mqtt-qos0-5k-128.txt"
|
||||
|
||||
python3 "$ROOT/scripts/build-optimal-report.py" "$OUT" \
|
||||
"$ROOT/results/20260912T045131Z" \
|
||||
"$ROOT/results/20260912T051237Z" \
|
||||
"$ROOT/results/20260912T053120Z"
|
||||
echo "exhaustive complete $OUT"
|
||||
132
scripts/latency.mjs
Normal file
132
scripts/latency.mjs
Normal file
|
|
@ -0,0 +1,132 @@
|
|||
#!/usr/bin/env node
|
||||
/**
|
||||
* Pub→sub round trip through the cluster (two connections).
|
||||
* Usage: NATS_URL=... node latency.mjs [count] [payloadBytes] [publishers] [ping|flood|reconnect]
|
||||
* ping = sequential publish-wait on persistent sockets (one-message RTT)
|
||||
* flood = publish the batch then drain (queueing under burst)
|
||||
* reconnect = connect, one publish, wait, close — measures handshake tax
|
||||
*/
|
||||
import { connect, headers } from "nats";
|
||||
|
||||
const url = process.env.NATS_URL || "nats://10.10.10.21:4222";
|
||||
const count = Number(process.argv[2] || 5000);
|
||||
const size = Number(process.argv[3] || 128);
|
||||
const pubs = Number(process.argv[4] || 1);
|
||||
const mode = process.argv[5] || "flood";
|
||||
const servers = url.split(",").map((s) => s.trim());
|
||||
const subject = `bench.lat.${process.pid}`;
|
||||
const payload = new Uint8Array(size);
|
||||
|
||||
function pct(sorted, p) {
|
||||
if (!sorted.length) return 0;
|
||||
const i = Math.min(sorted.length - 1, Math.floor((p / 100) * sorted.length));
|
||||
return sorted[i];
|
||||
}
|
||||
|
||||
const samples = [];
|
||||
if (mode === "reconnect") {
|
||||
const subNc = await connect({ servers, name: "lat-sub" });
|
||||
let resolveOne = null;
|
||||
const sub = subNc.subscribe(subject, { max: count });
|
||||
const consume = (async () => {
|
||||
for await (const m of sub) {
|
||||
const sent = Number(m.headers?.get("t") || 0);
|
||||
samples.push(Number(process.hrtime.bigint() / 1000n) - sent);
|
||||
resolveOne?.();
|
||||
}
|
||||
})();
|
||||
await subNc.flush();
|
||||
for (let i = 0; i < count; i++) {
|
||||
const got = new Promise((r) => {
|
||||
resolveOne = r;
|
||||
});
|
||||
const pubNc = await connect({ servers, name: `lat-re-${i}` });
|
||||
const h = headers();
|
||||
h.set("t", String(process.hrtime.bigint() / 1000n));
|
||||
pubNc.publish(subject, payload, { headers: h });
|
||||
await pubNc.flush();
|
||||
await got;
|
||||
await pubNc.close();
|
||||
}
|
||||
await consume;
|
||||
await subNc.close();
|
||||
} else if (mode === "ping") {
|
||||
const subNc = await connect({ servers, name: "lat-sub" });
|
||||
const pubNc = await connect({ servers, name: "lat-pub" });
|
||||
let resolveOne = null;
|
||||
const sub = subNc.subscribe(subject, { max: count });
|
||||
const consume = (async () => {
|
||||
for await (const m of sub) {
|
||||
const sent = Number(m.headers?.get("t") || 0);
|
||||
samples.push(Number(process.hrtime.bigint() / 1000n) - sent);
|
||||
resolveOne?.();
|
||||
}
|
||||
})();
|
||||
await subNc.flush();
|
||||
for (let i = 0; i < count; i++) {
|
||||
const got = new Promise((r) => {
|
||||
resolveOne = r;
|
||||
});
|
||||
const h = headers();
|
||||
h.set("t", String(process.hrtime.bigint() / 1000n));
|
||||
pubNc.publish(subject, payload, { headers: h });
|
||||
await got;
|
||||
}
|
||||
await consume;
|
||||
await pubNc.close();
|
||||
await subNc.close();
|
||||
} else {
|
||||
const subNc = await connect({ servers, name: "lat-sub" });
|
||||
const sub = subNc.subscribe(subject, { max: count });
|
||||
const done = (async () => {
|
||||
for await (const m of sub) {
|
||||
const sent = Number(m.headers?.get("t") || 0);
|
||||
if (sent) samples.push(Number(process.hrtime.bigint() / 1000n) - sent);
|
||||
}
|
||||
})();
|
||||
await subNc.flush();
|
||||
const per = Math.ceil(count / pubs);
|
||||
const publishers = [];
|
||||
for (let p = 0; p < pubs; p++) {
|
||||
publishers.push(
|
||||
(async () => {
|
||||
const nc = await connect({ servers, name: `lat-pub-${p}` });
|
||||
const n = p === pubs - 1 ? count - per * (pubs - 1) : per;
|
||||
for (let i = 0; i < n; i++) {
|
||||
const h = headers();
|
||||
h.set("t", String(process.hrtime.bigint() / 1000n));
|
||||
nc.publish(subject, payload, { headers: h });
|
||||
}
|
||||
await nc.flush();
|
||||
await nc.close();
|
||||
})(),
|
||||
);
|
||||
}
|
||||
await Promise.all(publishers);
|
||||
await done;
|
||||
await subNc.close();
|
||||
}
|
||||
|
||||
samples.sort((a, b) => a - b);
|
||||
const sum = samples.reduce((a, b) => a + b, 0);
|
||||
const us = (n) => `${(n / 1000).toFixed(3)}ms`;
|
||||
console.log(
|
||||
JSON.stringify({
|
||||
count: samples.length,
|
||||
pubs,
|
||||
size,
|
||||
mode,
|
||||
min_us: samples[0],
|
||||
avg_us: Math.round(sum / samples.length),
|
||||
p50_us: pct(samples, 50),
|
||||
p90_us: pct(samples, 90),
|
||||
p99_us: pct(samples, 99),
|
||||
max_us: samples[samples.length - 1],
|
||||
min: us(samples[0]),
|
||||
avg: us(sum / samples.length),
|
||||
p50: us(pct(samples, 50)),
|
||||
p90: us(pct(samples, 90)),
|
||||
p99: us(pct(samples, 99)),
|
||||
max: us(samples[samples.length - 1]),
|
||||
}),
|
||||
);
|
||||
62
scripts/lib-ct.sh
Executable file
62
scripts/lib-ct.sh
Executable file
|
|
@ -0,0 +1,62 @@
|
|||
# shellcheck shell=bash
|
||||
# Shared LXC bootstrap for NS1 Proxmox. Does not generate SSH keys if one exists.
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:$PATH"
|
||||
|
||||
ct_require_proxmox() {
|
||||
if [[ ! -d /etc/pve/nodes ]]; then
|
||||
echo "not a Proxmox host" >&2
|
||||
return 1
|
||||
fi
|
||||
command -v pct >/dev/null || { echo "pct missing" >&2; return 1; }
|
||||
}
|
||||
|
||||
ct_ensure() {
|
||||
local vmid="$1" hostname="$2" ip="$3"
|
||||
if [[ ! -f "$TEMPLATE" ]]; then
|
||||
echo "missing template $TEMPLATE" >&2
|
||||
return 1
|
||||
fi
|
||||
if ! sudo pct status "$vmid" >/dev/null 2>&1; then
|
||||
echo "pct create $vmid $hostname $ip/24"
|
||||
sudo pct create "$vmid" "$TEMPLATE" \
|
||||
--hostname "$hostname" \
|
||||
--memory "$MEMORY" --cores "$CORES" --swap 256 \
|
||||
--net0 "name=eth0,bridge=${BRIDGE},ip=${ip}/24,gw=${GW},type=veth" \
|
||||
--rootfs "${STORAGE}:${DISK}" \
|
||||
--unprivileged 1 --onboot 1 --nameserver "$DNS" \
|
||||
--features nesting=1 \
|
||||
--ostype ubuntu
|
||||
else
|
||||
echo "CT $vmid already exists"
|
||||
fi
|
||||
sudo pct start "$vmid" 2>/dev/null || true
|
||||
local i
|
||||
for i in $(seq 1 40); do
|
||||
sudo pct exec "$vmid" -- true 2>/dev/null && return 0
|
||||
sleep 2
|
||||
done
|
||||
echo "CT $vmid did not start" >&2
|
||||
return 1
|
||||
}
|
||||
|
||||
ct_bootstrap_user() {
|
||||
local vmid="$1"
|
||||
local pub=""
|
||||
[[ -f "$HOME/.ssh/id_ed25519.pub" ]] && pub="$(cat "$HOME/.ssh/id_ed25519.pub")"
|
||||
[[ -z "$pub" && -f "$HOME/.ssh/authorized_keys" ]] && pub="$(head -1 "$HOME/.ssh/authorized_keys")"
|
||||
[[ -n "$pub" ]] || { echo "no ssh public key" >&2; return 1; }
|
||||
sudo pct exec "$vmid" -- bash -lc "
|
||||
set -e
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
apt-get update -qq
|
||||
apt-get install -y --no-install-recommends openssh-server sudo curl ca-certificates xz-utils tar
|
||||
id $USER_NAME >/dev/null 2>&1 || useradd -m -s /bin/bash $USER_NAME
|
||||
echo '$USER_NAME ALL=(ALL) NOPASSWD:ALL' >/etc/sudoers.d/90-$USER_NAME
|
||||
chmod 440 /etc/sudoers.d/90-$USER_NAME
|
||||
install -d -m 700 -o $USER_NAME -g $USER_NAME /home/$USER_NAME/.ssh
|
||||
grep -qxF '$pub' /home/$USER_NAME/.ssh/authorized_keys 2>/dev/null || echo '$pub' >>/home/$USER_NAME/.ssh/authorized_keys
|
||||
chown $USER_NAME:$USER_NAME /home/$USER_NAME/.ssh/authorized_keys
|
||||
chmod 600 /home/$USER_NAME/.ssh/authorized_keys
|
||||
systemctl enable --now ssh
|
||||
"
|
||||
}
|
||||
102
scripts/maximize-ns1-study.sh
Executable file
102
scripts/maximize-ns1-study.sh
Executable file
|
|
@ -0,0 +1,102 @@
|
|||
#!/usr/bin/env bash
|
||||
# Maximize nats LXC resources + RAM-disk JetStream, run the NS1 study, then
|
||||
# put product streams back on ZFS. Cores/RAM/max_mem stay raised.
|
||||
# Must run on NS1.GEORGELAMBERT.ORG.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:/usr/local/bin:$PATH"
|
||||
|
||||
HOST="$(hostname -f 2>/dev/null || hostname)"
|
||||
case "$HOST" in
|
||||
NS1.GEORGELAMBERT.ORG|NS1|ns1.georgelambert.org|ns1) ;;
|
||||
*)
|
||||
echo "refusing: maximize-ns1-study.sh must run on NS1, got '$HOST'" >&2
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
|
||||
CORES="${CORES:-8}"
|
||||
MEMORY="${MEMORY:-16384}"
|
||||
TMPFS_SIZE="${TMPFS_SIZE:-8G}"
|
||||
NATS_VMS=(511 512 513)
|
||||
ALL_VMS=(510 511 512 513)
|
||||
|
||||
apply_resources() {
|
||||
local v
|
||||
for v in "${ALL_VMS[@]}"; do
|
||||
echo "pct set $v --cores $CORES --memory $MEMORY"
|
||||
sudo pct set "$v" --cores "$CORES" --memory "$MEMORY"
|
||||
done
|
||||
}
|
||||
|
||||
patch_max_mem() {
|
||||
local v
|
||||
for v in "${NATS_VMS[@]}"; do
|
||||
sudo pct exec "$v" -- sed -i -E 's/max_mem:[[:space:]]*[0-9]+[MmGg]/max_mem: 8G/' /etc/nats/nats.conf
|
||||
sudo pct exec "$v" -- grep -n max_mem /etc/nats/nats.conf
|
||||
done
|
||||
}
|
||||
|
||||
mount_tmpfs() {
|
||||
local v
|
||||
for v in "${NATS_VMS[@]}"; do
|
||||
sudo pct exec "$v" -- bash -lc "
|
||||
set -e
|
||||
systemctl stop nats-server
|
||||
mkdir -p /var/lib/nats/jetstream
|
||||
if ! mountpoint -q /var/lib/nats/jetstream; then
|
||||
mount -t tmpfs -o size=${TMPFS_SIZE} nats-js /var/lib/nats/jetstream
|
||||
fi
|
||||
chown nats:nats /var/lib/nats/jetstream
|
||||
chmod 755 /var/lib/nats/jetstream
|
||||
systemctl start nats-server
|
||||
mount | grep jetstream
|
||||
"
|
||||
done
|
||||
}
|
||||
|
||||
unmount_tmpfs() {
|
||||
local v
|
||||
for v in "${NATS_VMS[@]}"; do
|
||||
sudo pct exec "$v" -- bash -lc '
|
||||
set -e
|
||||
systemctl stop nats-server || true
|
||||
if mountpoint -q /var/lib/nats/jetstream; then
|
||||
umount /var/lib/nats/jetstream
|
||||
fi
|
||||
mkdir -p /var/lib/nats/jetstream
|
||||
chown nats:nats /var/lib/nats/jetstream
|
||||
systemctl start nats-server
|
||||
'
|
||||
done
|
||||
}
|
||||
|
||||
wait_cluster() {
|
||||
local n=0
|
||||
until sudo pct exec 511 -- curl -fsS --max-time 2 http://127.0.0.1:8222/varz >/dev/null 2>&1; do
|
||||
n=$((n + 1))
|
||||
[[ $n -lt 30 ]] || { echo "nats-a varz not up" >&2; return 1; }
|
||||
sleep 1
|
||||
done
|
||||
sleep 2
|
||||
}
|
||||
|
||||
restore_durable() {
|
||||
echo "restoring ZFS JetStream (product streams)"
|
||||
unmount_tmpfs
|
||||
wait_cluster
|
||||
bash "$ROOT/scripts/ensure-streams.sh" || true
|
||||
}
|
||||
|
||||
apply_resources
|
||||
patch_max_mem
|
||||
mount_tmpfs
|
||||
wait_cluster
|
||||
trap restore_durable EXIT
|
||||
|
||||
export JS_EXTRA_MEMORY=1
|
||||
export COMPARE_DIR="${COMPARE_DIR:-$ROOT/results/20260912T051237Z}"
|
||||
# recorded in host-before by appending after dump starts — study script reads pct config live
|
||||
bash "$ROOT/scripts/study-on-ns1.sh"
|
||||
|
||||
echo "maximize study finished; trap will restore ZFS jetstream"
|
||||
33
scripts/mqtt-probe.mjs
Executable file
33
scripts/mqtt-probe.mjs
Executable file
|
|
@ -0,0 +1,33 @@
|
|||
#!/usr/bin/env node
|
||||
/** MQTT QoS0 publish rate against nats-server MQTT gateway. */
|
||||
import mqtt from "mqtt";
|
||||
|
||||
const url = process.argv[2] || "mqtt://10.10.10.21:1883";
|
||||
const count = Number(process.argv[3] || 5000);
|
||||
const size = Number(process.argv[4] || 128);
|
||||
const payload = Buffer.alloc(size, 9);
|
||||
const topic = `bench/mqtt/${process.pid}`;
|
||||
|
||||
const c = mqtt.connect(url, { reconnectPeriod: 0, connectTimeout: 5000 });
|
||||
await new Promise((res, rej) => {
|
||||
c.on("connect", res);
|
||||
c.on("error", rej);
|
||||
});
|
||||
const t0 = process.hrtime.bigint();
|
||||
for (let i = 0; i < count; i++) {
|
||||
await new Promise((res, rej) => c.publish(topic, payload, { qos: 0 }, (err) => (err ? rej(err) : res())));
|
||||
}
|
||||
const ns = Number(process.hrtime.bigint() - t0);
|
||||
c.end(true);
|
||||
const sec = ns / 1e9;
|
||||
console.log(
|
||||
JSON.stringify({
|
||||
mode: "mqtt-qos0",
|
||||
count,
|
||||
size,
|
||||
url,
|
||||
secs: Number(sec.toFixed(3)),
|
||||
pubs_per_sec: Math.round(count / sec),
|
||||
mb_per_sec: Number(((count * size) / sec / 1e6).toFixed(2)),
|
||||
}),
|
||||
);
|
||||
202
scripts/ns1-study-methodology.md
Normal file
202
scripts/ns1-study-methodology.md
Normal file
|
|
@ -0,0 +1,202 @@
|
|||
## 4. Study methodology
|
||||
|
||||
### 4.1 Question
|
||||
|
||||
On the NS1 test stand, what message **throughput** and **delay** does the three-node `verae` JetStream cluster deliver at several loads, and which part of the stack is the limiter for product traffic (jobs, events, webhooks, archive)?
|
||||
|
||||
### 4.2 Hypotheses (stated before the run)
|
||||
|
||||
1. **H1 — Core vs JetStream.** Fire-and-forget core NATS is at least an order of magnitude faster than JetStream **file + replicas=3**, because durable publish waits for a majority disk replica.
|
||||
2. **H2 — JetStream parallelism.** Adding publishers does **not** linearly increase JetStream write rate once the replica log is saturated.
|
||||
3. **H3 — Quiet delay.** Sequential pub→sub round trip on `vmbr1` is well under 1 ms p99 when the consumer is waiting.
|
||||
4. **H4 — Burst delay.** If publishers dump a batch before the subscriber drains, observed delay is **queueing time**, roughly linear in backlog, not in cluster hop count.
|
||||
5. **H5 — Payload.** Moving 128 B → 1 KiB lowers message rate and raises byte rate on core NATS; JetStream in this size band stays replica/fsync bound.
|
||||
|
||||
### 4.3 Independent variables (what we changed)
|
||||
|
||||
| Factor | Levels |
|
||||
|--------|--------|
|
||||
| Transport | Core NATS pub/sub vs JetStream file replicas=3 |
|
||||
| Publisher count | 1, 2, 4, 8 |
|
||||
| Subscriber count | 0 (JS publish-only), 1, 2, 4, 8 |
|
||||
| Message count | 1k, 5k, 10k, 20k, 50k, 100k, 200k (by ladder step) |
|
||||
| Payload | 128 B, 1024 B |
|
||||
| Delay mode | **ping** (publish, wait, repeat) vs **flood** (publish all, then drain) |
|
||||
|
||||
### 4.4 Dependent variables (what we recorded)
|
||||
|
||||
| Metric | Instrument | Unit |
|
||||
|--------|------------|------|
|
||||
| Publish rate | `nats bench` 0.1.6 Pub stats | msgs/s, MB/s |
|
||||
| Subscribe rate | `nats bench` Sub stats | msgs/s, MB/s |
|
||||
| Aggregate | `nats bench` NATS Pub/Sub stats | msgs/s (fan-out counts both sides) |
|
||||
| Publisher spread | nats min/avg/max **msgs/s** | not delay |
|
||||
| One-way-ish RTT | `latency.mjs` header timestamp | min, avg, p50, p90, p99, max |
|
||||
| Host load | `/proc/loadavg` before and after | load average |
|
||||
| Broker counters | `http://127.0.0.1:8222/varz` inside each nats LXC | connections, in/out msgs, cpu, mem |
|
||||
|
||||
**Important:** nats CLI 0.1.6 min/avg/max are **rate spread across publishers**, not microseconds of delay. Delay is only `latency.mjs`.
|
||||
|
||||
### 4.5 Controls and constants
|
||||
|
||||
- Cluster name `verae`, three routes, client `:4222`, cluster `:6222`, monitor loopback `:8222`.
|
||||
- Client URL always the three-node list on `vmbr1` (never host `127.0.0.1:4222`, never `vmbr0`).
|
||||
- Bench client is LXC **510**, not a nats-* server.
|
||||
- JetStream bench stream name `benchstream`, **file** storage, **replicas=3**, deleted between JS loads (`nats stream rm --force`) so names do not collide.
|
||||
- Product streams were **not** the bench target (no load test on `ZAPIER_*` / `VERAE_ARCHIVE`).
|
||||
- No TLS, no nkeys, no account isolation (isolation is `vmbr1`).
|
||||
- Same nats CLI version (0.1.6) and `nats@2` Node client as the first ladder.
|
||||
|
||||
### 4.6 Procedure
|
||||
|
||||
1. Confirm this script is executing on **NS1.GEORGELAMBERT.ORG**. Refuse otherwise.
|
||||
2. Snapshot host load, memory, LXC configs, and each nats `varz`.
|
||||
3. From NS1, `pct exec 510` the core ladder (1p1s, 4p4s, 8p8s at 128 B; 4p4s at 1 KiB).
|
||||
4. Delete `benchstream`; JS ladder (1p, 4p, 4p×1 KiB, 2p2s pull) at replicas=3 file.
|
||||
5. Copy `latency.mjs` into 510; ping then flood at several batch sizes.
|
||||
6. Snapshot host/`varz` again.
|
||||
7. Parse logs on **this host**; draw charts; write HTML and PDF on **this host**.
|
||||
|
||||
No publish, subscribe, chart, or PDF process runs on the operator laptop for this study.
|
||||
|
||||
### 4.7 Instrumentation path
|
||||
|
||||
```text
|
||||
[NS1 host 70.88.205.138]
|
||||
study-on-ns1.sh (bash + python3)
|
||||
|
|
||||
| sudo pct exec 510
|
||||
v
|
||||
[LXC 510 verae-px-worker 10.10.10.20]
|
||||
nats bench / node latency.mjs
|
||||
|
|
||||
| NATS client protocol to
|
||||
v
|
||||
[LXC 511/512/513 10.10.10.21-23 :4222]
|
||||
nats-server -js cluster routes :6222
|
||||
```
|
||||
|
||||
The hypervisor issues the guest commands. The messages themselves never leave `vmbr1`.
|
||||
|
||||
### 4.8 Threats to validity
|
||||
|
||||
| Threat | Effect on numbers |
|
||||
|--------|-------------------|
|
||||
| **One physical host** | Three “replicas” share CPU, memory, and usually the same datastore. This measures process/LXC HA, not disk HA. |
|
||||
| **Shared load** | NS1 also runs Caddy, Forgejo, keep, fleet, portal, and other CTs. Load average during a run is part of the result, not noise to ignore. |
|
||||
| **Single bench client** | All publishers live in 510. Per-publisher rate spread is contention in that guest. |
|
||||
| **Short runs** | Seconds of traffic. No compaction, no multi-hour page-cache eviction, no snapshot during load. |
|
||||
| **No TLS/nkeys** | Production auth will cost CPU. Do not treat these rates as post-nkeys rates. |
|
||||
| **Fan-out aggregate** | Core aggregate msgs/s counts pub+sub. Do not compare that column to JetStream unique writes. |
|
||||
| **Flood ≠ RTT** | Mixing flood averages with ping p99 produces a fake “NATS is slow” story. |
|
||||
| **Lab only** | Not a Zapier HTTPS bench and not live `api.veraetime.net`. |
|
||||
|
||||
### 4.9 Ethics / safety
|
||||
|
||||
Bench uses throwaway subjects (`bench.core.*`, `bench.js.*`, `bench.lat.*`) and a throwaway stream. It does not purge product streams. Zapier cloud has no NATS socket.
|
||||
|
||||
---
|
||||
|
||||
## 5. Suggestions for fine-tuning
|
||||
|
||||
These follow from the method and from the first ladder on this stand (JetStream ~16k durable 128 B pubs/s; ping ~0.3 ms; flood hundreds of ms). Apply in order of leverage. Re-run **this NS1 study** after each change so the delta is measured the same way.
|
||||
|
||||
### 5.1 Treat JetStream as the product limiter
|
||||
|
||||
Product jobs/events/webhooks/archive are durable. Tuning core NATS to 2M msgs/s will not move a timestamp Zap. Put effort into **replica write path** and **consumer lag**, not core fan-out.
|
||||
|
||||
### 5.2 Split storage class by stream
|
||||
|
||||
| Stream | Suggested store | Why |
|
||||
|--------|-----------------|-----|
|
||||
| `ZAPIER_JOBS` | file, r=3 | Work queue; lose-a-job is bad |
|
||||
| `ZAPIER_EVENTS` | file r=3, or memory r=3 if events are rebuildable from job status | Hot waiters; measure both |
|
||||
| `ZAPIER_WEBHOOKS` | file, r=3, workqueue | HTTPS to Zapier is the slow consumer |
|
||||
| `ZAPIER_USAGE` | file, r=3, limits + max-age | Telemetry |
|
||||
| `VERAE_ARCHIVE` | file, r=3, on the **best disk** | Puts are larger and must survive |
|
||||
|
||||
Try `ZAPIER_EVENTS` as memory store in a maintenance window and re-run only the JS + ping/flood steps. If ping stays ~0.3 ms and durable events still ack at a higher rate, keep it; if a CT restart drops in-flight waiters, revert.
|
||||
|
||||
### 5.3 Give JetStream real disks
|
||||
|
||||
Today r=3 on three LXC guests on **one Proxmox host** is three files, one failure domain.
|
||||
|
||||
- Bind-mount a distinct SSD/NVMe (or ZFS dataset with its own vdev) into each nats LXC `store_dir`.
|
||||
- Set `sync: always` only on archive if you need it; default sync is often enough for jobs and is faster. Measure.
|
||||
- Do not put JetStream `store_dir` on the same busy rootfs as Forgejo/Caddy if we can avoid it.
|
||||
- When moving to three metal boxes: same configs, private NIC, one disk (or mirror) **per node**. That is the first change that makes r=3 mean “two boxes can die.”
|
||||
|
||||
### 5.4 Isolate the nats CTs from the rest of NS1
|
||||
|
||||
Host load on this box is often already several. Pin:
|
||||
|
||||
- `nats-a/b/c`: dedicated cores, no steal from keep/fleet Node processes.
|
||||
- Memory high enough that file-backed streams stay cache-hot for the working set.
|
||||
- `cpuunits` / cpuset in `pct config` so a Zapier-facing Node GC pause does not stall fsync.
|
||||
|
||||
Re-run this study after pinning; H1/H2 should move more than ping.
|
||||
|
||||
### 5.5 Consumer and mailbox tuning (delay H4)
|
||||
|
||||
Flood delay is backlog / consume_rate. Fine-tune the **waiters**, not the broker RTT.
|
||||
|
||||
- `jobs.events` and `webhooks.deliver`: raise `max_ack_pending` so a slow HTTPS hook does not stall the whole consumer; cap it so a poison message cannot unbounded-buffer RAM.
|
||||
- Pull consumers: larger batch, shorter `expires`, more pullers horizontally (fleet replica floors) instead of one fat subscriber.
|
||||
- Middleware should **not** flood-publish then wait; it already does per-job publish. Keep that. The flood test is the outage profile when a consumer is stopped.
|
||||
- Alert on **consumer lag** (pending + ack pending) from JetStream, not on ping RTT.
|
||||
|
||||
### 5.6 Publisher-side batching in middleware
|
||||
|
||||
A timestamp job is one small JSON. 16k msgs/s is ample. Still:
|
||||
|
||||
- Avoid per-byte publishes; one message per job/event.
|
||||
- Reuse NATS connections (connection churn showed up as publisher spread in the core 4p/8p runs).
|
||||
- Idempotent `msg id` / duplicate window sized to Verae retry window, not default-only.
|
||||
|
||||
### 5.7 nats-server knobs worth measuring (A/B with this script)
|
||||
|
||||
| Knob | Why try it |
|
||||
|------|------------|
|
||||
| `max_payload` | Keep default unless archive puts grow |
|
||||
| `write_deadline` | Slow consumer protection for webhooks |
|
||||
| `max_pending` | Bound memory on a stuck Zapier hook |
|
||||
| `max_connections` | Fleet workers + keep + middleware |
|
||||
| JetStream `max_file_store` / `max_memory_store` | Prevent one stream from filling the CT |
|
||||
| `max_outstanding_catchup` | Replica restart after a nats-c blip |
|
||||
| GOMAXPROCS = LXC cores | Do not overthread a 2-core CT |
|
||||
|
||||
Change **one** knob, re-run `study-on-ns1.sh`, compare JetStream 1p 128 B and ping p99.
|
||||
|
||||
### 5.8 Network
|
||||
|
||||
- Keep NATS off `vmbr0`. No change.
|
||||
- When on metal: dedicated NIC or VLAN for cluster `:6222` vs client `:4222` if possible (replication vs client load).
|
||||
- Check virtio queue counts on the LXC nics if core 1 KiB byte rate plateaus.
|
||||
|
||||
### 5.9 Security cost (when nkeys/mTLS flip)
|
||||
|
||||
`verae-nats-accounts` is still a sketch. Enabling accounts will add CPU on publish. Budget: re-run this exact study **after** creds are in every `NATS_URL`, and accept a drop on both core and JS. Do not flip without that measurement.
|
||||
|
||||
### 5.10 Operational fine-tuning (lag, not peak msgs/s)
|
||||
|
||||
1. Scrape `varz` / `jsz` from the host over `vmbr1` (not public). Monitor loopback `:8222` is invisible to Prometheus on NS1 unless we add a host-side proxy on `10.10.10.21:8222` bound only to `vmbr1`.
|
||||
2. Keep replica floors for webhook-deliver and job-poller — they are the flood defense.
|
||||
3. Backup/restore drill of JetStream **during idle**, then a short JS 1p run to see catchup cost.
|
||||
4. A 15–30 minute soak (not in this ladder) for page cache and compaction; add that as a third study when disks are dedicated.
|
||||
|
||||
### 5.11 What not to tune
|
||||
|
||||
- Do not chase core 8p8s aggregate. It is fan-out on a lab bridge.
|
||||
- Do not treat flood 400 ms as “cluster RTT.” Fix consumers.
|
||||
- Do not load-test on `ZAPIER_*` streams.
|
||||
- Do not bind client NATS to `0.0.0.0` on `vmbr0`.
|
||||
|
||||
### 5.12 Recommended next experiments (same method, one change each)
|
||||
|
||||
1. CPU pin nats-a/b/c → re-run JS 1p + ping.
|
||||
2. `ZAPIER_EVENTS`-shaped memory stream vs file (throwaway stream, same flags as this JS ladder).
|
||||
3. Distinct `store_dir` disks per node.
|
||||
4. nkeys on, same ladder.
|
||||
5. Three hardware boxes, same `cluster.env` IPs updated.
|
||||
|
||||
Each experiment should produce a new `results/<utc>/` on NS1 and a new progress-repo report so we can diff H1–H5 instead of arguing from memory.
|
||||
22
scripts/status.sh
Executable file
22
scripts/status.sh
Executable file
|
|
@ -0,0 +1,22 @@
|
|||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/cluster.env"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:$PATH"
|
||||
printf '%s\n' "$NODES" | awk 'NF==3 {print}' | while read -r vmid name ip; do
|
||||
st="$(sudo pct status "$vmid" 2>/dev/null || echo missing)"
|
||||
js="$(sudo pct exec "$vmid" -- curl -fsS --max-time 2 http://127.0.0.1:8222/varz 2>/dev/null || echo '{}')"
|
||||
echo "$vmid $name $ip $st"
|
||||
python3 -c "
|
||||
import json,sys
|
||||
try:
|
||||
d=json.loads(sys.argv[1])
|
||||
except Exception:
|
||||
print(' nats down')
|
||||
raise SystemExit
|
||||
print(' server_name', d.get('server_name'), 'cluster', (d.get('cluster') or {}).get('name'), 'routes', len((d.get('cluster') or {}).get('urls') or d.get('connect_urls') or []))
|
||||
print(' jetstream', bool(d.get('jetstream')), 'port', d.get('port'), 'host', d.get('host'))
|
||||
" "$js" 2>/dev/null || echo " nats down"
|
||||
done
|
||||
echo "host loopback still: $(ss -lnt | grep '127.0.0.1:4222' && echo up || echo down)"
|
||||
97
scripts/study-on-ns1.sh
Executable file
97
scripts/study-on-ns1.sh
Executable file
|
|
@ -0,0 +1,97 @@
|
|||
#!/usr/bin/env bash
|
||||
# Full message-speed study. Must run ON NS1.GEORGELAMBERT.ORG (70.88.205.138).
|
||||
# Orchestration, nats bench (via pct into LXC 510), charts, HTML, and PDF all
|
||||
# happen on this host. The laptop is not in the measurement path.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:/usr/local/bin:$PATH"
|
||||
|
||||
HOST="$(hostname -f 2>/dev/null || hostname)"
|
||||
case "$HOST" in
|
||||
NS1.GEORGELAMBERT.ORG|NS1|ns1.georgelambert.org|ns1) ;;
|
||||
*)
|
||||
echo "refusing: study-on-ns1.sh must run on NS1.GEORGELAMBERT.ORG (70.88.205.138), got '$HOST'" >&2
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/client.env"
|
||||
CLIENT_VMID="${CLIENT_VMID:-510}"
|
||||
STAMP="$(date -u +%Y%m%dT%H%M%SZ)"
|
||||
OUT="${BENCH_OUT:-$ROOT/results/$STAMP}"
|
||||
mkdir -p "$OUT"
|
||||
|
||||
dump_env() {
|
||||
local tag="$1"
|
||||
local f="$OUT/host-$tag.txt"
|
||||
{
|
||||
echo "execution_host=NS1.GEORGELAMBERT.ORG"
|
||||
echo "execution_ip=70.88.205.138"
|
||||
echo "hostname=$(hostname)"
|
||||
echo "utc=$(date -u +%Y-%m-%dT%H:%M:%SZ)"
|
||||
echo "whoami=$(whoami)"
|
||||
echo "pwd=$(pwd)"
|
||||
echo "uname=$(uname -a)"
|
||||
echo "nproc=$(nproc)"
|
||||
echo "loadavg=$(cat /proc/loadavg)"
|
||||
echo "client_vmid=$CLIENT_VMID"
|
||||
echo "nats_url=$NATS_URL"
|
||||
echo "js_extra_memory=${JS_EXTRA_MEMORY:-0}"
|
||||
echo "compare_dir=${COMPARE_DIR:-}"
|
||||
echo "--- nats 511 max_mem ---"
|
||||
sudo pct exec 511 -- grep max_mem /etc/nats/nats.conf || true
|
||||
echo "--- nats 511 jetstream mount ---"
|
||||
sudo pct exec 511 -- mount | grep jetstream || echo "jetstream on rootfs"
|
||||
echo "--- free ---"
|
||||
free -h
|
||||
echo "--- pct list ---"
|
||||
sudo pct list
|
||||
for v in 510 511 512 513; do
|
||||
echo "--- pct config $v ---"
|
||||
sudo pct config "$v" | grep -E '^(hostname|cores|memory|swap|rootfs|mp|net)' || true
|
||||
done
|
||||
} >"$f"
|
||||
python3 - "$OUT" "$tag" <<'PY'
|
||||
import json, sys, urllib.request
|
||||
from pathlib import Path
|
||||
out, tag = Path(sys.argv[1]), sys.argv[2]
|
||||
nodes = []
|
||||
for vmid, name in (("511", "nats-a"), ("512", "nats-b"), ("513", "nats-c")):
|
||||
raw = ""
|
||||
try:
|
||||
import subprocess
|
||||
raw = subprocess.check_output(
|
||||
["sudo", "pct", "exec", vmid, "--", "curl", "-fsS", "--max-time", "3", "http://127.0.0.1:8222/varz"],
|
||||
text=True,
|
||||
)
|
||||
d = json.loads(raw)
|
||||
nodes.append({
|
||||
"vmid": vmid,
|
||||
"name": name,
|
||||
"server_name": d.get("server_name"),
|
||||
"host": d.get("host"),
|
||||
"port": d.get("port"),
|
||||
"connections": d.get("connections"),
|
||||
"in_msgs": d.get("in_msgs"),
|
||||
"out_msgs": d.get("out_msgs"),
|
||||
"in_bytes": d.get("in_bytes"),
|
||||
"out_bytes": d.get("out_bytes"),
|
||||
"cpu": d.get("cpu"),
|
||||
"cores": d.get("cores"),
|
||||
"mem": d.get("mem"),
|
||||
"jetstream": bool(d.get("jetstream")),
|
||||
})
|
||||
except Exception as e:
|
||||
nodes.append({"vmid": vmid, "name": name, "error": str(e)})
|
||||
(out / f"varz-{tag}.json").write_text(json.dumps(nodes, indent=2) + "\n", encoding="utf-8")
|
||||
PY
|
||||
}
|
||||
|
||||
echo "NS1 study $STAMP out=$OUT"
|
||||
dump_env before
|
||||
BENCH_OUT="$OUT" CLIENT_VMID="$CLIENT_VMID" bash "$ROOT/scripts/bench.sh"
|
||||
dump_env after
|
||||
python3 "$ROOT/scripts/build-ns1-study-report.py" "$OUT" "${COMPARE_DIR:-}"
|
||||
echo "NS1 study complete $OUT"
|
||||
ls -la "$OUT"/nats-cluster-bench-ns1.* "$OUT"/charts 2>/dev/null || ls -la "$OUT"
|
||||
54
scripts/test.sh
Executable file
54
scripts/test.sh
Executable file
|
|
@ -0,0 +1,54 @@
|
|||
#!/usr/bin/env bash
|
||||
# Local syntax check always. Live cluster check when pct is present.
|
||||
set -euo pipefail
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
bash -n "$ROOT/scripts/lib-ct.sh"
|
||||
bash -n "$ROOT/scripts/create-cluster.sh"
|
||||
bash -n "$ROOT/scripts/status.sh"
|
||||
bash -n "$ROOT/scripts/bench.sh"
|
||||
bash -n "$ROOT/scripts/study-on-ns1.sh"
|
||||
bash -n "$ROOT/scripts/maximize-ns1-study.sh"
|
||||
bash -n "$ROOT/scripts/exhaustive-ns1-study.sh"
|
||||
grep -q 'host: {{IP}}' "$ROOT/conf/nats.conf.tmpl"
|
||||
grep -qv '0.0.0.0' "$ROOT/conf/nats.conf.tmpl"
|
||||
if [[ ! -d /etc/pve/nodes ]]; then
|
||||
echo "OK (syntax; not on Proxmox)"
|
||||
exit 0
|
||||
fi
|
||||
# shellcheck disable=SC1091
|
||||
. "$ROOT/cluster.env"
|
||||
export PATH="/usr/sbin:/usr/bin:/bin:$PATH"
|
||||
mapfile -t rows < <(printf '%s\n' "$NODES" | awk 'NF==3 {print}')
|
||||
ready=0
|
||||
for row in "${rows[@]}"; do
|
||||
# shellcheck disable=SC2086
|
||||
set -- $row
|
||||
vmid=$1 name=$2 ip=$3
|
||||
js="$(sudo pct exec "$vmid" -- curl -fsS --max-time 3 http://127.0.0.1:8222/varz 2>/dev/null || true)"
|
||||
echo "$js" | grep -q '"jetstream"' && ready=$((ready + 1)) || echo "not ready $name"
|
||||
done
|
||||
[[ $ready -eq 3 ]] || { echo "cluster not fully up ($ready/3)" >&2; exit 1; }
|
||||
|
||||
# nats CLI on first node
|
||||
first="$(echo "${rows[0]}" | awk '{print $1}')"
|
||||
sudo pct exec "$first" -- bash -lc '
|
||||
set -e
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
if [[ ! -x /usr/local/bin/nats ]]; then
|
||||
apt-get install -y --no-install-recommends unzip >/dev/null
|
||||
curl -fsSL https://github.com/nats-io/natscli/releases/download/v0.1.6/nats-0.1.6-linux-amd64.zip -o /tmp/natscli.zip
|
||||
rm -rf /tmp/natscli && mkdir -p /tmp/natscli
|
||||
unzip -o /tmp/natscli.zip -d /tmp/natscli >/dev/null
|
||||
BIN=$(find /tmp/natscli /tmp -maxdepth 3 -type f -name nats | head -1)
|
||||
test -n "$BIN"
|
||||
install -m 0755 "$BIN" /usr/local/bin/nats
|
||||
fi
|
||||
IP=$(hostname -I | awk "{print \$1}")
|
||||
export NATS_URL=nats://$IP:4222
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin
|
||||
nats stream rm VERAE_PX_TEST --force >/dev/null 2>&1 || true
|
||||
nats stream add VERAE_PX_TEST --subjects="verae.px.test" --replicas=3 --storage=file --retention=limits --discard=old --max-msgs=-1 --max-bytes=-1 --max-age=1h --dupe-window=2m --defaults
|
||||
nats pub verae.px.test cluster-ok
|
||||
nats stream info VERAE_PX_TEST
|
||||
'
|
||||
echo "OK live cluster (3/3 + replicas=3 stream)"
|
||||
54
scripts/udp-probe.mjs
Executable file
54
scripts/udp-probe.mjs
Executable file
|
|
@ -0,0 +1,54 @@
|
|||
#!/usr/bin/env node
|
||||
/** UDP echo RTT. server: node udp-probe.mjs server [port]
|
||||
* client: node udp-probe.mjs client <host> <count> [port] [size] */
|
||||
import dgram from "node:dgram";
|
||||
|
||||
const mode = process.argv[2] || "server";
|
||||
const port = Number(process.argv[mode === "server" ? 3 : 5] || 9999);
|
||||
|
||||
if (mode === "server") {
|
||||
const s = dgram.createSocket("udp4");
|
||||
s.on("message", (msg, rinfo) => s.send(msg, rinfo.port, rinfo.address));
|
||||
s.bind(port, "0.0.0.0", () => console.log(JSON.stringify({ mode: "udp-server", port })));
|
||||
} else {
|
||||
const host = process.argv[3];
|
||||
const count = Number(process.argv[4] || 1000);
|
||||
const size = Number(process.argv[6] || 128);
|
||||
const sock = dgram.createSocket("udp4");
|
||||
const payload = Buffer.alloc(size, 7);
|
||||
const samples = [];
|
||||
let i = 0;
|
||||
const sendOne = () => {
|
||||
const t0 = process.hrtime.bigint();
|
||||
const once = (msg) => {
|
||||
sock.off("message", once);
|
||||
samples.push(Number(process.hrtime.bigint() - t0) / 1000);
|
||||
i += 1;
|
||||
if (i >= count) {
|
||||
samples.sort((a, b) => a - b);
|
||||
const us = (n) => `${(n / 1000).toFixed(3)}ms`;
|
||||
const pct = (p) => samples[Math.min(samples.length - 1, Math.floor((p / 100) * samples.length))];
|
||||
const sum = samples.reduce((a, b) => a + b, 0);
|
||||
console.log(
|
||||
JSON.stringify({
|
||||
mode: "udp-ping",
|
||||
count: samples.length,
|
||||
size,
|
||||
host,
|
||||
min: us(samples[0]),
|
||||
avg: us(sum / samples.length),
|
||||
p50: us(pct(50)),
|
||||
p99: us(pct(99)),
|
||||
max: us(samples[samples.length - 1]),
|
||||
p50_us: Math.round(pct(50)),
|
||||
p99_us: Math.round(pct(99)),
|
||||
}),
|
||||
);
|
||||
sock.close();
|
||||
} else sendOne();
|
||||
};
|
||||
sock.on("message", once);
|
||||
sock.send(payload, port, host);
|
||||
};
|
||||
sendOne();
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue