Skip to content

Commit 591b94f

Browse files
committed
bench: live per-column cell-verdict watcher (same/faster/slower vs baseline)
Reads the engine's per-egress-column partial snapshot (run_grid_streaming flushes results/snapshots/<gw>.json every 6 cells) off a running field box over ssh and prints per-cell throughput deltas (frontier RPS@10ms) vs a baseline snapshot, with append-only numbering so each cell is reported once. Gives early same/faster/slower insight during a multi-hour grid instead of waiting for the full run. Reusable for any gateway/version (e.g. the next busbar bump).
1 parent ca0ef20 commit 591b94f

2 files changed

Lines changed: 186 additions & 0 deletions

File tree

‎cell-verdict.py‎

Lines changed: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
1+
#!/usr/bin/env python3
2+
"""Per-cell live verdict: compare a running gateway's partial snapshot to a baseline, cell by cell.
3+
4+
Emits one verdict line per NEWLY-measured cell (throughput headline = frontier RPS at the 10ms p99
5+
bound, higher-is-better), so an operator watching a 36-cell grid gets "cell N: same/better/worse
6+
(was X now Y) - proceeding" the moment each cell lands, instead of waiting out the whole grid.
7+
8+
Numbering is append-only and persisted in --state (a JSON list of reported "ingress>egress" keys):
9+
a cell is numbered the first time it is seen measured, and never renumbered, so cell 1/2/3 are stable
10+
even as later cells complete out of matrix order. Re-run against the same --state to get only what is
11+
new since last call.
12+
13+
cell-verdict.py --baseline <snap.json> --new <partial.json> --state <state.json> [--count-only]
14+
15+
A cell counts as "measured" once it has a frontier RPS at the 10ms p99 bound; latency (added p50/p99,
16+
c=1 p99, cpu us/req) rides along as corroboration. Absences print their own reason rather than a zero.
17+
"""
18+
import argparse
19+
import json
20+
import os
21+
import sys
22+
23+
BOUND_US = 10_000 # headline: frontier rung whose p99 bound is 10ms
24+
SAME_BAND_PCT = 3.0 # |delta| < this -> "same"; else better/worse (throughput, higher better)
25+
26+
27+
def load(path):
28+
try:
29+
with open(path) as f:
30+
return json.load(f)
31+
except Exception:
32+
return None
33+
34+
35+
def cells_in_order(doc):
36+
"""Yield (ingress, egress, cell) in the snapshot's own upstream/ingress insertion order."""
37+
if not doc:
38+
return
39+
ups = (doc.get("matrix") or {}).get("upstreams") or {}
40+
for egress, u in ups.items():
41+
for ingress, cell in ((u or {}).get("cells") or {}).items():
42+
yield ingress, egress, cell
43+
44+
45+
def frontier_rps(cell, bound_us=BOUND_US):
46+
if not cell:
47+
return None
48+
for r in ((cell.get("perf") or {}).get("frontier") or []):
49+
if r.get("p99_bound_us") == bound_us:
50+
return r.get("rps")
51+
return None
52+
53+
54+
def perf_num(cell, field):
55+
return ((cell.get("perf") or {}) or {}).get(field)
56+
57+
58+
def cell_map(doc):
59+
return {f"{ing}>{eg}": c for ing, eg, c in cells_in_order(doc)}
60+
61+
62+
def measured_keys_in_order(doc):
63+
"""Keys of cells that have a headline frontier RPS, in snapshot order."""
64+
out = []
65+
for ing, eg, c in cells_in_order(doc):
66+
if frontier_rps(c) is not None:
67+
out.append(f"{ing}>{eg}")
68+
return out
69+
70+
71+
def verdict_of(old_rps, new_rps):
72+
if new_rps is None:
73+
return "unmeasured", None
74+
if old_rps is None or old_rps == 0:
75+
return "new", None
76+
pct = (new_rps - old_rps) / old_rps * 100.0
77+
if abs(pct) < SAME_BAND_PCT:
78+
return "same", pct
79+
return ("better" if pct > 0 else "worse"), pct
80+
81+
82+
def main():
83+
ap = argparse.ArgumentParser()
84+
ap.add_argument("--baseline", required=True)
85+
ap.add_argument("--new", required=True)
86+
ap.add_argument("--state", required=True)
87+
ap.add_argument("--count-only", action="store_true")
88+
a = ap.parse_args()
89+
90+
new_doc = load(a.new)
91+
measured = measured_keys_in_order(new_doc) if new_doc else []
92+
if a.count_only:
93+
print(len(measured))
94+
return
95+
96+
base = cell_map(load(a.baseline))
97+
newm = cell_map(new_doc)
98+
99+
state = load(a.state) or []
100+
reported = list(state)
101+
reported_set = set(reported)
102+
103+
fresh = [k for k in measured if k not in reported_set]
104+
for k in fresh:
105+
idx = len(reported) + 1
106+
oc, nc = base.get(k), newm.get(k)
107+
o_rps, n_rps = frontier_rps(oc), frontier_rps(nc)
108+
v, pct = verdict_of(o_rps, n_rps)
109+
o_lat, n_lat = perf_num(oc, "added_latency_p50_us"), perf_num(nc, "added_latency_p50_us")
110+
o_cpu, n_cpu = perf_num(oc, "cpu_us_per_request"), perf_num(nc, "cpu_us_per_request")
111+
112+
def r(x):
113+
return f"{x:,.0f}" if isinstance(x, (int, float)) else "—"
114+
115+
pcts = f" ({pct:+.1f}%)" if pct is not None else ""
116+
lat = ""
117+
if o_lat is not None or n_lat is not None:
118+
lat = f"; added-lat p50 {r(o_lat)}->{r(n_lat)}us"
119+
cpu = ""
120+
if o_cpu is not None or n_cpu is not None:
121+
cpu = f"; cpu {r(o_cpu)}->{r(n_cpu)}us/req"
122+
# Human relay line (the operator copies this to the user verbatim).
123+
print(f"RELAY: cell {idx} ({k}): {v} — RPS@10ms was {r(o_rps)} now {r(n_rps)}{pcts}{lat}{cpu}")
124+
reported.append(k)
125+
126+
with open(a.state, "w") as f:
127+
json.dump(reported, f)
128+
129+
# Machine tail so the caller knows counts without reparsing.
130+
print(f"COUNT: measured={len(measured)} reported={len(reported)} new={len(fresh)}")
131+
132+
133+
if __name__ == "__main__":
134+
main()

