arborist/bench/watt_bench.py
russell@unturf.com 5b1cbeed80
fix(#000057): measure power STATES, not a duty-cycle blend; guarantee cache miss
fox 2026-05-21: 'gen 200W' was a bug — joules/window blends the ~400W
generation bursts with the sub-100W gaps (retrieval/verify/network) into
a power state the card never sits at. A card occupies DISTINCT states
(idle / middle-idle = resident-between-requests / generation), differing
per card×model×server.

watt_probe.classify_power_bands(): largest-gap split of the window
samples into a low band (serving floor) and high band (generation draw)
+ duty cycle. Data-derived, never hardcoded — tested at two scales. The
worker emits the decomposition + raw samples; RemoteProbe/LocalProbe
expose band_stats() uniformly.

energy_cogs: marginal now taken against the measured SERVING FLOOR (the
standing cost of being ready), not deep idle; the blend is kept but
labelled window_mean_w. Reports idle/serving-floor/gen-draw/duty.

Cache-miss certainty (fox's question): the arborist arm runs
burn_existing=True (force-deletes any live providence row before
inference) and asserts cache_hits==0 with a loud warning + real_inference
flag — so we time real generation, never a SQLite lookup. Solo has no
cache path. 11 tests (energy math + band split). Validated live on the
isolated 4090: solo gen 308W/70%-duty vs substrate 396W/8.6%-duty —
substrate marginal/tok is LOWER, gross/tok higher (it holds the card
longer for retrieval).
2026-05-21 10:58:56 -04:00

701 lines
32 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""GPU wattage harness for the #000057 benchmark matrix — energy as a
constraint-optimizer axis (fox 2026-05-20).
The quality matrix (CG%) is scored by `bench/control_sweep.py`. This
harness adds the **cost** side: how much GPU energy does config X cost
to answer a question, on a 3090 vs a 4090? It samples
`nvidia-smi power.draw` on the inference GPU while driving a small
representative subset, and reports mean/peak watts, total joules,
joules-per-question, and joules-per-completion-token.
**Run it ON the GPU box** (the machine whose GPU serves the model) so
`nvidia-smi` reads the GPU actually doing the inference. The orchestrator
box has no GPU; the 3090/4090 live on the inference boxes. The driven
endpoint can be localhost (same box) or the public URL — the power
draw is the same GPU either way.
It does NOT grade answers (energy cost is independent of correctness).
Answers + per-question timing are saved to JSONL so a later
`bench/score_with_code_judge.py` pass can compute quality-per-joule.
Usage (on the GPU box)::
python -m bench.watt_bench --models qwen-nothink --n 20 \
--gpu-label 4090 --endpoint http://localhost:8080/v1
# arborist arm (retrieval) energy:
python -m bench.watt_bench --arborist-ref qwen-nothink --n 20 \
--gpu-label 3090
Output: `bench/qa_results/watt_<gpu>_<ts>.{jsonl,json}` — per-question
rows + an aggregate energy report. Compare two GPUs by running once on
each with the same `--n` and fixture slice.
"""
from __future__ import annotations
import argparse
import json
import subprocess
import sys
import threading
import time
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
from bench.control_ab import _SOLO_SYS, _gold # noqa: E402
from bench.control_sweep import MODELS, VARIANTS # noqa: E402
# ----------------------------------------------------------- power sampler
class PowerSampler:
"""Background thread polling `nvidia-smi --query-gpu=power.draw`.
Records (timestamp, watts) samples. Integrates to joules via the
trapezoid rule over the sampled window. Degrades gracefully when
nvidia-smi is absent (``available=False``) so the harness runs —
and is testable — anywhere; only the energy fields go null."""
def __init__(self, gpu_index: int = 0, interval_s: float = 0.5):
self.gpu_index = gpu_index
self.interval_s = interval_s
self.samples: list[tuple[float, float]] = [] # (t, watts)
self._stop = threading.Event()
self._thread: threading.Thread | None = None
self.available = self._probe()
self.gpu_name = self._gpu_name() if self.available else "no-nvidia-smi"
def _probe(self) -> bool:
try:
subprocess.run(["nvidia-smi", "--version"],
capture_output=True, timeout=5)
return True
except Exception: # noqa: BLE001
return False
def _gpu_name(self) -> str:
try:
out = subprocess.run(
["nvidia-smi", f"--id={self.gpu_index}",
"--query-gpu=name", "--format=csv,noheader"],
capture_output=True, text=True, timeout=5).stdout.strip()
return out or "unknown"
except Exception: # noqa: BLE001
return "unknown"
def _read_watts(self) -> float | None:
try:
out = subprocess.run(
["nvidia-smi", f"--id={self.gpu_index}",
"--query-gpu=power.draw", "--format=csv,noheader,nounits"],
capture_output=True, text=True, timeout=5).stdout.strip()
return float(out.splitlines()[0])
except Exception: # noqa: BLE001
return None
def _loop(self):
while not self._stop.is_set():
w = self._read_watts()
if w is not None:
self.samples.append((time.time(), w))
self._stop.wait(self.interval_s)
def start(self):
if not self.available:
return
self._thread = threading.Thread(target=self._loop, daemon=True)
self._thread.start()
def stop(self):
self._stop.set()
if self._thread:
self._thread.join(timeout=5)
def stats(self) -> dict:
"""Mean/peak watts and integrated joules over the sampled window."""
if not self.samples:
return {"available": False, "n_samples": 0,
"mean_w": None, "peak_w": None, "joules": None,
"window_s": None}
ws = [w for _, w in self.samples]
mean_w = sum(ws) / len(ws)
peak_w = max(ws)
# Trapezoid integration of power over time → joules.
joules = 0.0
for (t0, w0), (t1, w1) in zip(self.samples, self.samples[1:]):
joules += (w0 + w1) / 2.0 * (t1 - t0)
window_s = self.samples[-1][0] - self.samples[0][0]
return {"available": True, "n_samples": len(self.samples),
"mean_w": round(mean_w, 1), "peak_w": round(peak_w, 1),
"joules": round(joules, 1), "window_s": round(window_s, 1)}
# ----------------------------------------------------------- cpu power (RAPL)
class CpuSampler:
"""CPU package energy via Intel RAPL (`/sys/class/powercap/intel-rapl`).
RAPL exposes a CUMULATIVE energy counter in microjoules
(``energy_uj``); the energy used over a window is just the
end-minus-start diff (handling the wrap at ``max_energy_range_uj``),
so we don't integrate — RAPL gives joules directly, more accurate
than sampling instantaneous power. Mean watts = joules / seconds.
``energy_uj`` is root-only by default on most kernels (the
PLATYPUS side-channel mitigation, CVE-2020-8694). When it isn't
readable we degrade to ``available=False`` and report null CPU
energy — the GPU path is unaffected. Pass ``--cpu-energy-cmd`` to
supply a privileged reader (e.g. a sudo rule) when one exists.
Sums all RAPL packages (multi-socket) under intel-rapl:N."""
def __init__(self, energy_cmd: str | None = None):
self.energy_cmd = energy_cmd # optional privileged reader template
self._pkgs = self._discover_packages()
self._t0: float | None = None
self._e0: int | None = None
self.available = self._probe()
def _discover_packages(self) -> list[Path]:
base = Path("/sys/class/powercap")
if not base.exists():
return []
return sorted(p / "energy_uj" for p in base.glob("intel-rapl:*")
if (p / "energy_uj").exists() and ":" in p.name
and p.name.count(":") == 1) # top-level packages only
def _read_one(self, path: Path) -> int | None:
if self.energy_cmd:
try:
out = subprocess.run(
self.energy_cmd.format(path=str(path)), shell=True,
capture_output=True, text=True, timeout=5).stdout.strip()
return int(out)
except Exception: # noqa: BLE001
return None
try:
return int(path.read_text().strip())
except Exception: # noqa: BLE001
return None
def _read_total(self) -> int | None:
if not self._pkgs:
return None
vals = [self._read_one(p) for p in self._pkgs]
if any(v is None for v in vals):
return None
return sum(vals)
def _probe(self) -> bool:
return self._read_total() is not None
def start(self):
if not self.available:
return
self._e0 = self._read_total()
self._t0 = time.time()
def stats(self) -> dict:
if not self.available or self._e0 is None:
return {"available": False, "joules": None, "mean_w": None,
"window_s": None, "n_packages": len(self._pkgs)}
e1 = self._read_total()
t1 = time.time()
if e1 is None:
return {"available": False, "joules": None, "mean_w": None,
"window_s": None, "n_packages": len(self._pkgs)}
# energy_uj wraps; if it went down assume a single wrap per pkg.
# max_energy_range_uj is per-package; approximate with the
# observed positive delta when no wrap, else mark suspect.
d_uj = e1 - self._e0
if d_uj < 0:
# wrap occurred — best-effort: skip rather than report garbage
return {"available": True, "joules": None, "mean_w": None,
"window_s": round(t1 - self._t0, 1),
"n_packages": len(self._pkgs), "note": "counter wrapped"}
joules = d_uj / 1_000_000.0
window = t1 - self._t0
return {"available": True, "joules": round(joules, 1),
"mean_w": round(joules / window, 1) if window else None,
"window_s": round(window, 1), "n_packages": len(self._pkgs)}
# ----------------------------------------------------------- probe abstraction
class LocalProbe:
"""GPU+CPU power sampled on THIS machine (watt_bench runs on the
GPU box). Wraps PowerSampler + CpuSampler behind the start/stop/
gpu_stats/cpu_stats interface RemoteProbe also implements."""
def __init__(self, gpu_index: int, interval: float, cpu_cmd: str | None):
self._g = PowerSampler(gpu_index, interval)
self._c = CpuSampler(cpu_cmd)
self.gpu_name = self._g.gpu_name
self.available = self._g.available
self.cpu_available = self._c.available
def start(self):
self._g.samples = []
self._g.start()
self._c.start()
def stop(self):
self._g.stop()
self._gs = self._g.stats()
self._cs = self._c.stats()
def gpu_stats(self) -> dict:
return self._gs
def cpu_stats(self) -> dict:
return self._cs
def band_stats(self) -> dict:
from bench.watt_probe import classify_power_bands
return classify_power_bands([w for _, w in self._g.samples])
class RemoteProbe:
"""GPU+CPU power sampled on a REMOTE worker box over SSH — the
laptop-driver / worker-reporter mode (fox 2026-05-20). The laptop
runs the workload (retrieval + judge + driving the worker's LLM
endpoint); this orchestrates ``bench/watt_probe.py`` on the worker
for the workload window so the worker reports its own power without
needing shards / arborist / a venv (the probe is stdlib-only).
start() scp's the probe (once) and launches it in --until-file mode
detached; stop() touches the stop-file, waits for flush, fetches the
JSON, and exposes it in the LocalProbe stats shape."""
def __init__(self, host: str, gpu_index: int, gpu_label: str,
interval: float):
self.host = host
self.gpu_index = gpu_index
self.gpu_label = gpu_label or "gpu"
self.interval = interval
tag = time.strftime("%H%M%S", time.gmtime())
self.stop_file = f"/tmp/wattprobe_{tag}.stop"
self.out_file = f"/tmp/wattprobe_{tag}.json"
self._report: dict | None = None
self.gpu_name = f"remote:{host}"
self.available = self._ensure_probe()
self.cpu_available = self.available # confirmed after stop()
def _ssh(self, cmd: str, timeout: int = 30):
return subprocess.run(
["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=10",
self.host, cmd],
capture_output=True, text=True, timeout=timeout)
def _ensure_probe(self) -> bool:
probe = str(Path(__file__).resolve().parent / "watt_probe.py")
try:
subprocess.run(
["scp", "-q", "-o", "BatchMode=yes", probe,
f"{self.host}:/tmp/watt_probe.py"],
check=True, timeout=30)
return True
except Exception: # noqa: BLE001
return False
def start(self):
self._ssh(f"rm -f {self.stop_file} {self.out_file}")
# Detach so ssh returns while the probe keeps sampling.
self._ssh(
f"nohup python3 /tmp/watt_probe.py --until-file {self.stop_file} "
f"--out {self.out_file} --gpu-index {self.gpu_index} "
f"--gpu-label {self.gpu_label} --interval {self.interval} "
f">/dev/null 2>&1 & echo started")
def stop(self):
self._ssh(f"touch {self.stop_file}")
time.sleep(self.interval * 3 + 1.0) # let the probe flush its JSON
r = self._ssh(f"cat {self.out_file}")
try:
self._report = json.loads(r.stdout)
self.gpu_name = self._report.get("gpu_name", self.gpu_name)
self.cpu_available = bool(self._report.get("cpu_rapl_available"))
except Exception: # noqa: BLE001
self._report = None
def gpu_stats(self) -> dict:
rp = self._report or {}
return {"available": rp.get("gpu_available", False),
"mean_w": rp.get("gpu_mean_w"), "peak_w": rp.get("gpu_peak_w"),
"joules": rp.get("gpu_joules"), "window_s": rp.get("window_s"),
"util_mean_pct": rp.get("gpu_util_mean_pct"),
"n_samples": rp.get("n_samples")}
def cpu_stats(self) -> dict:
rp = self._report or {}
return {"available": rp.get("cpu_rapl_available", False),
"mean_w": rp.get("cpu_mean_w"), "joules": rp.get("cpu_joules"),
"window_s": rp.get("window_s")}
def band_stats(self) -> dict:
return (self._report or {}).get("power_bands", {"available": False})
def watt_samples(self) -> list:
return (self._report or {}).get("watt_samples", [])
# ----------------------------------------------------------- driver
def _make_client(endpoint: str):
from arborist.qa.client import OpenAICompatibleClient
import os
return OpenAICompatibleClient(
base_url=endpoint,
api_key=os.environ.get("ARBORIST_LLM_API_KEY"))
def _solo_call(client, mkey: str, question: str) -> tuple[str, int]:
"""One solo completion. Returns (answer, completion_tokens). The
bench client returns only the string; for token counts we re-issue
via a thin httpx call when usage is wanted — but to keep the energy
path simple we estimate tokens as len//4 when usage isn't surfaced."""
cfg = MODELS[mkey]
ans = client.chat_completion(
[{"role": "system", "content": _SOLO_SYS},
{"role": "user", "content": question}],
model=cfg["model"], max_tokens=cfg["max_tokens"],
extra_body=cfg["extra"])
return ans, max(1, len(ans) // 4)
def energy_cogs(gpu_joules: float | None, window_s: float | None,
tot_tok: int, idle_mean_w: float | None,
price_per_kwh: float,
serving_floor_w: float | None = None,
gen_draw_w: float | None = None,
gen_duty: float | None = None) -> dict:
"""Energy cost-of-goods-sold from MEASURED power states — no hardcodes.
fox 2026-05-21: a card occupies DISTINCT power states (idle, middle-
idle = model resident between requests, generation), and they differ
for every card × model × inference-server. So every power number is
measured at runtime and the states are kept DISTINCT — we do NOT
collapse them into ``joules/window`` and call that "generation"
(that blend is a duty-cycle artifact, not a state the card sits at).
* idle_mean_w — deep idle (no-request baseline window).
* serving_floor_w — middle-idle: model resident, between requests
(low band of the driven window). The right floor
for marginal cost — the standing cost of being
ready to answer.
* gen_draw_w — the actual generation draw (high band), with
``gen_duty`` = fraction of the window generating.
* gross COGS — ALL measured joules / tokens: the all-in cost,
amortized across throughput (state-agnostic, so
this number was always correct).
* marginal COGS — joules ABOVE the serving floor / tokens: what
one more request's generation actually adds.
Falls back to idle floor if no band split.
kWh = J / 3.6e6; ``$/1k-tok`` is the unit that compares to API pricing.
Only ``price_per_kwh`` is an operator input. See docs/stock-v1-config.md.
"""
if not gpu_joules or not tot_tok or not window_s:
return {"available": False, "price_per_kwh": price_per_kwh}
J_PER_KWH = 3.6e6
gross_usd = gpu_joules / J_PER_KWH * price_per_kwh
out = {
"available": True,
"price_per_kwh": price_per_kwh,
# measured power states (distinct — not a blend)
"idle_w_measured": (round(idle_mean_w, 1)
if idle_mean_w is not None else None),
"serving_floor_w_measured": (round(serving_floor_w, 1)
if serving_floor_w is not None else None),
"gen_draw_w_measured": (round(gen_draw_w, 1)
if gen_draw_w is not None else None),
"gen_duty_cycle": gen_duty,
"window_mean_w": round(gpu_joules / window_s, 1), # blend, labelled
"window_s": round(window_s, 1),
"gross_joules": round(gpu_joules, 1),
"gross_usd": round(gross_usd, 6),
"gross_usd_per_1k_tok": round(gross_usd / tot_tok * 1000, 6),
}
# Marginal against the serving floor (preferred) or idle (fallback).
floor_w = serving_floor_w if serving_floor_w is not None else idle_mean_w
if floor_w is not None:
floor_joules = floor_w * window_s
marginal_joules = max(0.0, gpu_joules - floor_joules)
marginal_usd = marginal_joules / J_PER_KWH * price_per_kwh
out.update({
"marginal_floor": ("serving" if serving_floor_w is not None
else "idle"),
"marginal_floor_w": round(floor_w, 1),
"marginal_joules": round(marginal_joules, 1),
"marginal_usd": round(marginal_usd, 6),
"marginal_usd_per_1k_tok": round(marginal_usd / tot_tok * 1000, 6),
})
return out
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--fixture", default="bench/qa_questions_stale_map.json")
ap.add_argument("--n", type=int, default=20,
help="subset size — small & representative")
ap.add_argument("--models", default="qwen-nothink",
help="comma list for the solo arm (energy of the "
"model alone)")
ap.add_argument("--arborist-ref", default="",
help="if set, also measure the arborist (retrieval) "
"arm energy with this model as synthesis LLM")
ap.add_argument("--variants", default="plain")
ap.add_argument("--shards-dir",
default=str(Path.home() / ".arborist" / "shards"))
ap.add_argument("--gpu-index", type=int, default=0)
ap.add_argument("--gpu-label", default="",
help="human label for the GPU tier (e.g. 3090, 4090) "
"— stamped into the output for cross-GPU compare")
ap.add_argument("--remote-gpu-host", default="",
help="laptop-driver mode: run the workload HERE "
"(local shards + judge, driving --endpoint) but "
"sample power on this REMOTE worker box over SSH "
"via bench/watt_probe.py. The worker needs no "
"shards / arborist / venv. e.g. ai.foxhop.net")
ap.add_argument("--endpoint", default="",
help="override the model endpoint. In --remote-gpu-host "
"mode point this at the worker's endpoint (its "
"public URL or an SSH tunnel). Default = MODELS "
"entry's endpoint.")
ap.add_argument("--sample-interval", type=float, default=0.5)
ap.add_argument("--idle-baseline-s", type=float, default=5.0,
help="seconds to sample idle (warm-idle) power before "
"driving load — MEASURED, never hardcoded")
ap.add_argument("--price-per-kwh", type=float, default=0.33,
help="electricity rate for energy COGS (USD/kWh; "
"configurable site rate, fox default 0.33). The "
"only non-measured input — power states are all "
"measured per card/model/server at runtime.")
ap.add_argument("--answer-mode", choices=["quote", "claim_lattice"],
default="claim_lattice",
help="STOCK V.1 substrate-ON answer shape for the "
"arborist arm (frozen bench.stock_v1 policy).")
ap.add_argument("--cpu-energy-cmd", default="",
help="optional privileged reader for RAPL energy_uj "
"when it's root-only (PLATYPUS mitigation), e.g. "
"'sudo cat {path}'. {path} is substituted with "
"the intel-rapl energy_uj sysfs path. Omit to "
"read directly (works only if energy_uj is "
"world-readable).")
ap.add_argument("--out-dir", default="bench/qa_results")
a = ap.parse_args()
models = [m for m in a.models.split(",") if m in MODELS]
variants = [v for v in a.variants.split(",") if v in VARIANTS]
items = json.loads(Path(a.fixture).read_text())[:a.n]
shards_dir = Path(a.shards_dir)
ts = time.strftime("%Y-%m-%dT%H-%M-%SZ", time.gmtime())
label = a.gpu_label or "gpu"
outp = Path(a.out_dir) / f"watt_{label}_{ts}.jsonl"
outp.parent.mkdir(parents=True, exist_ok=True)
remote = a.remote_gpu_host.strip()
def _new_probe():
if remote:
return RemoteProbe(remote, a.gpu_index, label, a.sample_interval)
return LocalProbe(a.gpu_index, a.sample_interval,
a.cpu_energy_cmd or None)
probe0 = _new_probe()
mode = f"REMOTE via ssh {remote}" if remote else "LOCAL"
print(f"watt_bench [{mode}] — gpu_label={label!r} "
f"detected={probe0.gpu_name!r} "
f"gpu={'yes' if probe0.available else 'NO'}")
print(f" cpu energy: {'yes' if probe0.cpu_available else 'NO (RAPL root-only / unavailable)'}")
print(f" fixture={a.fixture} n={len(items)} models={models} "
f"variants={variants} arborist_ref={a.arborist_ref or 'none'}")
# Idle baseline.
idle = None
if probe0.available and a.idle_baseline_s > 0:
print(f" sampling idle baseline for {a.idle_baseline_s}s "
f"({'remote' if remote else 'local'}) …")
probe0.start()
time.sleep(a.idle_baseline_s)
probe0.stop()
ig = probe0.gpu_stats()
idle = {"mean_w": ig.get("mean_w"), "peak_w": ig.get("peak_w"),
"gpu_util_mean_pct": ig.get("util_mean_pct"),
"cpu_mean_w": probe0.cpu_stats().get("mean_w")}
print(f" idle: gpu mean={idle['mean_w']}W peak={idle['peak_w']}W "
f"util={idle['gpu_util_mean_pct']}% cpu={idle['cpu_mean_w']}W"
+ (" [util>5% ⇒ contaminated by live traffic]"
if (idle.get('gpu_util_mean_pct') or 0) > 5 else ""))
cells = [] # one entry per (arm, model, variant) measured
def run_cell(arm: str, mkey: str):
endpoint = a.endpoint or MODELS[mkey]["endpoint"]
client = _make_client(endpoint)
rows = []
cache_hits = 0 # arborist arm: must stay 0 (real inference, not lookup)
probe = _new_probe()
probe.start()
t0 = time.time()
try:
for variant in variants:
for it in items:
q = VARIANTS[variant](it["question"])
tq = time.time()
cache_status = "solo_no_cache"
if arm == "solo":
ans, ctoks = _solo_call(client, mkey, q)
else: # arborist
from arborist.qa.query import query
from bench.stock_v1 import (
assert_not_drifted as _assert_stock)
from bench.stock_v1 import policy_for as _stock_policy_for
gold = _gold(shards_dir, it.get("shard", ""),
it["target_root"])
# STOCK V.1 frozen substrate-ON policy (cost is
# measured for the SAME substrate the campaign
# grades). Reasoning refs (phase 3) layer the
# documented JSON overrides and skip the assert.
reasoning = bool(MODELS[mkey].get("reasoning"))
if not reasoning:
_assert_stock(a.answer_mode)
pol = _stock_policy_for(a.answer_mode)
if reasoning:
pol["claim_lattice_json_stop_sequences"] = []
pol["max_tokens"] = 8192
qa_db = Path("/tmp") / f"watt_{ts}_{mkey}.db"
# burn_existing force-deletes any matching live row
# before inference, so the arborist arm ALWAYS runs
# the LLM (real generation energy) and never times a
# cache lookup. cache_hits MUST stay 0 (fox 2026-05-21).
r = query(question=q, qa_db=qa_db,
chat_client=client,
model_id=MODELS[mkey]["model"],
shards_dir=shards_dir,
extra_body=MODELS[mkey]["extra"],
policy=pol, burn_existing=True)
ans = r.get("raw_answer") or r.get("answer_text") or ""
ctoks = max(1, len(ans) // 4)
cache_status = r.get("status", "?")
if (cache_status == "cache_hit"
or "cache_hit" in str(r.get("lookup_path", ""))):
cache_hits += 1
dt = time.time() - tq
rows.append({"arm": arm, "model": mkey,
"variant": variant,
"latency_s": round(dt, 2),
"answer_chars": len(ans),
"est_completion_tokens": ctoks,
"cache_status": cache_status})
finally:
client.close()
probe.stop()
st = probe.gpu_stats()
cpu_st = probe.cpu_stats()
bands = probe.band_stats()
elapsed = time.time() - t0
nq = len(rows)
tot_tok = sum(r["est_completion_tokens"] for r in rows)
gpu_j = st.get("joules")
cpu_j = cpu_st.get("joules")
total_j = (gpu_j or 0) + (cpu_j or 0) if (gpu_j or cpu_j) else None
cell = {
"arm": arm, "model": mkey, "variants": variants,
"n_questions": nq, "elapsed_s": round(elapsed, 1),
"gpu_power": st,
"cpu_power": cpu_st,
"gpu_label": label, "gpu_name": probe.gpu_name,
"idle_baseline": idle,
"gpu_joules_per_question": (round(gpu_j / nq, 1)
if gpu_j and nq else None),
"cpu_joules_per_question": (round(cpu_j / nq, 1)
if cpu_j and nq else None),
"total_joules_per_question": (round(total_j / nq, 1)
if total_j and nq else None),
"est_total_completion_tokens": tot_tok,
"gpu_joules_per_token": (round(gpu_j / tot_tok, 3)
if gpu_j and tot_tok else None),
"mean_latency_s": (round(sum(r["latency_s"] for r in rows) / nq, 2)
if nq else None),
# Measured power-state decomposition (idle / serving-floor /
# generation), data-derived, never hardcoded.
"power_bands": bands,
# Real-inference guard: arborist arm runs under burn_existing,
# so cache_hits MUST be 0 — else we'd be timing a SQLite lookup,
# not generation, and the energy number is meaningless.
"cache_hits": cache_hits,
"real_inference": (cache_hits == 0),
# Energy COGS — measured power STATES (not a duty-cycle blend);
# marginal taken against the serving floor (middle-idle), the
# standing cost of being ready. Only price_per_kwh is operator
# input. Window timestamps let a post-hoc load_monitor cross-ref
# flag organic-traffic contamination under non-isolation.
"energy_cogs": energy_cogs(
gpu_j, st.get("window_s"), tot_tok,
(idle or {}).get("mean_w"), a.price_per_kwh,
serving_floor_w=(bands.get("low_band_w")
if bands.get("bimodal") else None),
gen_draw_w=bands.get("high_band_w"),
gen_duty=bands.get("duty_cycle")),
"window_start_unix": round(t0, 3),
"window_end_unix": round(t0 + elapsed, 3),
"watt_samples": (probe.watt_samples()
if hasattr(probe, "watt_samples") else []),
}
cells.append(cell)
with open(outp, "a") as f:
for r in rows:
f.write(json.dumps({**r, "gpu_label": label}) + "\n")
print(f" [{arm}/{mkey}] n={nq} "
f"gpu={st.get('mean_w')}W/{st.get('peak_w')}peak "
f"cpu={cpu_st.get('mean_w')}W "
f"J/q gpu={cell['gpu_joules_per_question']} "
f"cpu={cell['cpu_joules_per_question']} "
f"tot={cell['total_joules_per_question']} "
f"lat={cell['mean_latency_s']}s")
if arm == "arborist" and cache_hits:
print(f" !! WARNING {cache_hits}/{nq} CACHE HITS — "
f"energy NOT a true generation cost (re-run; expected 0)")
b = bands or {}
if b.get("available"):
print(f" states: idle≈{(idle or {}).get('mean_w')}W · "
f"serving-floor {b.get('low_band_w')}W · "
f"gen {b.get('high_band_w')}W "
f"(duty {b.get('duty_cycle')}, peak {b.get('peak_w')}W)"
+ ("" if b.get("bimodal") else " [unimodal — no gen state]"))
cg = cell["energy_cogs"]
if cg.get("available"):
print(f" COGS @${cg['price_per_kwh']}/kWh: "
f"gross ${cg['gross_usd_per_1k_tok']}/1k-tok"
+ (f" · marginal ${cg['marginal_usd_per_1k_tok']}/1k-tok "
f"(vs {cg['marginal_floor']}-floor {cg['marginal_floor_w']}W)"
if "marginal_usd_per_1k_tok" in cg else
" · marginal n/a (no floor)"))
for mkey in models:
run_cell("solo", mkey)
if a.arborist_ref and a.arborist_ref in MODELS:
run_cell("arborist", a.arborist_ref)
report = {"ts": ts, "gpu_label": label, "gpu_name": probe0.gpu_name,
"mode": "remote" if remote else "local",
"remote_gpu_host": remote or None,
"gpu_available": probe0.available,
"cpu_energy_available": probe0.cpu_available,
"idle_baseline": idle, "fixture": a.fixture,
"n_items": len(items), "cells": cells}
rp = outp.with_suffix(".json").with_name(f"watt_{label}_{ts}.json")
rp.write_text(json.dumps(report, indent=2))
print(f"\n per-question rows: {outp}")
print(f" energy report: {rp}")
if not probe0.available:
print(" NOTE: GPU power unavailable — energy fields null. Run ON "
"the GPU box, or use --remote-gpu-host to sample a worker.")
return 0
if __name__ == "__main__":
raise SystemExit(main())