#!/usr/bin/env python3 """client.py -- THE ENGINE-AGNOSTIC SWEEP: one barrier, levels, runs, spread, medians, the A-B-A anchor, CPU. What this file is, for a reader who opened it cold (SPEC Δ6, the bench-client lane §3.14; PLAN-v2 §2.6; M-9, N-4, N-11, N-12). Given ONE engine that is up and proven (arms.py boots it), this fires the same kind of callers at it that sweep.py fires at the seats, by the same mechanism, and turns its streams into the figures of PLAN-v2 §2.6: * `fire` -- every stream connects, then waits on ONE threading.Barrier; the barrier's action stamps time.monotonic() and every stream's clock starts there. "Arrivals together" is the workload's name (N-4). * `level` -- 1 warm-up and N scored runs of one (engine, prompt, level); a run with any failed stream is VOID (trial_void) and stays out of every median. EVERY run is printed, void ones included. A level with any VOID (gates.void_rollup: a failed request, a hidden token, the all-layers void) carries NO scored figure: `unscore_cell` moves its medians to `unscored_medians`, kept as a reading and never read as a result. * the medians are over the scored, non-void runs; `spread` = (max - min) / median of aggregate_wall over them (PLAN-v2 §2.6); a level with fewer than two such runs is thin. * after each run, outside the timed window: R-7x per stream, the hidden-token gate (gates.hidden_tokens), the finish-reason flags with the completion-token total, and the nonce-leak flag; after the level (its log span read), every retokenize-method surplus is replayed once with logprobs and the exact count printed beside it (`replay_hidden`; the bench-client fold's fairness I-1). * CPU (N-12): the client pins itself to the SMT siblings of the highest-numbered physical core (os.sched_setaffinity, read from /sys/devices/system/cpu/cpu*/topology) and records the pin; per run, utime + stime from /proc//stat for the client and every process of the engine's tree (stdlib, no pidstat). * the headline guard: a cell whose prompt is not headline_eligible (D's finish-reason read) prints `not headline: finished at tokens in D` on every aggregate it writes; a cell whose path is not `raw` can never feed a headline field. """ import glob import hashlib import json import os import threading import time import engine as E import framing import gates import nonces as N import parity import prereg import seatlib as L _CLK_TCK = os.sysconf("SC_CLK_TCK") if hasattr(os, "sysconf") else 100 #: The expected count a road passes when no parity count applies to its requests (TT road (d)). R7X_NOT_APPLICABLE = "n/a" # ------------------------------------------------------------------------------------------------ # # CPU: the pin and the per-process read # ------------------------------------------------------------------------------------------------ # def cpu_topology(sys_root="/sys/devices/system/cpu"): """{cpu: (core_id, siblings list text)} for every online cpu the kernel lists (package-local core ids).""" out = {} for path in glob.glob(os.path.join(sys_root, "cpu[0-9]*")): try: cpu = int(os.path.basename(path)[3:]) with open(os.path.join(path, "topology", "core_id"), encoding="utf-8") as fh: core = int(fh.read().strip()) with open(os.path.join(path, "topology", "thread_siblings_list"), encoding="utf-8") as fh: sib = fh.read().strip() except (OSError, ValueError): continue out[cpu] = (core, sib) return out def parse_cpu_list(text): cpus = set() for part in (text or "").split(","): part = part.strip() if "-" in part: a, b = part.split("-", 1) cpus.update(range(int(a), int(b) + 1)) elif part: cpus.add(int(part)) return sorted(cpus) def client_cpus(sys_root="/sys/devices/system/cpu"): """The SMT siblings of the highest-numbered physical core (its core_id), or [] when the topology is unread.""" topo = cpu_topology(sys_root) if not topo: return [] top = max(core for core, _ in topo.values()) sib = next(s for cpu, (core, s) in sorted(topo.items()) if core == top) return parse_cpu_list(sib) def pin_client(sys_root="/sys/devices/system/cpu"): """Pin THIS process (and the threads it starts) to client_cpus(); the record says what was pinned or why not.""" want = client_cpus(sys_root) rec = {"want": want, "rule": L.RULES["cpu_per_process"], "pinned": None, "engines": "launched on every CPU (engine.ALL_CPUS, released in each launch's preexec_fn): only the client " "is pinned"} if not want: rec["why"] = "the CPU topology was unread; the client is not pinned" return rec try: os.sched_setaffinity(0, set(want)) rec["pinned"] = sorted(os.sched_getaffinity(0)) except (OSError, AttributeError, ValueError) as e: rec["why"] = f"sched_setaffinity failed: {type(e).__name__}: {e}" return rec def proc_cpu_s(pid, proc="/proc"): """utime + stime of one process in seconds, or None when unreadable.""" try: with open(f"{proc}/{pid}/stat", encoding="utf-8", errors="replace") as fh: rest = fh.read().rsplit(")", 1)[1].split() return (int(rest[11]) + int(rest[12])) / float(_CLK_TCK) except (OSError, IndexError, ValueError): return None def proc_name(pid, proc="/proc"): argv = E.program_argv(E.proc_argv(pid)) return os.path.basename(argv[0]) if argv else str(pid) def cpu_snapshot(pids): return {int(p): proc_cpu_s(p) for p in pids} def cpu_delta(before, after, wall_s): """{name(pid): {cpu_s, cpu_pct_of_one_core}} over a run, from two snapshots; None where a read was missing.""" out = {} for pid, b in before.items(): a = after.get(pid) key = f"{proc_name(pid)}({pid})" if a is None or b is None: out[key] = {"cpu_s": None, "cpu_pct_of_one_core": None, "why": "unread before or after the run"} continue d = a - b out[key] = {"cpu_s": round(d, 3), "cpu_pct_of_one_core": round(100.0 * d / wall_s, 1) if wall_s else None} return out # ------------------------------------------------------------------------------------------------ # # one run # ------------------------------------------------------------------------------------------------ # def fire(engine, bodies, path_kind, url_path=None): """Every body at `engine`, released together behind ONE barrier. bodies = [(rec, body dict)]; returns (records, wall_s, release_stamp). A failure is a record, never an exception (law 4).""" import http.client n = len(bodies) release = {"t": None} barrier = threading.Barrier(n, action=lambda: release.__setitem__("t", time.monotonic())) path = url_path or (engine.raw_path() if path_kind == "raw" else engine.chat_path()) out = [None] * n def one(k, rec, body): data = json.dumps(body).encode("utf-8") conn = http.client.HTTPConnection("127.0.0.1", engine.port, timeout=E.STREAM_TIMEOUT_S) try: try: conn.connect() except OSError as e: try: barrier.wait(timeout=E.BARRIER_TIMEOUT_S) except threading.BrokenBarrierError: pass rec.update(fail_kind="transport", error=f"connect: {type(e).__name__}: {e}") return barrier.wait(timeout=E.BARRIER_TIMEOUT_S) t0 = release["t"] try: conn.request("POST", path, body=data, headers={"Content-Type": "application/json"}) resp = conn.getresponse() except (OSError, http.client.HTTPException) as e: rec.update(fail_kind="transport", error=f"{type(e).__name__}: {e}", t_end_s=time.monotonic() - t0) return rec["http_status"] = resp.status if resp.status != 200: txt = resp.read(600).decode("utf-8", "replace") rec.update(fail_kind="http_5xx" if resp.status >= 500 else "http_4xx", error=f"HTTP {resp.status}: {txt[:300]}", t_end_s=time.monotonic() - t0) return try: engine.decode_stream(resp, t0, rec) except (OSError, http.client.HTTPException, ValueError) as e: rec["error"] = f"the stream broke: {type(e).__name__}: {e}" rec["t_end_s"] = time.monotonic() - t0 except threading.BrokenBarrierError: rec.update(fail_kind="transport", error="the barrier broke before release") finally: conn.close() out[k] = rec threads = [threading.Thread(target=one, args=(k, rec, body), daemon=True) for k, (rec, body) in enumerate(bodies)] for t in threads: t.start() for t in threads: t.join(timeout=E.STREAM_TIMEOUT_S + E.BARRIER_TIMEOUT_S + 30) for k in range(n): if out[k] is None: out[k] = dict(bodies[k][0], fail_kind="transport", error="the stream thread never returned") out[k]["stream"] = k + 1 wall = max((r.get("t_end_s") or 0.0) for r in out) if out else 0.0 return out, round(wall, 4), release["t"] def classify(rec, engine, expected_count): """Status and fail kind from what the stream carried; a failed stream carries no speed (engine.unscore).""" reasons = [] if rec.get("fail_kind") in ("transport", "http_4xx", "http_5xx"): reasons.append(rec["fail_kind"]) elif rec.get("repeat_abort"): reasons.append("repeat_abort") else: if rec.get("finish_reason") is None: reasons.append("truncated_stream") if rec.get("completion_tokens") is None: reasons.append("no_usage") if expected_count == R7X_NOT_APPLICABLE: rec["r7x"] = ("n/a: this road's engine renders the prompt itself (Ollama's /v1/completions builds its " "GenerateRequest without Raw, openai/openai.go L812-871 at v0.34.4), so no parity count " "applies; the count it read is kept") else: why = parity.r7x(engine.prompt_count(rec), expected_count) if rec.get("completion_tokens") is not None and why: reasons.append("token_count_mismatch") rec["r7x"] = why rec["fail_reasons"] = reasons rec["status"] = "failed" if reasons else "ok" rec["fail_kind"] = reasons[0] if reasons else None E.finish_record(rec) if reasons: E.unscore(rec) return rec def run_figures(streams): """Per run: the ok streams' figures (seatlib.RULES), None where no ok stream has one.""" ok = [s for s in streams if s.get("status") == "ok"] dec = [s.get("decode_tok_s") for s in ok if s.get("decode_tok_s") is not None] ttft = [s.get("ttft_ms") for s in ok if s.get("ttft_ms") is not None] ends = [s["t_end_s"] for s in ok if s.get("t_end_s")] ctoks = sum(s.get("completion_tokens") or 0 for s in ok) srv = [((s.get("server_clock") or {}).get("decode_tok_s_server")) for s in ok] srv = [x for x in srv if x is not None] cached = [s.get("cached_tokens") for s in ok if s.get("cached_tokens") is not None] return {"streams_ok": len(ok), "streams_failed": len(streams) - len(ok), "decode_tok_s_p50": L.rnd(L.median(dec)), "ttft_ms_p50": L.rnd(L.median(ttft)), "ttft_ms_p95": L.rnd(L.p95_nearest_rank(ttft)), "aggregate_sum_tok_s": L.rnd(sum(dec)) if ok else None, "aggregate_wall_tok_s": L.rnd(ctoks / max(ends)) if ends else None, "completion_tokens_total": ctoks, "decode_tok_s_server_p50": L.rnd(L.median(srv)), "cached_tokens_max": max(cached) if cached else None, "rules": {"aggregate_wall": L.RULES["aggregate_wall"], "timing": L.RULES["engine_timing"]}} def run_once(engine, cell, n, warmup, strings, expected, nonce_end, samplers=None, cpu_pids=None, road=None, url_path=None, path_kind="raw", max_tokens=None, keep_text=False): """One run: bodies built from the run's nonce strings, fired behind one barrier, then the gates (outside the timed window). `strings` = [(nonce key, nonce text, the string)] for this run's requests. `max_tokens` is the registered 256 unless an arm registers its own (arm Q: the six-way's 8); `keep_text` keeps each stream's whole visible text (arm Q parses it; the pen reads it as model text).""" bodies = [] mt = prereg.MAX_TOKENS if max_tokens is None else int(max_tokens) rec_path = "v1-completions" if (road or {}).get("body") == "v1-completions" else path_kind for key, ntext, text in strings: rec = E.new_record(engine.key, engine.posture["key"], rec_path, key, ntext) if path_kind == "chat": body = engine.chat_row_body(text, max_tokens=mt) elif road and road.get("body") == "v1-completions": body = engine.v1_completions_body(text, max_tokens=mt) else: body = engine.raw_body(text, max_tokens=mt, **(road or {}).get("options", {})) bodies.append((rec, body)) power = ups = None if samplers: power, ups = samplers() pids = list(cpu_pids() if callable(cpu_pids) else (cpu_pids or [])) cpu0 = cpu_snapshot([os.getpid()] + pids) utc = L.utc_now() if power: power.start() ups.start() try: streams, wall, _rel = fire(engine, bodies, path_kind, url_path) finally: if power: power.halt() ups.halt() cpu1 = cpu_snapshot(list(cpu0)) # ---- outside the timed window: R-7x, the gates, the flags ---- for (key, _nt, text), s in zip(strings, streams): s["sent_sha256"] = sent_sha(text) classify(s, engine, expected.get(text) if isinstance(expected, dict) else expected) s["expected_prompt_tokens"] = expected.get(text) if isinstance(expected, dict) else expected s["server_clock"] = engine.server_clock(s) cached, why = engine.cached_tokens(s) s["cached_tokens"], s["cached_why"] = cached, why leak = N.nonce_leak(cached, nonce_end) if key is not None else None if leak: s["flags"].append(leak) if s["status"] == "ok": gates.hidden_tokens(s, engine) if not keep_text: s.pop("completion_text", None) figs = run_figures(streams) flagged, total, flags = gates.finish_flags(streams) run = {"n": n, "warmup": warmup, "utc": utc, "wall_s": wall, "streams": streams, "trial_void": any(s.get("status") != "ok" for s in streams), "finish_reason_flags": flagged, "finish_flags": flags, "completion_tokens_total": total, "hidden": {"max_surplus": max([s.get("hidden_surplus") or 0 for s in streams] or [0]), "method": sorted({s.get("hidden_method") for s in streams if s.get("hidden_method")})}, "cpu": cpu_delta(cpu0, cpu1, wall), "stop": None} run.update(figs) if power: run["power"] = power.reduce() run["power_samples"] = list(power.csv_rows(0.0)) run["ups"] = ups.reduce() run["stop"] = power.temp_tripped or ups.tripped return run #: The figures `medians` reduces; a VOID or refused cell holds each as None in `medians` (unscore_cell). MEDIAN_FIELDS = ("decode_tok_s_p50", "ttft_ms_p50", "ttft_ms_p95", "aggregate_sum_tok_s", "aggregate_wall_tok_s", "decode_tok_s_server_p50", "completion_tokens_total") def unscore_cell(cell, why): """A VOID (or refused) cell carries no scored figure anywhere a reader looks for one: its medians move to `unscored_medians` (a reading, kept) and `medians` holds n_scored 0, every figure None, and the why.""" if cell.get("unscored_medians") is None: cell["unscored_medians"] = cell.get("medians") kept = cell.get("unscored_medians") or {} m = {"n_scored": 0, "n_void": kept.get("n_void"), "thin": True, "rule": L.RULES["level_void"], "unscored_why": why, "spread": None, "spread_rule": L.RULES["spread"]} for f in MEDIAN_FIELDS: m[f] = {"median": None, "values": []} cell["medians"] = m cell["unscored_why"] = why return cell def sent_sha(text): """The first 16 hex of the sha256 of what a request sent (the string replay_hidden finds it by).""" return hashlib.sha256((text or "").encode("utf-8")).hexdigest()[:16] def replay_hidden(engine, stream, text, max_tokens=None): """The fold's fairness I-1: ONE replay of a retokenize-method surplus -- the same string at level 1 with logprobs (engine.raw_body(..., logprobs=True)), after the level's runs and outside every timed window -- and the exact count on the replay (engine.replay_exact_surplus: the generated tokens that sent no text). Its verdict: confirmed (the replay's exact count >= the surplus, on the same text), partly (0 < exact < surplus), not-confirmed (0 on the same text: re-tokenizing the visible text over-counted), unread (the replay failed, or its text differs: greedy decoding at level 1 need not repeat a level-16 output). The VOID stands whatever it reads.""" rec = {"verdict": "unread", "retokenize": stream.get("hidden_surplus"), "exact": None, "same_text": None, "rule": L.RULES["hidden_confirm"]} if text is None: return dict(rec, why="the sent string was not kept") r = E.new_record(engine.key, engine.posture["key"], "raw", stream.get("nonce"), stream.get("nonce_text")) mt = prereg.MAX_TOKENS if max_tokens is None else int(max_tokens) out, _wall, _rel = fire(engine, [(r, engine.raw_body(text, max_tokens=mt, logprobs=True))], "raw") st = classify(out[0], engine, stream.get("expected_prompt_tokens")) if st["status"] != "ok": return dict(rec, why=f"the replay failed: {st.get('fail_kind')}") exact = engine.replay_exact_surplus(st) same = st.get("completion_sha256") == stream.get("completion_sha256") rec.update(exact=exact, same_text=same, replay_completion_sha256=st.get("completion_sha256")) if exact is None: return dict(rec, why="the replay carried no logprobs") if not same: return dict(rec, why="the replay's text differs from the stream's (greedy decoding at level 1)") retok = int(stream.get("hidden_surplus") or 0) verdict = "confirmed" if exact >= retok else ("not-confirmed" if exact == 0 else "partly") return dict(rec, verdict=verdict, why=None) def medians(runs): """The level's medians over its scored, non-void runs, its spread, and what was left out.""" scored = [r for r in runs if not r["warmup"] and not r["trial_void"]] out = {"n_scored": len(scored), "n_void": sum(1 for r in runs if not r["warmup"] and r["trial_void"]), "thin": len(scored) < L.THIN_BELOW_RUNS, "rule": L.RULES["median"]} for f in MEDIAN_FIELDS: vals = [r.get(f) for r in scored if r.get(f) is not None] out[f] = {"median": L.rnd(L.median(vals)), "values": vals} out["spread"] = spread([r.get("aggregate_wall_tok_s") for r in scored]) out["spread_rule"] = L.RULES["spread"] return out def spread(values): """(max - min) / median of aggregate_wall over the scored runs (PLAN-v2 §2.6); None with fewer than two.""" v = [x for x in values if x is not None] if len(v) < 2: return None m = L.median(v) return None if not m else round((max(v) - min(v)) / m, 4) def run_strings(arm, prompt_id, prompt_text, lvl, n, path_kind="raw", road=None): """[(nonce key, nonce text, what is sent)] for run `n` of a level: on the raw path the composed string (framing.compose), on a chat row the user text (nonce + prompt, the engine renders it); a road marked `no_nonce` sends the prompt bare (arm P, PR-2). The same (arm, prompt, level, run, request) gives the same string on every engine.""" out = [] for r in range(int(lvl)): if road and road.get("no_nonce"): out.append((None, None, framing.compose("", prompt_text) if path_kind == "raw" else prompt_text)) continue ntext = N.nonce(arm, prompt_id, lvl, n, r) text = framing.compose(ntext, prompt_text) if path_kind == "raw" else framing.user_text(ntext, prompt_text) out.append((N.key(arm, prompt_id, lvl, n, r), ntext, text)) return out def scored_runs_for(level): """PLAN-v2 §2.6: 10 scored runs at levels 1 and 2, 3 above (prereg.SCORED_RUNS).""" return prereg.SCORED_RUNS[int(level)] def level(engine, arm, prompt_id, prompt_text, lvl, expected, nonce_end, *, scored=None, warmups=prereg.WARMUPS, samplers=None, cpu_pids=None, road=None, path_kind="raw", headline_eligible=None, printer=print, stop_check=None, url_path=None, log_offset_fn=None, strings_for=None, pre_run=None, max_tokens=None, keep_text=False): """One (engine, prompt, level) cell: its runs (every one printed), medians, spread, voids and flags. `expected` is {string: parity count}; `scored` defaults to the registered count for the level. `stop_check()` returns a stop line to end the level before its next run. PR-2 (SPEC Δ7), each optional and composable: `strings_for(n)` gives run n's [(key, nonce text, string)] when the requests of one run differ (arm P's new-question shape, arm Q's six-way waves; default: `run_strings`); `pre_run(n)` runs BEFORE run n, outside its timed window, and returns a stream record (arm P's warm request): a pre-run that failed voids the level (`warm_failed`, PLAN-v2 §2.8: any request fails); `max_tokens` and `keep_text` pass to run_once. Every run also records `log_span`, the engine log's byte offsets it spans (the pre-run included): arm P's slot-log excerpts read them.""" scored = scored_runs_for(lvl) if scored is None else scored cell = {"engine": engine.key, "posture": engine.posture["key"], "prompt": prompt_id, "level": int(lvl), "path": path_kind, "arm": arm, "road": (road or {}).get("name"), "runs": [], "voids": [], "flags": [], "headline_eligible": headline_eligible, "not_headline_why": None, "halted": None} if path_kind != "raw": cell["not_headline_why"] = "a chat row: as a user calls it, never a headline field (PLAN-v2 §2.4)" elif headline_eligible is not None and headline_eligible is not True: cell["not_headline_why"] = headline_eligible if isinstance(headline_eligible, str) else "not headline eligible" offset0 = engine.log_offset() first_run_offset = None sent = {} for n in range(warmups + scored): if stop_check: line = stop_check() if line: cell["halted"] = line break if first_run_offset is None: first_run_offset = engine.log_offset() span0 = engine.log_offset() pre = pre_run(n) if pre_run else None strings = strings_for(n) if strings_for else run_strings(arm, prompt_id, prompt_text, lvl, n, path_kind, road) sent.update((sent_sha(x[2]), x[2]) for x in strings) run = run_once(engine, cell, n, n < warmups, strings, expected, nonce_end, samplers, cpu_pids, road, url_path, path_kind, max_tokens=max_tokens, keep_text=keep_text) run["log_span"] = [span0, engine.log_offset()] if pre is not None: run["pre"] = pre if pre.get("status") != "ok": run["trial_void"] = True cell["runs"].append(run) printer(print_run(cell, run)) if run.get("stop"): cell["halted"] = f"{run['stop'].get('kind')}: {run['stop'].get('verdict')}" break start = first_run_offset if first_run_offset is not None else offset0 before = engine.runner_log_since(0, start) # the log up to the level's first run (byte offsets) span = engine.runner_log_since(start) # the level's own span: a second load line here is a reload layers = gates.all_layers(before, span, engine_key=engine.key) cell["all_layers"] = layers # after the level's own span is read: replay every retokenize surplus once (fairness I-1; the VOID stands) if path_kind == "raw" and getattr(engine, "replays_hidden", False): for r in cell["runs"]: for s in r["streams"]: flagged = s.get("hidden_method") == "retokenize" and (s.get("hidden_surplus") or 0) > 0 if flagged and s.get("path") == "raw": s["hidden_confirm"] = replay_hidden(engine, s, sent.get(s.get("sent_sha256")), max_tokens) cell["voids"] = gates.void_rollup(cell["runs"], layers, engine.key, prompt_id, lvl) pre_failed = [r["n"] for r in cell["runs"] if r.get("pre") is not None and r["pre"].get("status") != "ok"] if pre_failed: cell["voids"].append({"fail_kind": "warm_failed", "runs": pre_failed, "sentence": f"VOID: {engine.key} {prompt_id} L{lvl}: the warm request before run(s) " f"{', '.join(map(str, pre_failed))} failed (PLAN-v2 §2.8: any failed request " f"voids the level)"}) cell["medians"] = medians(cell["runs"]) cell["medians"]["completion_tokens_total_beside"] = True if cell["not_headline_why"]: cell["medians"]["not_headline"] = cell["not_headline_why"] cell["voided"] = bool(cell["voids"]) if cell["voided"]: unscore_cell(cell, "; ".join(v["sentence"] for v in cell["voids"])) for r in cell["runs"]: for s in r["streams"]: cell["flags"].extend(dict(f, run=r["n"], stream=s.get("stream")) for f in s.get("flags") or []) return cell def print_run(cell, run): tag = "warm-up" if run["warmup"] else f"run {run['n']}" return (f" {cell['engine']} {cell['posture']} {cell['prompt']} L{cell['level']} {tag}: ok " f"{run['streams_ok']}/{run['streams_ok'] + run['streams_failed']} decode p50 {run['decode_tok_s_p50']} tok/s " f"ttft p50 {run['ttft_ms_p50']} ms wall {run['aggregate_wall_tok_s']} tok/s completion tokens " f"{run['completion_tokens_total']}" + (" · VOID" if run["trial_void"] else "") + (f" · hidden {run['hidden']['max_surplus']}" if run["hidden"]["max_surplus"] else "") + (f" · {run['finish_reason_flags']} not-length" if run["finish_reason_flags"] else "")) def anchor(first_rows, anchor_rows): """The A-B-A drift (N-11): (anchor - first) / first of the median aggregate_wall at each level; printed, never used to correct anything. rows = {level: median aggregate_wall}.""" out = {} for lv in sorted(set(first_rows) | set(anchor_rows)): f, a = first_rows.get(lv), anchor_rows.get(lv) out[f"L{lv}"] = {"first": f, "anchor": a, "drift": round((a - f) / f, 4) if f and a is not None else None} return {"by_level": out, "rule": L.RULES["aba_drift"], "corrects": "nothing"}