‎watch-cells.sh‎

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
#!/usr/bin/env bash
2+
# Wait for the NEXT measured cell(s) on a running field box, emit their verdict lines, then exit.
3+
#
4+
# Blocks (polling the box's partial snapshot over ssh) until at least one new cell has a headline
5+
# frontier RPS, or the run finishes, or a STOP file appears - then prints and returns. Designed to be
6+
# relaunched in a loop by an operator/agent: each return delivers the newly-completed cells so they
7+
# can be relayed one at a time, until RUN-DONE.
8+
#
9+
# IP=1.2.3.4 KEY=~/.cache/gateway-bench/gateway-bench-key.pem GW=busbar \
10+
# BASELINE=results/snapshots/result_busbar-151_....json \
11+
# STATE=/tmp/busbar.state.json WORK=/tmp/busbarwatch ./watch-cells.sh
12+
#
13+
# Exit: 0 new cells emitted (RELAY: lines) | 2 run finished (RUN-DONE) | 3 STOP file | 4 unreachable-too-long
14+
set -uo pipefail
15+
HERE="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
16+
: "${IP:?need IP}"; : "${KEY:?need KEY}"; : "${GW:?need GW}"; : "${BASELINE:?need BASELINE}"
17+
: "${STATE:?need STATE}"; : "${WORK:?need WORK}"
18+
POLL="${POLL:-45}"; MAXUNREACH="${MAXUNREACH:-40}"
19+
SSH=(ssh -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o ConnectTimeout=12 -i "$KEY")
20+
mkdir -p "$WORK"; NEW="$WORK/${GW}.partial.json"
21+
unreach=0
22+
while :; do
23+
[ -f "$WORK/STOP" ] && { echo "STOP file present"; exit 3; }
24+
# newest partial snapshot path on the box (engine writes/updates it as cells land)
25+
remote="$("${SSH[@]}" ubuntu@"$IP" "ls -t ~/benchmarking/results/snapshots/result_${GW}_*.json 2>/dev/null | head -1" 2>/dev/null)"
26+
done_rc="$("${SSH[@]}" ubuntu@"$IP" 'cat ~/benchmarking/.run-done 2>/dev/null' 2>/dev/null)"
27+
if [ -z "$remote" ] && [ -z "$done_rc" ]; then
28+
# box up but nothing written yet, or a transient ssh miss
29+
if ! "${SSH[@]}" ubuntu@"$IP" true 2>/dev/null; then
30+
unreach=$((unreach+1)); [ "$unreach" -ge "$MAXUNREACH" ] && { echo "UNREACHABLE for $unreach polls"; exit 4; }
31+
else
32+
unreach=0
33+
fi
34+
sleep "$POLL"; continue
35+
fi
36+
unreach=0
37+
if [ -n "$remote" ]; then
38+
rsync -az --timeout=60 -e "${SSH[*]}" "ubuntu@$IP:$remote" "$NEW" 2>/dev/null || true
39+
fi
40+
out="$(python3 "$HERE/cell-verdict.py" --baseline "$HERE/$BASELINE" --new "$NEW" --state "$STATE" 2>/dev/null)"
41+
relay="$(printf '%s\n' "$out" | grep '^RELAY:' || true)"
42+
if [ -n "$relay" ]; then
43+
printf '%s\n' "$out"
44+
exit 0
45+
fi
46+
if [ -n "$done_rc" ]; then
47+
echo "RUN-DONE=$done_rc"
48+
printf '%s\n' "$out" # COUNT: line for final tally
49+
exit 2
50+
fi
51+
sleep "$POLL"
52+
done

0 commit comments

Comments
 (0)