diff --git a/evals/postmortem/README.md b/evals/postmortem/README.md new file mode 100644 index 0000000000..742f28dff3 --- /dev/null +++ b/evals/postmortem/README.md @@ -0,0 +1,117 @@ +# Post-mortem harness: forensics + live A/B for the 1,393-agent run fixes + +The scripts that produced every number in tracking issue #103563 and the "Independent review +(round 2)" sections of its 13 PRs. Two halves: + +- **`forensics/`** reads a *copy* of a Hermes `state.db` (plus rotated `agent.log*` and git) and + recomputes the *observed* figures for any run: where the money went, per-call cache behaviour, + nested-delegate timeouts, batch-join delivery delay, tool friction, `/goal` loop behaviour, and the + post-open rework inventory. It needs no model calls and no network. +- **`live_ab/`** and **`review_probes/`** exercise the real code paths of a checkout (real + `AIAgent`, dispatch, judge, scanner, SDK) against local fake providers to show what each fix + *does*. `run.py` runs them against one or two checkouts and prints PASS/FAIL side by side. + `review_probes/` are the probes the independent `/review` wrote; each reproduced a defect in the + first version of a PR and the fixed head must pass it. + +Everything is labeled **OBSERVED** (from usage rows / logs / git) or **MODELED** (a reconstruction or +replay). Do not add the modeled figures to the observed ones; see §5 of #103563 for why they overlap. + +## Requirements + +- A Hermes checkout with its venv (`.venv/bin/python`). NeMo Relay is not required for the forensics; + it is what the run itself used for the wire captures in `live_ab/cache_prefix_wire.py`. +- For forensics: a **copy** of `~/.hermes/state.db` (never point at the live file; `sqlite3 state.db ".backup copy.db"` + or `cp` while Hermes is idle) and, optionally, the rotated `~/.hermes/logs/agent.log*`. +- For `--live` probes: real credentials in `HERMES_HOME` (they spend cents per run). + +## Forensics: recompute the observed numbers for YOUR run + +```bash +cd +P=.venv/bin/python +# 1. cost buckets, depth/duration shares, context-size reconstruction, excess-cache-write proxy, cap replay +$P -m evals.postmortem.forensics.tokens --db state_copy.db [--root ] [--cap 200000 --floor 65000] +# 2. per-call cache behaviour from the logs (coverage fraction is printed first; quote nothing without it) +$P -m evals.postmortem.forensics.logcalls --db state_copy.db --logs "$HOME/.hermes/logs/agent.log*" +# 3. delegation: timeouts, orphaned children, polling, batch-join delay, truncated summaries +$P -m evals.postmortem.forensics.delegation --db state_copy.db +# 4. tool friction: hardline false blocks, foreground refusals, whole-file rewrites, output volume +$P -m evals.postmortem.forensics.tools --db state_copy.db +# 5. /goal loop: nudges, parked barrier, notification counts +$P -m evals.postmortem.forensics.goal_loop --db state_copy.db +# 6. rework inventory for a large PR (git only) +$P -m evals.postmortem.forensics.rework --repo . --base --open --head +``` + +Each writes `postmortem_out/.json` and prints a summary. `--root` defaults to the top-level +session with the most descendants; compression-rollover children are excluded from the population so +cost buckets are disjoint. Pricing is fitted from `estimated_cost_usd`, so dollars match what that +Hermes recorded (an estimator, not an invoice). + +### Reference output (the #102117 run, `state_copy.db` of 2026-09-04) + +| lane | prints | +|---|---| +| tokens | `1394 sessions, $19,302.59`; buckets cache_write $11,159.76 · cache_read $3,587.17 · output $4,555.48; depth-2 65%; >60 min 61%; MODELED excess cache-write proxy ~$9.1k; sawtooth cap 200K: prompt volume ×0.76 (reconstructed sizes) | +| logcalls | coverage 22,489/93,284 (24.1%); median prompt 229,648, p90 408,794, >200K 59%; hit ratio 93.9%; strict plateau 22.5% of uncached input, non-advancing 74.8%; sawtooth cap 200K on real sizes ×0.498 | +| delegation | 266 orchestrators; 332 delegate_task timeouts in 234 sessions (all nested); 93 results carried; 242.6 h sleep after first timeout; those sessions' lifetime $4,034.69 (includes real work); summaries truncated 123/220; batch-join withheld child-hours root 233 / depth-1 300 / depth-2 52 | +| tools | 96,855 tool calls; hardline blocks 579 (566 "malformed" class); foreground timeout refusals 475 (303 asked 900 s); `&` 185, nohup 24; write_file 8,188 calls / 92.8M chars, 661 rewrites of a file read this session >20k; patch 4,623 | +| goal_loop | 5 nudges (2 within 180 s of a "waiting" turn); 34 batch notices; 48 bg-process notices; final barrier `waiting_on_session=proc_…` parked 201 min | +| rework | surface at open: 1,703 names / 341 modules, 951 methods / 156, 126 test defs / 52 files; 125 post-open commits (simplify 55, review-fix 39, fix 19, …) | + +The two sawtooth figures differ on purpose: `tokens` replays *reconstructed* per-call sizes for all +93k calls (×0.76); `logcalls` replays *real* per-call sizes for the 24% of calls in the logs (×0.50, +the peak-concurrency window). The tracking issue quotes the second and says so. + +## Live A/B: what each fix does + +```bash +# offline probes (fake providers, temp HERMES_HOME), main vs a branch or integration checkout: +.venv/bin/python -m evals.postmortem.run --repo /path/to/main --compare /path/to/branch +# add the probes that make real provider calls (cents): +.venv/bin/python -m evals.postmortem.run --repo /path/to/branch --live +# one PR only: +.venv/bin/python -m evals.postmortem.run --repo . --only 103492 +``` + +| probe | PR | expects on the fixed head | +|---|---|---| +| `live_ab/hardline_scanner_matrix.py` | #103492 | 11-case matrix `ALL OK` (546-block class allowed; newline/`;`/`&&`/`|` hidden `reboot` blocked as itself) | +| `review_probes/scanner_bypass_probe.py` | #103492 | public guard `approved: False`, 0 callbacks, harmless Bash witness not executed | +| `live_ab/subagent_context_cap.py` | #103513 | child trigger `200,000` on a 1M model; parent untouched | +| `review_probes/context_cap_probe.py` | #103513 | cap holds through repeated compression + persistence; config validation | +| `live_ab/nested_delegate_deadline.py` | #103486 | 40 s deadline + 75 s leaf: result delivered (main: lost) | +| `review_probes/deadline_probe.py` | #103486 | same through actual dispatch | +| `live_ab/auth_stampede.py 40` | #103526 | `401s=0` (main: 40) | +| `review_probes/credential_identity_probe.py pr` | #103526 | explicit account-A key stays A (v1: became B) | +| `live_ab/batch_failure_notice.py` | #103549 | `TASK_FAILURE_NOTICE` at t+0.3 s, `BATCH_FINAL` after | +| `review_probes/notice_delivery_probe.py` | #103549 | gateway receives notice, notice, final; busy-parent final claim succeeds | +| `review_probes/cache_estimator_probe.py` | #103476 | preflight ≈ wire estimate; `should_compress` agrees | +| `live_ab/cache_prefix_wire.py B` (live) | #103476 | 0 mutated prefixes across 6 calls | +| `live_ab/goal_judge_wait.py 3` (live) | #103534 | `wait` ×3 on the run's "waiting on workers" response (main: `continue` ×3) | +| `review_probes/goal_scope_probe.py` | #103496/#103534 | judge sees own processes; delegation WAIT lifts on batch return | +| `review_probes/goal_repaste_probe.py` | #103553 | near-whole re-paste → pointer; `ship the API` ≠ `ship the UI` | +| `review_probes/rewrite_hint_probe.py` | #103551 | remote backend: no host-derived hint; FIFO: returns; 460 KB repeated-line file: skipped, not 22 s | +| `review_probes/finalizer_schedule_probe.py` | #103507 | pytest plugin: `-p evals.postmortem.review_probes.finalizer_schedule_probe --finalizer-probe=consumer-first` on `tests/e2e/test_relay_native_openai_stream.py` → 2 passed | + +## What is NOT here, and why + +- **The trajectories themselves.** The run's `state.db` contains 51,956 absolute home paths, 5,341 + e-mail addresses, private IPs, chat/user ids, and real-shaped API keys and JWTs in tool output. It + is not publishable, and this harness is written so it does not need to be: run it on your own DB. +- **Hand classification.** The root-cause classes of the 72 rework commits (dropped symbol vs semantic + drift vs compat fallout) were labeled by hand in the original audit; `forensics/rework.py` reproduces + the mechanical inventory that labeling started from and stops there. +- **A single aggregate saving.** By design. Each lane prints its own number with its own caveat. + +## Evidence bundle + +The original lane reports, the independent review, and the JSON this harness recomputes on the run's DB +are in a secret gist linked from #103563 (no trajectories; see "What is NOT here"). + +## Provenance + +Forensic lanes: five parallel Hermes subagents (2026-09-04), rewritten here to take `--db`/`--root` +instead of hard-coded paths. `live_ab/`: the primary agent's per-PR A/Bs. `review_probes/`: the +independent `/review` subagent's probes (2026-09-05), adapted to take paths from the command line; +their findings and the fixes are in each PR's "Independent review (round 2)" section. diff --git a/evals/postmortem/__init__.py b/evals/postmortem/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/evals/postmortem/forensics/__init__.py b/evals/postmortem/forensics/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/evals/postmortem/forensics/common.py b/evals/postmortem/forensics/common.py new file mode 100644 index 0000000000..6d3d8f17a1 --- /dev/null +++ b/evals/postmortem/forensics/common.py @@ -0,0 +1,186 @@ +"""Shared loaders for the post-mortem forensics: point at ANY Hermes ``state.db`` (a copy, never the live +file) and get the run tree, the in-run session set, fitted pricing and message iterators. + +Nothing here knows about a particular run. The root is discovered as the session with the most +descendants unless ``--root`` is given; compression-rollover children (a child whose ``id`` the parent's +``compaction`` metadata names as its continuation) are excluded from the tree so cost populations stay +disjoint. Pricing is fitted by least squares from ``sessions`` usage columns to ``estimated_cost_usd``, so +the recomputed dollars match what THAT Hermes recorded, not an invoice. + +Usage from a lane script:: + + from evals.postmortem.forensics.common import Run + run = Run.from_args() # --db, --root, --out + for sid in run.in_run: # ordered session ids + ... + run.write("q1_cost.json", data) +""" +from __future__ import annotations + +import argparse +import collections +import json +import sqlite3 +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any, Dict, Iterable, Iterator, List, Optional + +USAGE_COLS = ("input_tokens", "cache_read_tokens", "cache_write_tokens", "output_tokens") + + +def _lstsq(rows: List[List[float]], y: List[float]) -> List[float]: + """Ordinary least squares without numpy (4 unknowns): normal equations solved by Gaussian elimination.""" + n = len(rows[0]) + ata = [[sum(r[i] * r[j] for r in rows) for j in range(n)] for i in range(n)] + atb = [sum(r[i] * yy for r, yy in zip(rows, y)) for i in range(n)] + m = [row[:] + [b] for row, b in zip(ata, atb)] + for c in range(n): + piv = max(range(c, n), key=lambda r: abs(m[r][c])) + m[c], m[piv] = m[piv], m[c] + if abs(m[c][c]) < 1e-12: + continue + for r in range(n): + if r != c: + f = m[r][c] / m[c][c] + m[r] = [a - f * b for a, b in zip(m[r], m[c])] + return [m[i][n] / m[i][i] if abs(m[i][i]) > 1e-12 else 0.0 for i in range(n)] + + +@dataclass +class Run: + db_path: Path + out_dir: Path + root: str + sessions: Dict[str, Dict[str, Any]] + depth: Dict[str, int] + in_run: List[str] # root + descendants, rollover excluded, dispatch order + price_per_token: Dict[str, float] # fitted: USD per token for each USAGE_COLS entry + _conn: sqlite3.Connection = field(repr=False) + + # ── construction ────────────────────────────────────────────────────────────────────────── + @classmethod + def parser(cls, description: str = "") -> argparse.ArgumentParser: + ap = argparse.ArgumentParser(description=description) + ap.add_argument("--db", required=True, help="path to a COPY of ~/.hermes/state.db") + ap.add_argument("--root", default=None, help="root session id (default: the session with the most descendants)") + ap.add_argument("--out", default="postmortem_out", help="directory for JSON/markdown outputs") + return ap + + @classmethod + def from_args(cls, argv: Optional[List[str]] = None, description: str = "") -> "Run": + a = cls.parser(description).parse_args(argv) + return cls.open(a.db, root=a.root, out=a.out) + + @classmethod + def open(cls, db: str, *, root: Optional[str] = None, out: str = "postmortem_out") -> "Run": # noqa: C901 + conn = sqlite3.connect(f"file:{db}?mode=ro", uri=True) + conn.row_factory = sqlite3.Row + sessions = {r["id"]: dict(r) for r in conn.execute("SELECT * FROM sessions")} + children: Dict[Optional[str], List[str]] = collections.defaultdict(list) + for sid, s in sessions.items(): + children[s.get("parent_session_id")].append(sid) + rollover = cls._rollover_ids(conn, sessions) + if root is None: + def size(sid: str) -> int: + n, stack = 0, [sid] + while stack: + cur = stack.pop() + for c in children[cur]: + if c not in rollover: + n += 1; stack.append(c) + return n + root = str(max((sid for sid in sessions if sessions[sid].get("parent_session_id") is None), key=size)) + depth: Dict[str, int] = {} + order: List[str] = [] + stack = [(root, 0)] + while stack: + sid, d = stack.pop() + depth[sid] = d; order.append(sid) + stack.extend((c, d + 1) for c in sorted(children[sid], key=lambda x: sessions[x].get("started_at") or 0, reverse=True) if c not in rollover) + order.sort(key=lambda s: sessions[s].get("started_at") or 0) + price = cls._fit_pricing([sessions[s] for s in order]) + out_dir = Path(out); out_dir.mkdir(parents=True, exist_ok=True) + return cls(Path(db), out_dir, root, sessions, depth, order, price, conn) + + @staticmethod + def _rollover_ids(conn: sqlite3.Connection, sessions: Dict[str, Dict[str, Any]]) -> set: + """Children created by in-place compression rollover, not by delegation: a child whose ``source`` + is a top-level surface (cli/tui/telegram/...), not ``subagent``, whose parent shares that source, + and which started within a few seconds of the parent ending. Their whole later lifetime belongs to + the continuing conversation, not to the fan-out, so they are kept out of the run population.""" + out = set() + for sid, s in sessions.items(): + pid = s.get("parent_session_id") + if not pid or pid not in sessions or (s.get("source") or "") == "subagent": + continue + parent = sessions[pid] + if (s.get("source") or "") != (parent.get("source") or ""): + continue + try: + gap = float(s.get("started_at") or 0) - float(parent.get("ended_at") or 0) + except (TypeError, ValueError): + continue + if -5.0 <= gap <= 5.0: + out.add(sid) + return out + + @staticmethod + def _fit_pricing(rows: Iterable[Dict[str, Any]]) -> Dict[str, float]: + X, y = [], [] + for s in rows: + cost = s.get("estimated_cost_usd") + if not cost: + continue + X.append([float(s.get(c) or 0) for c in USAGE_COLS]); y.append(float(cost)) + if len(X) < 8: + return {c: 0.0 for c in USAGE_COLS} + # Columns with negligible mass (e.g. input_tokens on cache-heavy Anthropic routes) make the + # normal equations ill-conditioned; fit only columns carrying >0.1% of all tokens. + mass = [sum(r[i] for r in X) for i in range(len(USAGE_COLS))] + keep = [i for i, m in enumerate(mass) if m > 0.001 * sum(mass)] + coef = _lstsq([[r[i] for i in keep] for r in X], y) + price = {c: 0.0 for c in USAGE_COLS} + for i, k in enumerate(keep): + price[USAGE_COLS[k]] = max(0.0, coef[i]) + return price + + # ── accessors ───────────────────────────────────────────────────────────────────────────── + def cost(self, sid: str) -> float: + return float(self.sessions[sid].get("estimated_cost_usd") or 0.0) + + def messages(self, sid: str, cols: str = "*") -> List[Dict[str, Any]]: + return [dict(r) for r in self._conn.execute(f"SELECT {cols} FROM messages WHERE session_id=? ORDER BY id", (sid,))] + + def iter_messages(self, sids: Iterable[str], cols: str = "*") -> Iterator[Dict[str, Any]]: + for sid in sids: + yield from self.messages(sid, cols) + + def system_prompt_len(self, sid: str) -> int: + h = self.sessions[sid].get("system_prompt_hash") + if not h: + return 0 + r = self._conn.execute("SELECT length(prompt) AS n FROM system_prompts WHERE hash=?", (h,)).fetchone() + return int(r["n"]) if r else 0 + + def by_depth(self) -> Dict[int, List[str]]: + out: Dict[int, List[str]] = collections.defaultdict(list) + for sid in self.in_run: + out[self.depth[sid]].append(sid) + return dict(out) + + def write(self, name: str, data: Any) -> Path: + p = self.out_dir / name + p.write_text(json.dumps(data, indent=1, default=str) if not name.endswith(".md") else str(data), encoding="utf-8") + return p + + def summary(self) -> Dict[str, Any]: + tot = {c: sum(float(self.sessions[s].get(c) or 0) for s in self.in_run) for c in USAGE_COLS} + return { + "root": self.root, "sessions": len(self.in_run), "children": len(self.in_run) - 1, + "by_depth": {d: len(v) for d, v in sorted(self.by_depth().items())}, + "api_calls": sum(int(self.sessions[s].get("api_call_count") or 0) for s in self.in_run), + "cost_usd": round(sum(self.cost(s) for s in self.in_run), 2), + "usage_tokens": tot, + "fitted_price_per_million": {c: round(p * 1e6, 4) for c, p in self.price_per_token.items()}, + "cost_by_bucket_usd": {c: round(tot[c] * self.price_per_token[c], 2) for c in USAGE_COLS}, + } diff --git a/evals/postmortem/forensics/delegation.py b/evals/postmortem/forensics/delegation.py new file mode 100644 index 0000000000..0ebe5beef5 --- /dev/null +++ b/evals/postmortem/forensics/delegation.py @@ -0,0 +1,115 @@ +"""Lane 2: delegation behaviour — nested-delegate timeouts, orphaned children, polling cost, batch-join +delivery delay, and summary truncation. All OBSERVED from ``messages`` in the run tree. + + python -m evals.postmortem.forensics.delegation --db state_copy.db + +Reproduces: how many ``delegate_task`` tool results were deadline timeouts and in how many sessions; +how many nested calls ever returned a result; what those sessions spent after their first timeout +(lifetime cost, with the caveat that it includes real work); total ``sleep N`` seconds; child results +withheld by batch joins (concurrent child-hours, NOT critical path); and how many parent-side summaries +carried the truncation footer. +""" +from __future__ import annotations + +import collections +import json +import re +import statistics +from typing import Any, Dict, List + +from evals.postmortem.forensics.common import Run + +_TIMEOUT = re.compile(r"timed out after ([\d.]+)s") +_SLEEP = re.compile(r"\bsleep\s+(\d+)") +_TRUNC = re.compile(r"\[SUMMARY TRUNCATED\]|middle omitted|trimmed to protect the parent") + + +def main(argv=None) -> int: + run = Run.from_args(argv, (__doc__ or "").split("\n\n")[0]) + parent_of = {s: run.sessions[s].get("parent_session_id") for s in run.in_run} + children_of: Dict[str, List[str]] = collections.defaultdict(list) + for s, p in parent_of.items(): + if p: + children_of[p].append(s) + orchestrators = [s for s in run.in_run if children_of.get(s)] + + timeouts_by_sess: Dict[str, int] = collections.Counter() + delegate_results = delegate_ok = 0 + sleep_seconds_after_timeout = 0 + first_timeout_ts: Dict[str, float] = {} + truncated_summaries = total_summaries = 0 + for sid in orchestrators: + for m in run.messages(sid, "role, tool_name, content, tool_calls, timestamp"): + if m["role"] == "tool" and (m.get("tool_name") or "") == "delegate_task": + c = m.get("content") or "" + delegate_results += 1 + if _TIMEOUT.search(c) and "delegate_task" in c: + timeouts_by_sess[sid] += 1 + first_timeout_ts.setdefault(sid, float(m.get("timestamp") or 0)) + elif '"status"' in c and ("completed" in c or "summary" in c): + delegate_ok += 1 + total_summaries += 1 + if _TRUNC.search(c): + truncated_summaries += 1 + if m["role"] == "user" and (m.get("content") or "").startswith("[ASYNC DELEGATION"): + # background batches re-enter as a user-role completion block, one entry per child + c = m["content"] + entries = c.count("--- ✓ TASK") + c.count("--- ✗ TASK") + c.count("--- ⚠ TASK") or (1 if "COMPLETE" in c else 0) + total_summaries += entries + truncated_summaries += c.count("[SUMMARY TRUNCATED]") + if m["role"] == "assistant" and sid in first_timeout_ts and float(m.get("timestamp") or 0) >= first_timeout_ts[sid]: + for mt in _SLEEP.finditer(m.get("tool_calls") or ""): + sleep_seconds_after_timeout += int(mt.group(1)) + + timeout_sessions = list(timeouts_by_sess) + nested_timeout_sessions = [s for s in timeout_sessions if run.depth[s] >= 1] + lifetime_cost = sum(run.cost(s) for s in timeout_sessions) + lifetime_cache_write = sum(float(run.sessions[s].get("cache_write_tokens") or 0) for s in timeout_sessions) * run.price_per_token["cache_write_tokens"] + + # batch-join delivery delay: children of one parent dispatched within 60 s of each other = one batch + withheld: Dict[int, List[float]] = collections.defaultdict(list) + for p, kids in children_of.items(): + kids = sorted(kids, key=lambda k: float(run.sessions[k].get("started_at") or 0)) + batch: List[str] = [] + def flush(batch: List[str]) -> None: + if len(batch) < 2: + return + ends = [float(run.sessions[k].get("ended_at") or 0) for k in batch] + last = max(ends) + withheld[run.depth.get(str(p), 0)].extend(last - e for e in ends if e) + for k in kids: + if batch and float(run.sessions[k].get("started_at") or 0) - float(run.sessions[batch[-1]].get("started_at") or 0) > 60: + flush(batch); batch = [] + batch.append(k) + flush(batch) + + report: Dict[str, Any] = { + "observed": { + "orchestrators": len(orchestrators), + "delegate_task_results": delegate_results, + "delegate_task_timeouts": sum(timeouts_by_sess.values()), + "sessions_with_timeouts": len(timeout_sessions), + "nested_sessions_with_timeouts": len(nested_timeout_sessions), + "delegate_task_results_that_carried_a_result": delegate_ok, + "sleep_hours_after_first_timeout": round(sleep_seconds_after_timeout / 3600, 1), + "timeout_sessions_lifetime_cost_usd": round(lifetime_cost, 2), + "timeout_sessions_lifetime_cache_write_usd": round(lifetime_cache_write, 2), + "lifetime_cost_caveat": "whole-session spend; includes legitimate work after the timeout, and its cache writes overlap the tokens lane's excess proxy", + "summaries_truncated": truncated_summaries, "summaries_total": total_summaries, + "batch_join_withheld_child_hours_by_parent_depth": {d: round(sum(v) / 3600, 1) for d, v in sorted(withheld.items())}, + "batch_join_withheld_minutes_median_by_parent_depth": {d: round(statistics.median(v) / 60, 1) for d, v in sorted(withheld.items()) if v}, + "withheld_caveat": "concurrent child-hours held back from the parent; not critical-path time", + } + } + path = run.write("delegation.json", report) + o = report["observed"] + print(f"[delegation] {o['orchestrators']} orchestrators; delegate_task results {o['delegate_task_results']:,}: timeouts {o['delegate_task_timeouts']} in {o['sessions_with_timeouts']} sessions " + f"({o['nested_sessions_with_timeouts']} nested); carried a result: {o['delegate_task_results_that_carried_a_result']}") + print(f"[delegation] after first timeout: sleep {o['sleep_hours_after_first_timeout']} h; those sessions' lifetime ${o['timeout_sessions_lifetime_cost_usd']:,} (cache writes ${o['timeout_sessions_lifetime_cache_write_usd']:,}; includes real work)") + print(f"[delegation] summaries truncated {o['summaries_truncated']}/{o['summaries_total']}; batch-join withheld child-hours by parent depth {o['batch_join_withheld_child_hours_by_parent_depth']}") + print(f"[delegation] wrote {path}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/forensics/goal_loop.py b/evals/postmortem/forensics/goal_loop.py new file mode 100644 index 0000000000..87200f1ab7 --- /dev/null +++ b/evals/postmortem/forensics/goal_loop.py @@ -0,0 +1,70 @@ +"""Lane 4: /goal loop behaviour — nudges, judge verdicts, parked time. OBSERVED from the root session's +messages and the persisted goal state. + + python -m evals.postmortem.forensics.goal_loop --db state_copy.db + +Reproduces: how many ``[Continuing toward your standing goal]`` nudges fired and how soon after the +previous assistant turn; how many of those turns had just said they were waiting; the goal's final +persisted wait barrier (what it was parked on, since when); and the count of async-batch and +background-process notifications that re-entered the root. +""" +from __future__ import annotations + +import json +import re +from typing import Any, Dict, List + +from evals.postmortem.forensics.common import Run + +_WAITING = re.compile(r"\b(waiting|wait(s)? on|nothing (else |new )?to dispatch|no action needed|standing by|until .* (finish|return|complete)|in flight|still running)\b", re.I) + + +def main(argv=None) -> int: + run = Run.from_args(argv, (__doc__ or "").split("\n\n")[0]) + root = run.root + msgs = run.messages(root, "role, content, timestamp") + nudges: List[Dict[str, Any]] = [] + batch_notices = bgproc_notices = 0 + last_assistant = None + for m in msgs: + c = m.get("content") or "" + if m["role"] == "assistant": + last_assistant = m + elif m["role"] == "user": + if c.startswith("[Continuing toward your standing goal]"): + gap = float(m["timestamp"]) - float(last_assistant["timestamp"]) if last_assistant else None + nudges.append({"ts": m["timestamp"], "gap_s": round(gap, 1) if gap is not None else None, + "prev_turn_said_waiting": bool(last_assistant and _WAITING.search(last_assistant.get("content") or ""))}) + elif c.startswith("[ASYNC DELEGATION"): + batch_notices += 1 + elif c.startswith("[IMPORTANT: Background process"): + bgproc_notices += 1 + goal_state: Dict[str, Any] = {} + try: + row = run._conn.execute("SELECT value FROM state_meta WHERE key=?", (f"goal:{root}",)).fetchone() + if row: + goal_state = json.loads(row[0]) + except Exception: + pass + parked = {k: goal_state.get(k) for k in ("status", "waiting_on_pid", "waiting_on_session", "waiting_until", "waiting_since", "waiting_reason", "last_verdict", "turns_used")} + if goal_state.get("waiting_since"): + # The barrier is never auto-cleared unless a turn re-evaluates it, so its age at session end is + # the time the loop sat parked (the state may have been cleared by the user afterwards). + end = float(run.sessions[root].get("ended_at") or 0) + parked["parked_minutes_until_session_end"] = round((end - float(goal_state["waiting_since"])) / 60, 1) if end else None + report = {"observed": { + "nudges": len(nudges), "nudges_within_180s_of_a_waiting_turn": sum(1 for n in nudges if n["prev_turn_said_waiting"] and (n["gap_s"] or 1e9) < 180), + "nudge_detail": nudges, "async_batch_notices": batch_notices, "background_process_notices": bgproc_notices, + "final_goal_state": parked, + }} + path = run.write("goal_loop.json", report) + o = report["observed"] + print(f"[goal] nudges {o['nudges']}, of which {o['nudges_within_180s_of_a_waiting_turn']} fired <180 s after a turn that said it was waiting; " + f"batch notices {o['async_batch_notices']}, bg-process notices {o['background_process_notices']}") + print(f"[goal] final goal state: {o['final_goal_state']}") + print(f"[goal] wrote {path}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/forensics/logcalls.py b/evals/postmortem/forensics/logcalls.py new file mode 100644 index 0000000000..c8457bd9c3 --- /dev/null +++ b/evals/postmortem/forensics/logcalls.py @@ -0,0 +1,136 @@ +"""Lane 1b: per-call cache behaviour from ``agent.log`` (OBSERVED provider usage per call). + +Hermes logs one line per API call:: + + ... INFO [] agent.conversation_loop: API call #N: model=... in= out= total=... latency=..s cache=/ + +Given the rotated logs, this reproduces: prompt-size distribution, cache hit-ratio buckets, the +"plateau" signature of a broken cache prefix (hit count stuck at the previous call's breakpoint while +``in`` grows), the share of uncached input those plateaus explain, and the sawtooth replay of an absolute +context cap on REAL per-call prompt sizes (the figure the tokens lane reported as -49%). + + python -m evals.postmortem.forensics.logcalls --db state_copy.db --logs ~/.hermes/logs/agent.log* + +Coverage caveat: logs rotate; report the fraction of the run's calls that were found before quoting +anything from this lane, and treat extrapolations as upper bounds. +""" +from __future__ import annotations + +import collections +import glob +import re +import statistics +from typing import Any, Dict, List + +from evals.postmortem.forensics.common import Run + +_LINE = re.compile( + r"^(\S+ \S+) INFO \[(\S+)\] agent\.conversation_loop: API call #(\d+): model=(\S+) provider=\S+ " + r"in=(\d+) out=(\d+) total=\d+ latency=([\d.]+)s cache=(\d+)/(\d+)" +) + + +def parse_logs(paths: List[str], sids: set) -> List[Dict[str, Any]]: + calls: List[Dict[str, Any]] = [] + for path in paths: + try: + fh = open(path, encoding="utf-8", errors="replace") + except OSError: + continue + with fh: + for line in fh: + m = _LINE.match(line) + if m and m.group(2) in sids: + calls.append({"ts": m.group(1), "sid": m.group(2), "n": int(m.group(3)), "model": m.group(4), + "inp": int(m.group(5)), "out": int(m.group(6)), "lat": float(m.group(7)), + "hit": int(m.group(8))}) + return calls + + +def sawtooth_real(by_sid: Dict[str, List[Dict[str, Any]]], cap: int, floor: int) -> float: + """Replay on REAL prompt sizes: appended = in[i] - in[i-1]; compress to floor when above cap.""" + total = 0.0 + for calls in by_sid.values(): + calls = sorted(calls, key=lambda c: c["n"]) + ctx = None + for i, c in enumerate(calls): + appended = c["inp"] if i == 0 else max(0, c["inp"] - calls[i - 1]["inp"]) + ctx = c["inp"] if ctx is None else ctx + appended + if ctx > cap: + ctx = float(floor) + total += ctx + return total + + +def main(argv=None) -> int: + ap = Run.parser((__doc__ or "").split("\n\n")[0]) + ap.add_argument("--logs", nargs="+", required=True, help="agent.log files (globs ok)") + ap.add_argument("--cap", type=int, default=200_000) + ap.add_argument("--floor", type=int, default=65_000) + a = ap.parse_args(argv) + run = Run.open(a.db, root=a.root, out=a.out) + paths = sorted(p for g in a.logs for p in glob.glob(g)) + calls = parse_logs(paths, set(run.in_run)) + total_calls = run.summary()["api_calls"] + if not calls: + print("[logcalls] no matching API-call lines found in", paths); return 1 + inp = [c["inp"] for c in calls] + hit_total, in_total = sum(c["hit"] for c in calls), sum(inp) + buckets = collections.Counter() + for c in calls: + r = c["hit"] / c["inp"] if c["inp"] else 0 + buckets["<50%" if r < .5 else "50-90%" if r < .9 else "90-97%" if r < .97 else "97-99%" if r < .99 else ">=99%"] += 1 + by_sid: Dict[str, List[Dict[str, Any]]] = collections.defaultdict(list) + for c in calls: + by_sid[c["sid"]].append(c) + # Two definitions, reported separately (an independent review caught them being conflated): + # strict plateau: hit count EXACTLY equal to the previous call's (prefix stuck at the old breakpoint) + # non-advancing: hit count did not grow while input did (includes partial misses of other causes) + strict_pairs = strict_uncached = loose_pairs = loose_uncached = pairs = uncached_total = 0 + for cs in by_sid.values(): + cs.sort(key=lambda c: c["n"]) + for prev, cur in zip(cs, cs[1:]): + pairs += 1 + u = max(0, cur["inp"] - cur["hit"]); uncached_total += u + if cur["inp"] > prev["inp"]: + if cur["hit"] == prev["hit"]: + strict_pairs += 1; strict_uncached += u + if cur["hit"] <= prev["hit"]: + loose_pairs += 1; loose_uncached += u + price = run.price_per_token + real = sum(inp); capped = sawtooth_real(by_sid, a.cap, a.floor) + report = { + "coverage": {"calls_found": len(calls), "run_api_calls": total_calls, "fraction": round(len(calls) / total_calls, 4) if total_calls else None, + "first_ts": min(c["ts"] for c in calls), "last_ts": max(c["ts"] for c in calls), + "note": "rotated logs; extrapolations from this window are upper bounds"}, + "observed": { + "prompt_tokens": {"median": int(statistics.median(inp)), "p90": sorted(inp)[int(.9 * len(inp))], "max": max(inp)}, + "share_of_calls_above": {t: round(sum(1 for x in inp if x > t) / len(inp), 3) for t in (150_000, 200_000, 300_000)}, + "cache_hit_ratio_overall": round(hit_total / in_total, 4), + "hit_ratio_buckets": dict(buckets), + "strict_plateau": {"pairs_share": round(strict_pairs / pairs, 4) if pairs else None, + "share_of_uncached_input": round(strict_uncached / uncached_total, 4) if uncached_total else None}, + "non_advancing_hit": {"pairs_share": round(loose_pairs / pairs, 4) if pairs else None, + "share_of_uncached_input": round(loose_uncached / uncached_total, 4) if uncached_total else None}, + "uncached_input_usd_in_window": round(uncached_total * price["cache_write_tokens"], 2), + }, + "modeled": { + "sawtooth_on_real_prompt_sizes": { + "cap": a.cap, "floor": a.floor, "prompt_tokens_ratio": round(capped / real, 3) if real else None, + "caveats": ["compression-call cost excluded", "read/write ratio assumed unchanged", "window only"], + } + }, + } + path = run.write("logcalls.json", report) + o = report["observed"] + print(f"[logcalls] coverage {len(calls):,}/{total_calls:,} calls ({report['coverage']['fraction']:.1%}) from {len(paths)} file(s)") + print(f"[logcalls] OBSERVED median prompt {o['prompt_tokens']['median']:,} p90 {o['prompt_tokens']['p90']:,}; >200K: {o['share_of_calls_above'][200000]:.0%}; hit ratio {o['cache_hit_ratio_overall']:.1%}") + print(f"[logcalls] OBSERVED strict plateau (hit unchanged): {o['strict_plateau']['pairs_share']:.1%} of pairs, {o['strict_plateau']['share_of_uncached_input']:.1%} of uncached input; " + f"non-advancing hit: {o['non_advancing_hit']['pairs_share']:.1%} / {o['non_advancing_hit']['share_of_uncached_input']:.1%}; uncached in window ${o['uncached_input_usd_in_window']:,}") + print(f"[logcalls] MODELED sawtooth cap {a.cap}/{a.floor} on real sizes: prompt volume x{report['modeled']['sawtooth_on_real_prompt_sizes']['prompt_tokens_ratio']}") + print(f"[logcalls] wrote {path}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/forensics/rework.py b/evals/postmortem/forensics/rework.py new file mode 100644 index 0000000000..612fb8acb8 --- /dev/null +++ b/evals/postmortem/forensics/rework.py @@ -0,0 +1,76 @@ +"""Lane 5: post-open rework inventory for a large PR — reproducible part only. + + python -m evals.postmortem.forensics.rework --repo . --base --open --head + +Reports (OBSERVED from git): the public-surface drop at PR open (via ``scripts/ci/check_public_surface.py``), +then every commit between the opening SHA and the merged head grouped by subject prefix +(``fix``, ``review-fix``, ``revert``, ``test``, ``docs``, ``simplify``, other) with files/insertions/deletions +and whether tests were touched. Root-cause classification of each commit (dropped symbol vs semantic +drift vs compat fallout ...) was done BY HAND in the original audit and is not reproducible here; the +counts this script prints are the mechanical inventory that classification started from. +""" +from __future__ import annotations + +import argparse +import collections +import json +import re +import subprocess +import sys +from pathlib import Path + + +def _git(repo: str, *args: str) -> str: + return subprocess.run(["git", "-C", repo, *args], capture_output=True, text=True, encoding="utf-8", errors="replace").stdout + + +def main(argv=None) -> int: + ap = argparse.ArgumentParser(description=(__doc__ or "").split("\n\n")[0]) + ap.add_argument("--repo", default=".") + ap.add_argument("--base", required=True, help="merge-base of the PR") + ap.add_argument("--open", required=True, help="head SHA when the PR was opened / first reviewed") + ap.add_argument("--head", required=True, help="final merged SHA") + ap.add_argument("--out", default="postmortem_out") + a = ap.parse_args(argv) + out = Path(a.out); out.mkdir(parents=True, exist_ok=True) + + # scripts/ci/check_public_surface.py ships with #103541; look for it in this tree, then the target repo. + candidates = [Path(__file__).resolve().parents[3] / "scripts" / "ci" / "check_public_surface.py", + Path(a.repo).resolve() / "scripts" / "ci" / "check_public_surface.py"] + checker = next((c for c in candidates if c.exists()), None) + if checker is None: + surface_line = "(scripts/ci/check_public_surface.py not found; merge #103541 or pass a checkout that has it)" + else: + surface = subprocess.run([sys.executable, str(checker), "--base", a.base, "--head", a.open, "--json", str(out / "surface_at_open.json")], + cwd=a.repo, capture_output=True, text=True, encoding="utf-8", errors="replace") + surface_line = (surface.stdout.splitlines() or [f"(check_public_surface failed: {surface.stderr.strip()[:120]})"])[0] + + log = _git(a.repo, "log", "--no-merges", "--format=%H%x00%s", f"{a.open}..{a.head}") + groups = collections.defaultdict(list) + for line in log.splitlines(): + if not line.strip(): + continue + sha, subj = line.split("\x00", 1) + m = re.match(r"^([a-z-]+)(\(|:|!)", subj) + known = {"fix", "review-fix", "revert", "test", "docs", "simplify", "feat", "chore", "ci"} + prefix = m.group(1) if m and m.group(1) in known else "other" + files = _git(a.repo, "show", "--format=", "--name-only", sha).split() + stat = _git(a.repo, "show", "--format=", "--shortstat", sha) + ins = int((re.search(r"(\d+) insertion", stat) or [0, 0])[1]); dele = int((re.search(r"(\d+) deletion", stat) or [0, 0])[1]) + tests = [f for f in files if f.startswith("tests/") or "/tests/" in f] + groups[prefix].append({"sha": sha[:10], "subject": subj[:120], "files": len(files), "ins": ins, "del": dele, + "test_files": len(tests), "src_only": len(tests) == 0 and len(files) > 0}) + summary = {g: {"commits": len(v), "files": sum(x["files"] for x in v), "ins": sum(x["ins"] for x in v), "del": sum(x["del"] for x in v), + "touching_no_tests": sum(1 for x in v if x["src_only"])} for g, v in groups.items()} + report = {"surface_at_open": surface_line, "post_open_commits_by_prefix": summary, "commits": groups, + "note": "root-cause classes were hand-labeled in the original audit; not reproduced here"} + (out / "rework.json").write_text(json.dumps(report, indent=1), encoding="utf-8") + print(f"[rework] surface at open: {surface_line}") + total = sum(v["commits"] for v in summary.values()) + print(f"[rework] {total} post-open commits by prefix: " + ", ".join(f"{g} {v['commits']} ({v['touching_no_tests']} w/o tests)" for g, v in sorted(summary.items(), key=lambda kv: -kv[1]['commits']))) + print(f"[rework] wrote {out / 'rework.json'}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/forensics/tokens.py b/evals/postmortem/forensics/tokens.py new file mode 100644 index 0000000000..2312ed0238 --- /dev/null +++ b/evals/postmortem/forensics/tokens.py @@ -0,0 +1,122 @@ +"""Lane 1: where the money went, context sizes, and the compression-cap counterfactual. + +Reproduces (for any run): cost by bucket and depth, the share of calls above context thresholds, +the excess-cache-write proxy, and the sawtooth replay of an absolute subagent context cap. + + python -m evals.postmortem.forensics.tokens --db state_copy.db [--cap 200000 --floor 65000] + +Every figure is labeled in the output as OBSERVED (from usage rows) or MODELED (reconstruction / +replay). The per-call context reconstruction estimates tokens from message character counts +(chars/3.5) plus the system-prompt and tool-schema sizes; it is a diagnostic, not a measurement. +""" +from __future__ import annotations + +import collections +import statistics +from typing import Any, Dict, List + +from evals.postmortem.forensics.common import Run + +CHARS_PER_TOKEN = 3.5 +TOOL_SCHEMA_TOKENS = 12_000 + + +def _per_call_context(run: Run, sid: str) -> List[Dict[str, Any]]: + """Reconstructed prompt size at each assistant turn (MODELED).""" + ctx = run.system_prompt_len(sid) / CHARS_PER_TOKEN + TOOL_SCHEMA_TOKENS + calls, appended = [], 0.0 + for m in run.messages(sid, "role, length(coalesce(content,'')) AS lc, length(coalesce(tool_calls,'')) AS ltc"): + if m["role"] == "assistant": + calls.append({"ctx": ctx, "appended_since_last": appended}) + out = (m["lc"] + m["ltc"]) / CHARS_PER_TOKEN + ctx += out; appended = out + else: + ctx += m["lc"] / CHARS_PER_TOKEN; appended += m["lc"] / CHARS_PER_TOKEN + return calls + + +def sawtooth(calls: List[Dict[str, Any]], cap: int, floor: int) -> float: + """Prompt tokens if the session compressed to ``floor`` whenever ctx exceeded ``cap`` (MODELED).""" + total, ctx = 0.0, None + for c in calls: + ctx = c["ctx"] if ctx is None else ctx + c["appended_since_last"] + if ctx > cap: + ctx = float(floor) + total += ctx + return total + + +def main(argv=None) -> int: + ap = Run.parser((__doc__ or "").split("\n\n")[0]) + ap.add_argument("--cap", type=int, default=200_000) + ap.add_argument("--floor", type=int, default=65_000) + a = ap.parse_args(argv) + run = Run.open(a.db, root=a.root, out=a.out) + price = run.price_per_token + report: Dict[str, Any] = {"summary": run.summary()} + + # OBSERVED: cost by depth, by duration bucket, top-N concentration + by_depth = collections.defaultdict(float) + for sid in run.in_run: + by_depth[run.depth[sid]] += run.cost(sid) + costs = sorted((run.cost(s) for s in run.in_run), reverse=True) + total = sum(costs) or 1.0 + long_sessions = [s for s in run.in_run if (float(run.sessions[s].get("ended_at") or 0) - float(run.sessions[s].get("started_at") or 0)) > 3600] + report["observed"] = { + "cost_by_depth_usd": {d: round(v, 2) for d, v in sorted(by_depth.items())}, + "top30_share": round(sum(costs[:30]) / total, 3), + "sessions_over_60min": len(long_sessions), + "sessions_over_60min_cost_share": round(sum(run.cost(s) for s in long_sessions) / total, 3), + } + + # MODELED: per-call context distribution, excess cache-write proxy, sawtooth replay + all_calls, per_session = [], {} + est_total = actual_total = appended_total = 0.0 + for sid in run.in_run: + calls = _per_call_context(run, sid) + if not calls: + continue + per_session[sid] = calls + all_calls.extend(c["ctx"] for c in calls) + appended_total += sum(c["appended_since_last"] for c in calls) + calls[0]["ctx"] + est_total += sum(c["ctx"] for c in calls) + s = run.sessions[sid] + actual_total += float(s.get("cache_read_tokens") or 0) + float(s.get("cache_write_tokens") or 0) + float(s.get("input_tokens") or 0) + n = len(all_calls) or 1 + thresholds = {t: round(sum(1 for c in all_calls if c > t) / n, 3) for t in (100_000, 150_000, 200_000, 300_000)} + cache_write_tokens = sum(float(run.sessions[s].get("cache_write_tokens") or 0) for s in run.in_run) + excess = max(0.0, cache_write_tokens - appended_total) + base_prompt = sum(sum(c["ctx"] for c in v) for v in per_session.values()) + capped_prompt = sum(sawtooth(v, a.cap, a.floor) for v in per_session.values()) + ctx_spend = sum(float(run.sessions[s].get(c) or 0) * price[c] for s in run.in_run for c in ("cache_read_tokens", "cache_write_tokens")) + report["modeled"] = { + "note": "reconstructed from message sizes at chars/3.5; diagnostic, not a measurement", + "calls_reconstructed": len(all_calls), + "median_prompt_tokens": int(statistics.median(all_calls)) if all_calls else 0, + "share_of_calls_above": thresholds, + "reconstruction_vs_actual_ratio": round(est_total / actual_total, 3) if actual_total else None, + "excess_cache_write_proxy": { + "ideal_tokens_if_only_appended": int(appended_total), "actual_cache_write_tokens": int(cache_write_tokens), + "excess_tokens": int(excess), "excess_usd": round(excess * price["cache_write_tokens"], 2), + "caveats": ["token counts estimated", "ideal omits initial system/tool prefix writes", "reasoning accumulation omitted"], + }, + "sawtooth_replay": { + "cap": a.cap, "floor": a.floor, + "prompt_tokens_ratio_capped_vs_actual": round(capped_prompt / base_prompt, 3) if base_prompt else None, + "context_spend_usd_observed": round(ctx_spend, 2), + "context_spend_usd_if_capped": round(ctx_spend * capped_prompt / base_prompt, 2) if base_prompt else None, + "caveats": ["excludes compression-call cost", "assumes read/write ratio unchanged", "ignores re-reads and quality effects"], + }, + } + path = run.write("tokens.json", report) + o, m = report["observed"], report["modeled"] + print(f"[tokens] {run.summary()['sessions']} sessions, ${run.summary()['cost_usd']:,} | buckets {run.summary()['cost_by_bucket_usd']}") + print(f"[tokens] OBSERVED depth-2 share {by_depth.get(2,0)/total:.0%}, >60min share {o['sessions_over_60min_cost_share']:.0%}, top30 {o['top30_share']:.1%}") + print(f"[tokens] MODELED calls >150K ctx {m['share_of_calls_above'][150000]:.0%}; excess cache-write proxy ${m['excess_cache_write_proxy']['excess_usd']:,}") + print(f"[tokens] MODELED sawtooth cap {a.cap}: prompt volume x{m['sawtooth_replay']['prompt_tokens_ratio_capped_vs_actual']} -> context spend ${m['sawtooth_replay']['context_spend_usd_observed']:,} -> ${m['sawtooth_replay']['context_spend_usd_if_capped']:,}") + print(f"[tokens] wrote {path}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/forensics/tools.py b/evals/postmortem/forensics/tools.py new file mode 100644 index 0000000000..f824bc2fed --- /dev/null +++ b/evals/postmortem/forensics/tools.py @@ -0,0 +1,100 @@ +"""Lane 3: tool-layer friction — hardline false blocks, foreground refusals, whole-file rewrites, output +volume by tool. All OBSERVED from ``messages`` (role='tool') in the run tree. + + python -m evals.postmortem.forensics.tools --db state_copy.db +""" +from __future__ import annotations + +import collections +import json +import re +from typing import Any, Dict + +from evals.postmortem.forensics.common import Run + +_HARDLINE = "BLOCKED (hardline)" +_MALFORMED = "command parser limit or malformed executable payload" +_FG_TIMEOUT = re.compile(r"Foreground timeout (\d+)s exceeds the maximum") +_FG_AMP = "Foreground command uses '&' backgrounding" +_FG_WRAP = "shell-level background wrappers" +_DEADLINE = re.compile(r"timed out after ([\d.]+)s") + + +def main(argv=None) -> int: + run = Run.from_args(argv, (__doc__ or "").split("\n\n")[0]) + by_tool_bytes: Dict[str, int] = collections.Counter() + by_tool_calls: Dict[str, int] = collections.Counter() + hardline = malformed = fg_timeout = fg_amp = fg_wrap = tool_deadline = 0 + fg_timeout_requested: Dict[int, int] = collections.Counter() + write_calls = big_writes = rewrite_of_existing = 0 + write_chars = 0 + read_paths_by_sess: Dict[str, set] = collections.defaultdict(set) + for sid in run.in_run: + for m in run.messages(sid, "role, tool_name, content, tool_calls"): + if m["role"] == "assistant" and m.get("tool_calls"): + try: + tcs = json.loads(m["tool_calls"]) + except ValueError: + continue + for tc in tcs: + fn = (tc.get("function") or {}) + name = fn.get("name") or "" + by_tool_calls[name] += 1 + if name in ("read_file", "write_file", "patch"): + try: + args = json.loads(fn.get("arguments") or "{}") + except ValueError: + args = {} + if name == "read_file" and args.get("path"): + read_paths_by_sess[sid].add(args["path"]) + if name == "write_file": + write_calls += 1 + content = args.get("content") or "" + write_chars += len(content) + if len(content) > 20_000: + big_writes += 1 + if args.get("path") in read_paths_by_sess[sid]: + rewrite_of_existing += 1 + continue + if m["role"] != "tool": + continue + c = m.get("content") or "" + by_tool_bytes[m.get("tool_name") or "?"] += len(c) + if _HARDLINE in c: + hardline += 1 + if _MALFORMED in c: + malformed += 1 + mt = _FG_TIMEOUT.search(c) + if mt: + fg_timeout += 1; fg_timeout_requested[int(mt.group(1))] += 1 + if _FG_AMP in c: + fg_amp += 1 + if _FG_WRAP in c: + fg_wrap += 1 + if (m.get("tool_name") or "") == "terminal" and _DEADLINE.search(c): + tool_deadline += 1 + top_bytes = sorted(by_tool_bytes.items(), key=lambda kv: -kv[1])[:8] + report: Dict[str, Any] = {"observed": { + "tool_results": sum(by_tool_calls.values()), + "output_bytes_by_tool_top8": {k: v for k, v in top_bytes}, + "hardline_blocks": hardline, "hardline_blocks_malformed_class": malformed, + "foreground_timeout_refusals": fg_timeout, "foreground_timeout_requested_top": dict(fg_timeout_requested.most_common(5)), + "foreground_ampersand_refusals": fg_amp, "foreground_wrapper_refusals": fg_wrap, + "terminal_deadline_kills": tool_deadline, + "write_file": {"calls": write_calls, "chars": write_chars, "writes_over_20k_chars": big_writes, + "rewrites_of_a_file_read_this_session_over_20k": rewrite_of_existing, + "output_usd_at_fitted_price": round(write_chars / 3.5 * run.price_per_token["output_tokens"], 2)}, + "patch_calls": by_tool_calls.get("patch", 0), + }} + path = run.write("tools.json", report) + o = report["observed"] + print(f"[tools] {o['tool_results']:,} tool calls; hardline blocks {o['hardline_blocks']} ({o['hardline_blocks_malformed_class']} 'malformed' class); " + f"foreground timeout refusals {o['foreground_timeout_refusals']} (asked: {o['foreground_timeout_requested_top']}); '&' {o['foreground_ampersand_refusals']}, nohup {o['foreground_wrapper_refusals']}") + w = o["write_file"] + print(f"[tools] write_file {w['calls']:,} calls, {w['chars']/1e6:.1f}M chars (~${w['output_usd_at_fitted_price']:,}); >20k: {w['writes_over_20k_chars']}, of which rewrites of a file read this session: {w['rewrites_of_a_file_read_this_session_over_20k']}; patch calls {o['patch_calls']:,}") + print(f"[tools] wrote {path}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/live_ab/__init__.py b/evals/postmortem/live_ab/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/evals/postmortem/live_ab/auth_stampede.py b/evals/postmortem/live_ab/auth_stampede.py new file mode 100644 index 0000000000..1719faa023 --- /dev/null +++ b/evals/postmortem/live_ab/auth_stampede.py @@ -0,0 +1,79 @@ +"""Live A/B of the Nous hourly-expiry stampede, no real network. + +Local server: accepts bearer FRESH, returns 401 {"type":"authentication_error", "Your API key is invalid, +blocked or out of funds..."} for any other bearer. resolve_nous_runtime_credentials is patched to return +FRESH (standing in for the auth store the keepalive/peers have refreshed). N agents are built holding a +STALE JWT that expires in 30 s and each fires one API call concurrently. Count 401s the server saw. + +Usage: python stampede_ab.py +""" +import base64, json, os, sys, tempfile, threading, time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + +root, n = sys.argv[1], int(sys.argv[2]) +sys.path.insert(0, root) +os.environ["HERMES_HOME"] = tempfile.mkdtemp(prefix="hh-") +os.environ["HERMES_STREAM_RETRIES"] = "0" + + +def jwt(exp, sub="acct-A"): + # `sub` matters: pre-expiry adoption (#103526 round 2) only swaps to a key for the SAME account. + b = lambda o: base64.urlsafe_b64encode(json.dumps(o).encode()).rstrip(b"=").decode() # noqa: E731 + return f"{b({'alg':'none'})}.{b({'exp':exp,'sub':sub})}.s" + + +STALE, FRESH = jwt(time.time() + 30), jwt(time.time() + 3600) +hits = {"401": 0, "200": 0} +lock = threading.Lock() + + +class H(BaseHTTPRequestHandler): + def log_message(self, *a): pass + def do_POST(self): + self.rfile.read(int(self.headers.get("content-length", 0))) + auth = self.headers.get("authorization", "") + if not self.path.endswith("/chat/completions"): + self.send_response(404); self.send_header("content-length", "0"); self.end_headers(); return + if auth == f"Bearer {FRESH}": + body = json.dumps({"id": "x", "object": "chat.completion", "created": 0, "model": "m", + "choices": [{"index": 0, "message": {"role": "assistant", "content": "ok"}, "finish_reason": "stop"}], + "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}}).encode() + with lock: hits["200"] += 1 + self.send_response(200) + else: + if os.environ.get("TRACE401"): + import traceback; sys.stderr.write("401 path: "+self.path+"\n") + body = json.dumps({"error": {"type": "authentication_error", "message": "Your API key is invalid, blocked or out of funds. Please go visit the portal to sort that out: https://portal.nousresearch.com "}}).encode() + with lock: hits["401"] += 1 + self.send_response(401) + self.send_header("content-type", "application/json"); self.send_header("content-length", str(len(body))); self.end_headers(); self.wfile.write(body) + + +srv = ThreadingHTTPServer(("127.0.0.1", 0), H); threading.Thread(target=srv.serve_forever, daemon=True).start() +base = f"http://127.0.0.1:{srv.server_address[1]}/v1" + +import hermes_cli.auth as auth_mod +auth_mod.resolve_nous_runtime_credentials = lambda **kw: {"api_key": FRESH, "base_url": base} +import hermes_cli.nous_auth_keepalive as ka +ka.start_nous_auth_keepalive = lambda **kw: None # thread itself is out of scope here; we test adoption + +from run_agent import AIAgent + +agents = [] +for i in range(n): + a = AIAgent(api_key=STALE, base_url=base, provider="nous", model="test/model", quiet_mode=True, + skip_context_files=True, skip_memory=True) + a.api_mode = "chat_completions" + a._interrupt_requested = False + agents.append(a) + +results = [] +def go(a): + try: + results.append(a.chat("hi")[:20]); results.append("KEY:"+a.api_key[:5]) + except Exception as e: + results.append(f"ERR {type(e).__name__}") + +ts = [threading.Thread(target=go, args=(a,)) for a in agents] +t0 = time.time(); [t.start() for t in ts]; [t.join(120) for t in ts] +print("keys:", sorted(set(r for r in results if r.startswith("KEY"))));print(f"agents={n} server_401={hits['401']} server_200={hits['200']} errors={sum(r.startswith('ERR') for r in results)} wall={time.time()-t0:.1f}s") diff --git a/evals/postmortem/live_ab/batch_failure_notice.py b/evals/postmortem/live_ab/batch_failure_notice.py new file mode 100644 index 0000000000..d713e8f3bf --- /dev/null +++ b/evals/postmortem/live_ab/batch_failure_notice.py @@ -0,0 +1,42 @@ +"""Live: detached delegate_task batch, one child fails immediately, siblings run ~8 s. Does the parent's +completion queue see the per-task failure notice BEFORE the consolidated batch result? +Usage: python notice_live.py """ +import json, os, sys, tempfile, threading, time +root = sys.argv[1]; sys.path.insert(0, root) +os.environ["HERMES_HOME"] = tempfile.mkdtemp(prefix="hh-") +os.environ["HERMES_STREAM_RETRIES"] = "0" + +from run_agent import AIAgent +import tools.delegate_tool as dt +import tools.delegate_tool_dispatch as dd +from tools.process_registry import process_registry + +parent = AIAgent(api_key="k", base_url="https://example.com/v1", provider="test-provider", model="test/model", + quiet_mode=True, skip_context_files=True, skip_memory=True) +parent.api_mode = "chat_completions" + +# Stub the child run: task 0 fails instantly, others take 8 s and succeed. +orig_run = dd._Batch.run_child if hasattr(dd, "_Batch") else None +def fake_run_child(self, idx, task, child): + if idx == 0: + return {"task_index": 0, "status": "error", "error": "simulated 401 authentication_error", "api_calls": 0, "duration_seconds": 0.1} + time.sleep(8) + return {"task_index": idx, "status": "completed", "summary": f"ok {idx}", "api_calls": 1, "duration_seconds": 8.0} +dd._Batch.run_child = fake_run_child + +# Background dispatch requires an async-capable session; emulate a CLI session key. +import tools.async_delegation as ad +t0 = time.time() +out = dt.delegate_task(tasks=[{"goal": "fail fast on a simulated 401 error"}, {"goal": "slow successful worker number one"}, {"goal": "slow successful worker number two"}], background=True, parent_agent=parent) +print("dispatch:", json.dumps(json.loads(out))[:160]) +seen = [] +deadline = time.time() + 30 +while time.time() < deadline: + for evt, text in process_registry.drain_notifications(session_key=""): + kind = "TASK_FAILURE_NOTICE" if evt.get("task_failure_notice") else ("BATCH_FINAL" if evt.get("is_batch") else evt.get("type")) + seen.append((round(time.time() - t0, 1), kind)) + print(f" t+{time.time()-t0:5.1f}s {kind}: {text.splitlines()[0][:100]}") + if any(k == "BATCH_FINAL" for _, k in seen): + break + time.sleep(0.3) +print("ORDER:", seen) diff --git a/evals/postmortem/live_ab/cache_prefix_live.py b/evals/postmortem/live_ab/cache_prefix_live.py new file mode 100644 index 0000000000..e5e3554302 --- /dev/null +++ b/evals/postmortem/live_ab/cache_prefix_live.py @@ -0,0 +1,49 @@ +"""Live A/B for the thinking-strip cache miss (F0). Runs a short real tool loop through AIAgent on +Fable 5.1 via Nous and prints per-call cache hit ratios from agent.log. ~10 calls, well under $1. +Arm A = current code. Arm B = HERMES_KEEP_ALL_THINKING=1 monkeypatch of _manage_thinking_signatures +that passes thinking blocks back unchanged for the Nous/Anthropic route.""" +import os, sys, re, time, subprocess, json +# LIVE: real provider calls (cents). Usage: python cache_prefix_live.py +os.environ.setdefault("HERMES_HOME", os.path.expanduser("~/.hermes")) +sys.path.insert(0, sys.argv[1]) +arm = sys.argv[2] if len(sys.argv) > 2 else "A" +import agent.anthropic_message_convert as amc +if arm == "B": + _orig = amc._manage_thinking_signatures + def _keep_all(result, base_url, model): + # keep-all model on a signature-validating route: pass every thinking block back unchanged, + # only drop cache_control on thinking blocks and the internal flag. + for idx, m in amc._assistant_block_lists(result): + for b in m["content"]: + if amc._block_type(b) in amc._THINKING_TYPES: + b.pop("cache_control", None) + m.pop("_thinking_signature_invalidated", None) + amc._manage_thinking_signatures = _keep_all + # the adapter imports the name at call time via module attribute? verify binding + import agent.anthropic_adapter as ad + if hasattr(ad, "_manage_thinking_signatures"): + ad._manage_thinking_signatures = _keep_all +from run_agent import AIAgent +from hermes_cli.runtime_provider import resolve_runtime_provider +rt = resolve_runtime_provider(requested="nous", target_model="anthropic/claude-fable-5.1") +sid = f"f0ab_{arm}_{int(time.time())}" +ag = AIAgent(model="anthropic/claude-fable-5.1", provider="nous", base_url=rt.get("base_url"), api_key=rt.get("api_key"), + api_mode=rt.get("api_mode"), session_id=sid, quiet_mode=True, + enabled_toolsets=["file", "terminal"], platform="cli", max_iterations=12, + skip_context_files=True, skip_memory=True, reasoning_config={"enabled": True, "effort": "medium"}) +task = ("In /tmp/f0ab_work (create it), do these steps ONE tool call at a time, no parallel calls: " + "1) write a.txt with 'alpha', 2) write b.txt with 'beta', 3) read a.txt, 4) read b.txt, " + "5) run `ls -la /tmp/f0ab_work`, 6) run `wc -c /tmp/f0ab_work/*`, 7) run `cat /tmp/f0ab_work/a.txt`, " + "then reply with one line: DONE.") +t0 = time.time() +r = ag.run_conversation(task) +print("final:", (r.get("final_response") or "")[:80], "| wall", round(time.time() - t0, 1), "s") +time.sleep(1) +log = subprocess.run(f"grep -h '\\[{sid}\\]' ~/.hermes/logs/agent.log | grep 'API call #'", shell=True, capture_output=True).stdout.decode("utf-8", "replace") +rows = re.findall(r"API call #(\d+): .*in=(\d+) out=(\d+) .*cache=(\d+)/(\d+) \((\d+)%\)", log) +tot_in = tot_c = 0 +for n, i, o, c, ct, p in rows: + i, c = int(i), int(c); tot_in += i; tot_c += c + print(f" call {n:>2} in={i:>7} out={o:>5} cached={c:>7} ({p}%) uncached={i-c}") +print(f"ARM {arm}: calls={len(rows)} input={tot_in} cached={tot_c} uncached={tot_in-tot_c} hit={100*tot_c/max(tot_in,1):.1f}%") +print(json.dumps({"arm": arm, "calls": len(rows), "input": tot_in, "cached": tot_c})) diff --git a/evals/postmortem/live_ab/cache_prefix_wire.py b/evals/postmortem/live_ab/cache_prefix_wire.py new file mode 100644 index 0000000000..127350eab5 --- /dev/null +++ b/evals/postmortem/live_ab/cache_prefix_wire.py @@ -0,0 +1,79 @@ +"""Definitive F0 test: capture consecutive wire payloads on a real Fable 5.1 tool loop and diff the +message prefix between call N and N+1. If Hermes strips prior-turn thinking, call N+1's messages[:k] +will NOT equal call N's messages (prefix divergence) even though the conversation only grew. +Also reports cache hit per call. Cost: a handful of calls.""" +import os, sys, re, time, json, copy, subprocess +# LIVE: makes ~6 real calls to the configured provider (a few cents). Usage: +# python cache_prefix_wire.py [--hermes-home DIR] (default HERMES_HOME: the real one, for credentials) +sys.path.insert(0, sys.argv[1]) +os.environ.setdefault("HERMES_HOME", os.path.expanduser("~/.hermes")) +if "--hermes-home" in sys.argv: + os.environ["HERMES_HOME"] = sys.argv[sys.argv.index("--hermes-home") + 1] +arm = sys.argv[2] if len(sys.argv) > 2 else "A" +import agent.anthropic_message_convert as amc +if arm == "B": + def _keep_all(result, base_url, model): + for idx, m in amc._assistant_block_lists(result): + for b in m["content"]: + if amc._block_type(b) in amc._THINKING_TYPES: b.pop("cache_control", None) + m.pop("_thinking_signature_invalidated", None) + amc._manage_thinking_signatures = _keep_all +# Capture every outbound Anthropic-format payload +captured = [] +import agent.anthropic_adapter as ad +_orig_convert = None +for name in ("convert_messages_to_anthropic", "to_anthropic_messages", "convert_to_anthropic", "build_anthropic_messages"): + if hasattr(amc, name): _orig_convert = (name, getattr(amc, name)); break +if _orig_convert is None: + # find the public converter: the function that calls _manage_thinking_signatures at line ~703 + import inspect + for name, fn in inspect.getmembers(amc, inspect.isfunction): + try: + if "_manage_thinking_signatures(result" in inspect.getsource(fn): _orig_convert = (name, fn); break + except Exception: pass +name, fn = _orig_convert +def _wrapped(*a, **k): + out = fn(*a, **k) + try: + msgs = out[1] if isinstance(out, tuple) else out # (system, messages) + captured.append(copy.deepcopy(msgs)) + except Exception: pass + return out +setattr(amc, name, _wrapped) +if hasattr(ad, name): setattr(ad, name, _wrapped) +print("hooked converter:", name) +from run_agent import AIAgent +from hermes_cli.runtime_provider import resolve_runtime_provider +rt = resolve_runtime_provider(requested="nous", target_model="anthropic/claude-fable-5.1") +sid = f"f0wire_{arm}_{int(time.time())}" +ag = AIAgent(model="anthropic/claude-fable-5.1", provider="nous", base_url=rt.get("base_url"), api_key=rt.get("api_key"), + api_mode=rt.get("api_mode"), session_id=sid, quiet_mode=True, enabled_toolsets=["file", "terminal"], + platform="cli", max_iterations=10, skip_context_files=True, skip_memory=True, + reasoning_config={"enabled": True, "effort": "medium"}) +task = ("Work in /tmp/f0wire (create it). Before EACH tool call, think carefully for a moment about edge cases. " + "Steps, one tool call each: 1) write notes.md with a 5-line summary of what a Python context manager is; " + "2) read it back; 3) run `wc -l /tmp/f0wire/notes.md`; 4) append one more line to notes.md explaining __exit__ return values; " + "5) read it back; then reply DONE.") +ag.run_conversation(task) +time.sleep(1) +log = subprocess.run(f"grep -h '\\[{sid}\\]' ~/.hermes/logs/agent.log | grep 'API call #'", shell=True, capture_output=True).stdout.decode("utf-8", "replace") +rows = re.findall(r"API call #(\d+): .*in=(\d+) out=(\d+) .*cache=(\d+)/(\d+) \((\d+)%\)", log) +for r in rows: print(f" call {r[0]:>2} in={int(r[1]):>6} out={int(r[2]):>5} cached={int(r[3]):>6} ({r[5]}%) uncached={int(r[1])-int(r[3])}") +print(f"ARM {arm}: captured {len(captured)} payloads") +# prefix divergence check +def sig(m): # message signature: role + block types + text/thinking lengths + signature presence + c = m.get("content") + if isinstance(c, list): + return (m.get("role"), tuple((b.get("type"), len(str(b.get("text", b.get("thinking", b.get("input", ""))))), bool(b.get("signature"))) for b in c)) + return (m.get("role"), str(c)[:200]) +div = 0 +for i in range(1, len(captured)): + prev, cur = captured[i-1], captured[i] + k = len(prev) + same = [sig(a) == sig(b) for a, b in zip(prev, cur[:k])] + if not all(same): + j = same.index(False); div += 1 + print(f" payload {i}: prefix DIVERGED at message {j}/{k}: prev={sig(prev[j])} cur={sig(cur[j])}") +print(f"ARM {arm}: {div} of {len(captured)-1} consecutive payloads had a mutated prefix (thinking blocks in prev assistant msgs: " + f"{sum(1 for m in captured[-1] if m.get('role')=='assistant' and isinstance(m.get('content'),list) and any(b.get('type')=='thinking' for b in m['content']))} of " + f"{sum(1 for m in captured[-1] if m.get('role')=='assistant')} assistant msgs in final payload)") diff --git a/evals/postmortem/live_ab/goal_judge_wait.py b/evals/postmortem/live_ab/goal_judge_wait.py new file mode 100644 index 0000000000..bf55692390 --- /dev/null +++ b/evals/postmortem/live_ab/goal_judge_wait.py @@ -0,0 +1,19 @@ +"""Live judge A/B on the run's real "waiting" response shape. Usage: python judge_ab.py [n]""" +import os, sys +root = sys.argv[1]; n = int(sys.argv[2]) if len(sys.argv) > 2 else 3 +sys.path.insert(0, root) +os.environ.setdefault("HERMES_HOME", os.path.expanduser("~/.hermes")) # LIVE: real auxiliary judge calls (cents) +from hermes_cli import goals +goal = ("Simplify the hermes-agent codebase by >=30% LOC with zero behavior change, decomposing every god file, " + "through parallel subagent waves; integrate, run the full test suite, and open one PR.") +response = ("Round 3 status: 4 worker batches are still running (deleg_9f1a2b, deleg_77c0de, deleg_a1b2c3, deleg_0e9f88; " + "8 workers on adapters/gateway god files, median 38 min in). Round-2 integration merged clean; full suite " + "green on the integration branch. Nothing else is dispatchable until these return: their summaries decide " + "which files round 4 takes. Waiting on the batch-complete notifications; no action needed from me now.") +kw = {"active_delegations": 4} if "active_delegations" in goals.judge_goal.__code__.co_varnames else {} +verdicts = [] +for _ in range(n): + v, reason, pf, directive, tf = goals.judge_goal(goal, response, **kw) + verdicts.append((v, directive.get("seconds") if directive else None)) + print(v, directive, "|", reason[:90]) +print("SUMMARY", verdicts) diff --git a/evals/postmortem/live_ab/hardline_scanner_matrix.py b/evals/postmortem/live_ab/hardline_scanner_matrix.py new file mode 100644 index 0000000000..7e14e17dc5 --- /dev/null +++ b/evals/postmortem/live_ab/hardline_scanner_matrix.py @@ -0,0 +1,26 @@ +import sys; sys.path.insert(0, sys.argv[1] if len(sys.argv) > 1 else ".") # usage: python hardline_scanner_matrix.py +from tools.approval_detection import detect_hardline_command as d +cases = { + # the reviewer's witnesses + "newline-hidden reboot in quoted $(grep)": ('echo "$(grep -P \'safe\' /dev/null\nreboot)"', True), + "grep with backtick operand": ("grep -e `echo needle` file", False), + # the original class (must stay fixed) + "canonical sed/grep/cut": ('sed -n "$(grep -n X f | cut -d: -f1),+3p" f', False), + "grep -c inside $()": ('echo "$(grep -c x f)"', False), + # controls: other separators inside the substitution + "; hidden reboot": ('echo "$(grep x f; reboot)"', True), + "&& hidden reboot": ('echo "$(grep x f && reboot)"', True), + "| hidden shutdown": ('echo "$(grep x f | shutdown -h now)"', True), + "nested $( ) newline reboot": ('echo "$(echo $(grep x f)\nreboot)"', True), + "backtick newline reboot": ('echo "`grep x f\nreboot`"', True), + # data that must stay allowed + "reboot as quoted grep pattern": ("grep -F 'sudo reboot' notes.md", False), + "commit msg with reboot on a data line": ('git commit -m "fix\nsudo reboot handling"', False), +} +bad = 0 +for name, (cmd, want_block) in cases.items(): + got, why = d(cmd) + ok = got == want_block + bad += not ok + print(("OK " if ok else "FAIL") + f" {name}: blocked={got} ({why})") +print("ALL OK" if not bad else f"{bad} FAILURES") diff --git a/evals/postmortem/live_ab/nested_delegate_deadline.py b/evals/postmortem/live_ab/nested_delegate_deadline.py new file mode 100644 index 0000000000..f9629799ef --- /dev/null +++ b/evals/postmortem/live_ab/nested_delegate_deadline.py @@ -0,0 +1,45 @@ +"""Live A/B for the nested-delegate deadline. A depth-1 orchestrator child dispatches one leaf that runs a +~460 s task (longer than the 420 s sequential deadline). On main the orchestrator's delegate_task call returns +'timed out after 420.0s' and the leaf runs on as an orphan; on the branch the call blocks and returns the +real result. Uses glm-5.3 via Nous for cost. Deadline shortened via config to keep the run short.""" +import os, sys, json, time, re, subprocess +# Usage: python nested_delegate_deadline.py (run once per ref; LIVE: a couple of real child calls) +root = sys.argv[1]; arm = os.path.basename(os.path.normpath(root)) +sys.path.insert(0, root) +# Temp HERMES_HOME with the real auth + a config that shortens the generic sequential deadline to 40 s, so the +# run takes ~1.5 min instead of 8. The fix exempts delegate_task from this deadline entirely, so the shortened +# value is exactly what main will hit. +import shutil, tempfile, yaml +home = tempfile.mkdtemp(prefix="dl_home_"); os.environ["HERMES_HOME"] = home +real_home = os.environ.get("HERMES_HOME_SOURCE", os.path.expanduser("~/.hermes")) # credentials are copied from here into a temp home +shutil.copy(f"{real_home}/auth.json", f"{home}/auth.json") +cfg = yaml.safe_load(open(f"{real_home}/config.yaml", encoding="utf-8")) or {} +cfg.setdefault("timeouts", {}).setdefault("tools", {})["sequential_call"] = 40 +cfg.setdefault("delegation", {})["orchestrator_enabled"] = True +yaml.safe_dump(cfg, open(f"{home}/config.yaml", "w", encoding="utf-8")) +import agent.tool_executor as te +assert te.__file__.startswith(root) +from agent.deadline import resolve_timeout +print("effective sequential deadline:", resolve_timeout("tools.sequential_call", default=te._resolve_concurrent_tool_timeout())) +from run_agent import AIAgent +from hermes_cli.runtime_provider import resolve_runtime_provider +MODEL = "z-ai/glm-5.3-flash" +rt = resolve_runtime_provider(requested="nous", target_model=MODEL) +sid = f"dl_{arm}_{int(time.time())}" +ag = AIAgent(model=MODEL, provider="nous", base_url=rt.get("base_url"), api_key=rt.get("api_key"), api_mode=rt.get("api_mode"), + session_id=sid, quiet_mode=True, enabled_toolsets=["terminal", "delegation"], platform="cli", max_iterations=8, + skip_context_files=True, skip_memory=True) +# Make this agent a depth-1 orchestrator exactly as delegate_tool_child_run does for a real nested orchestrator: +# at depth>0 delegate_task runs synchronously inside the tool call, which is the path under the deadline. +ag._delegate_depth = 1 +ag._delegate_role = "orchestrator" +task = ("Use delegate_task exactly once (not background) with a single task whose goal is: " + "\"Run the shell command `sleep 75 && echo LEAF_DONE_MARKER` with the terminal tool (background=false is fine, it is under the tool timeout), " + "then reply with the exact text the command printed.\" " + "When delegate_task returns, reply with ONE line: RESULT: followed by the child's summary text verbatim (or the error text if it errored).") +t0 = time.time(); r = ag.run_conversation(task); wall = time.time() - t0 +final = (r.get("final_response") or "").strip() +print(f"ARM {arm}: wall={wall:.0f}s final={final[:300]!r}") +timed_out = "timed out" in final.lower() +got_marker = "LEAF_DONE_MARKER" in final +print(json.dumps({"arm": arm, "delegate_timed_out": timed_out, "orchestrator_got_leaf_result": got_marker, "wall_s": round(wall)})) diff --git a/evals/postmortem/live_ab/subagent_context_cap.py b/evals/postmortem/live_ab/subagent_context_cap.py new file mode 100644 index 0000000000..b18223d4ac --- /dev/null +++ b/evals/postmortem/live_ab/subagent_context_cap.py @@ -0,0 +1,32 @@ +"""Live: build a real child through delegate_tool's spawn path (real imports, temp HERMES_HOME) and read the +trigger it resolves on a 1M-window model. Run against main and the branch.""" +import os, sys, tempfile, shutil +root = sys.argv[1] +sys.path.insert(0, root) +home = tempfile.mkdtemp(prefix="hh-") +os.environ["HERMES_HOME"] = home +os.environ["HERMES_STREAM_RETRIES"] = "0" +try: + from run_agent import AIAgent + import tools.delegate_tool as dt + parent = AIAgent(api_key="k", base_url="https://example.com/v1", provider="test-provider", + model="anthropic/claude-fable-5.1", quiet_mode=True, skip_context_files=True, skip_memory=True) + # find the child-construction function by name + fn = [getattr(dt, n) for n in dir(dt) if n.startswith("_") and "child" in n.lower() and callable(getattr(dt, n)) and "spawn" in (getattr(dt, n).__doc__ or "").lower() + n.lower()] + import inspect + cands = [n for n, f in inspect.getmembers(dt, inspect.isfunction) if "AIAgent(" in inspect.getsource(f)] + print("constructor fn:", cands) + f = getattr(dt, cands[0]) + sig = inspect.signature(f); print("sig:", sig) + kwargs = {} + for name, p in sig.parameters.items(): + if name == "parent_agent": kwargs[name] = parent + elif name == "goal": kwargs[name] = "hi" + elif p.default is inspect._empty: kwargs[name] = None + child = f(**kwargs) + child = child[0] if isinstance(child, tuple) else child + cc = child.context_compressor + print(f"window={cc.context_length:,} threshold_percent={cc.threshold_percent} trigger={cc.threshold_tokens:,} cap={cc.threshold_tokens_cap}") + print(f"parent trigger={parent.context_compressor.threshold_tokens:,}") +finally: + shutil.rmtree(home, ignore_errors=True) diff --git a/evals/postmortem/review_probes/__init__.py b/evals/postmortem/review_probes/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/evals/postmortem/review_probes/cache_estimator_probe.py b/evals/postmortem/review_probes/cache_estimator_probe.py new file mode 100644 index 0000000000..47ecd61881 --- /dev/null +++ b/evals/postmortem/review_probes/cache_estimator_probe.py @@ -0,0 +1,38 @@ +"""#103476: preflight estimate vs Anthropic wire size with retained thinking + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import os,sys,tempfile,copy,json,types,subprocess +from pathlib import Path +sys.path.insert(0,sys.argv[1] if len(sys.argv)>1 else os.getcwd());os.environ['HERMES_HOME']=tempfile.mkdtemp(prefix='cache-boundary-') # usage: +from agent.anthropic_message_convert import convert_messages_to_anthropic +from agent.context_compressor import ContextCompressor +from agent.turn_context import _preflight_request_tokens +from agent.model_metadata import estimate_messages_tokens_rough +model='anthropic/claude-fable-5.1';url='https://inference-api.nousresearch.com/v1' +cc=ContextCompressor(model=model,provider='nous',base_url=url,api_mode='anthropic_messages',config_context_length=1000000,threshold_tokens_cap=200000) +rows=[{'role':'user','content':'Investigate repository.'}] +for i in range(48): + text='reasoning text ' * 1800 + rows.extend([{'role':'assistant','content':'tool','reasoning':text,'reasoning_details':[{'type':'thinking','thinking':text,'signature':f'signature-{i}'}],'tool_calls':[{'id':f't{i}','type':'function','function':{'name':'terminal','arguments':'{}'}}]},{'role':'tool','tool_call_id':f't{i}','content':f'Result {i}'}]) +a=types.SimpleNamespace(api_mode='anthropic_messages',provider='nous',model=model,base_url=url,tools=[],_usage_anchor=None) +pre=_preflight_request_tokens(a,rows,'');wire=convert_messages_to_anthropic(rows,base_url=url,model=model)[1];wire_est=estimate_messages_tokens_rough(wire) +print(json.dumps({'head':subprocess.check_output(['git','rev-parse','HEAD'],text=True, encoding='utf-8', errors='replace').strip(),'preflight':pre,'wire_estimate':wire_est,'preflight_should_compress':cc.should_compress(pre),'wire_estimate_should_compress':cc.should_compress(wire_est)})) +# Local signature-prefix checker validates the documented contract only; it is not the provider. +def prefixes(messages): + before=[]; result={} + for msg in messages: + nonthinking=[] + for b in msg['content'] if isinstance(msg['content'],list) else [{'type':'text','text':msg['content']}]: + if b.get('type')=='thinking':result[b['signature']]=json.dumps(before,sort_keys=True) + else:nonthinking.append(b) + before.append({'role':msg['role'],'content':nonthinking}) + return result +before=prefixes(wire) +cc._generate_summary=lambda *a,**k:'Task: investigate repository. Continue reviewing remaining files.' +out=cc.compress(copy.deepcopy(rows),current_tokens=wire_est,force=True) +after=prefixes(convert_messages_to_anthropic(out,base_url=url,model=model)[1]) +print('retained',list(after),'prefix_changed',[k for k,v in after.items() if before.get(k)!=v]) +print('retained with unchanged system; actual agent rebuilds system too, so this is lower bound') diff --git a/evals/postmortem/review_probes/context_cap_probe.py b/evals/postmortem/review_probes/context_cap_probe.py new file mode 100644 index 0000000000..57b3300f6a --- /dev/null +++ b/evals/postmortem/review_probes/context_cap_probe.py @@ -0,0 +1,95 @@ +"""#103513: child cap through repeated compression and persistence + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import os, sys, tempfile, json, socket, copy +from pathlib import Path +from types import SimpleNamespace +root, tag = sys.argv[1:3] +sys.path.insert(0, root) +os.chdir(root) +for k in list(os.environ): + if any(s in k for s in ('API_KEY','TOKEN','SECRET')) or k.startswith('HERMES_'): + os.environ.pop(k, None) +home = tempfile.mkdtemp(prefix='cap-review-') +os.environ['HERMES_HOME'] = home +os.environ['HERMES_DISABLE_REDACTION'] = 'true' +import yaml +cfg = {'model': {'default': 'anthropic/claude-fable-5.1', 'provider':'openai-compat', 'base_url':'http://127.0.0.1:1/v1', 'context_length':1000000}, 'compression':{'threshold':0.85}, 'delegation': {}} +if len(sys.argv)>3: + cfg['delegation']['compression_threshold_tokens'] = json.loads(sys.argv[3]) +Path(home,'config.yaml').write_text(yaml.safe_dump(cfg), encoding='utf-8') +def blocked(*a, **kw): + raise RuntimeError('Network disabled in unpaid cap probe') +socket.socket.connect = blocked +socket.create_connection = blocked +from run_agent import AIAgent +import tools.delegate_tool as dt +import agent.context_compressor as mod +from agent.model_metadata import estimate_messages_tokens_rough +from hermes_state import SessionDB +from unittest.mock import patch +print('IDENTITY',json.dumps({'tag':tag,'tree':root,'delegate':dt.__file__,'compressor':mod.__file__,'cap_present':hasattr(dt,'_apply_child_compression_cap'),'home':home}),flush=True) +db=SessionDB(Path(home,'state.db')) +parent=AIAgent(api_key='test-key',base_url='http://127.0.0.1:1/v1',provider='openai-compat',model='anthropic/claude-fable-5.1',enabled_toolsets=[],quiet_mode=True,skip_context_files=True,skip_memory=True,save_trajectories=False,session_db=db) +child=dt._build_child_agent(task_index=0,goal='Review compression cap; continue the task.',context=None,toolsets=[],model=None,max_iterations=10,task_count=1,parent_agent=parent) +cc=child.context_compressor +print('SPAWN',json.dumps({'parent':parent.context_compressor.threshold_tokens,'child':cc.threshold_tokens,'cap':cc.threshold_tokens_cap,'tail':cc.tail_token_budget,'enabled':child.compression_enabled}),flush=True) +if len(sys.argv)>3: + child.close(); parent.close(); db.close(); sys.exit(0) +# Real trigger and real compressor transitions, no hand-written compression behavior. +base=[{'role':'user','content':'Continue reviewing this project and preserve the current goal.'}] +for i in range(48): + base.append({'role':'assistant' if i%2==0 else 'user','content':f'fixture {i} '+('A documented implementation detail with evidence and constraints. '*360)}) +base.append({'role':'user','content':'Continue the review and report findings.'}) +summary_calls=[] +def local_summary(**kw): + assert kw['task']=='compression' + summary_calls.append(len(kw['messages'][0]['content'])) + return SimpleNamespace(choices=[SimpleNamespace(message=SimpleNamespace(content='## Current task\nContinue reviewing the project. Preserve evidence and report findings.\n## Decisions\nNo changes have been applied.\n## Next steps\nInspect remaining details.',reasoning=None,reasoning_content=None),finish_reason='stop')],usage=None) +messages=base +for cycle in range(4): + tokens=estimate_messages_tokens_rough(messages) + cc.update_from_response({'prompt_tokens': tokens,'completion_tokens':0}) + should=cc.should_compress(tokens) + before=len(messages) + if should: + with patch.object(mod,'call_llm',side_effect=local_summary): + messages=cc.compress(messages,current_tokens=tokens) + after=estimate_messages_tokens_rough(messages) + print('CYCLE',json.dumps({'cycle':cycle,'before_tokens':tokens,'trigger':cc.threshold_tokens,'should':should,'before_messages':before,'after_tokens':after,'after_messages':len(messages),'compressions':cc.compression_count,'summary_calls':len(summary_calls),'blocked':cc.should_compress_info(tokens),'cap':cc.threshold_tokens_cap}),flush=True) + cc.update_from_response({'prompt_tokens':after,'completion_tokens':0}) + if cycle < 3: + messages=messages+copy.deepcopy(base[1:]) +print('SUMMARY_INPUT_CHARS',summary_calls,flush=True) +# Model switching must keep the cap, including ratio-lower small windows. +for window in [128000,1000000]: + cc.update_model(child.model,window,provider=child.provider,base_url=child.base_url,api_mode=child.api_mode) + print('MODEL_SWITCH',window,cc.threshold_tokens,cc.threshold_tokens_cap,flush=True) +if hasattr(dt,'_apply_child_compression_cap'): + for raw in ['200k','default',True,False,0,None,200000.9,float('inf')]: + temp=SimpleNamespace(context_compressor=mod.ContextCompressor(model=child.model,threshold_percent=.85,config_context_length=1000000,quiet_mode=True)) + try: + dt._apply_child_compression_cap(temp,{'compression_threshold_tokens':raw}) + print('CONFIG',repr(raw),temp.context_compressor.threshold_tokens,temp.context_compressor.threshold_tokens_cap,flush=True) + except Exception as exc: + print('CONFIG_EXCEPTION',repr(raw),type(exc).__name__,str(exc),flush=True) +# Accounting policy: the cap changes the threshold, not what route-aware pressure counts. +from agent.turn_context import _preflight_request_tokens, _agent_stale_thinking_on_wire +history=[{'role':'user','content':'review'}] +for i in range(24): + history += [{'role':'assistant','content':'observed','reasoning_content':'reasoning detail '*4000},{'role':'user','content':'continue'}] +for provider, model in [('openai-compat','test-model'),('deepseek','deepseek-chat')]: + route=SimpleNamespace(provider=provider,model=model,base_url='http://127.0.0.1:1/v1',api_mode='chat_completions',tools=[]) + pressure=_preflight_request_tokens(route,history,'') + print('ACCOUNTING',provider,_agent_stale_thinking_on_wire(route),pressure,cc.should_compress(pressure),flush=True) +# Already-cached legacy tail is an explicit compatibility edge. +if hasattr(dt,'_apply_child_compression_cap'): + legacy=SimpleNamespace(context_compressor=mod.ContextCompressor(model=child.model,threshold_percent=.85,config_context_length=1000000,tail_mode='legacy',quiet_mode=True)) + old_tail=legacy.context_compressor.tail_token_budget + dt._apply_child_compression_cap(legacy,{}) + print('LEGACY_RESOLVED',old_tail,legacy.context_compressor.tail_token_budget,legacy.context_compressor.threshold_tokens,flush=True) +child.close();parent.close();db.close() +print('DONE',tag,flush=True) diff --git a/evals/postmortem/review_probes/credential_identity_probe.py b/evals/postmortem/review_probes/credential_identity_probe.py new file mode 100644 index 0000000000..4dd63ddfc9 --- /dev/null +++ b/evals/postmortem/review_probes/credential_identity_probe.py @@ -0,0 +1,113 @@ +"""#103526: explicit account-A key must not be replaced by the singleton's account-B key (real SDK -> loopback) + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import os, tempfile, sys, json, time, base64, threading, importlib.util +from pathlib import Path +from datetime import datetime, timezone +from http.server import ThreadingHTTPServer, BaseHTTPRequestHandler +ROOT = Path(sys.argv[1]); MODE = sys.argv[2] +sys.path.insert(0, str(ROOT)) +home = Path(tempfile.mkdtemp(prefix='pr103526-probe-')) +os.environ.update(HOME=str(home), HERMES_HOME=str(home/'hermes'), XDG_CONFIG_HOME=str(home/'config'), CODEX_HOME=str(home/'codex')) +for k in list(os.environ): + if any(x in k for x in ('API_KEY','TOKEN','NOUS_','SECRET')): os.environ.pop(k, None) +(home/'hermes').mkdir() +(home/'hermes'/'config.yaml').write_text('nous:\n keepalive_interval_seconds: 0\nmemory:\n memory_enabled: false\n', encoding='utf-8') +def guard(event,args): + if event == 'socket.connect' and isinstance(args[1], tuple) and args[1][0] not in ('127.0.0.1','::1'): + raise RuntimeError('External network forbidden by review probe') +sys.addaudithook(guard) +def jwt(sub, ttl): + def part(v): return base64.urlsafe_b64encode(json.dumps(v).encode()).decode().rstrip('=') + return part({'alg':'none'})+'.'+part({'sub':sub,'scope':'inference:invoke','exp':int(time.time()+ttl)})+'.sig' +def claims(token): + return json.loads(base64.urlsafe_b64decode(token.split('.')[1]+'===')) +records=[] +class Handler(BaseHTTPRequestHandler): + def do_POST(self): + raw_body=self.rfile.read(int(self.headers.get('Content-Length',0))) + if self.path == '/api/oauth/token': + records.append({'path':self.path,'refresh':True}) + raw=json.dumps({'access_token':refresh_reply,'refresh_token':'fixture-rotated','expires_in':3600,'token_type':'Bearer','scope':'inference:invoke'}).encode() + self.send_response(200);self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(raw)));self.end_headers();self.wfile.write(raw);return + body=json.loads(raw_body or '{}') + bearer=self.headers.get('Authorization','').removeprefix('Bearer ') + records.append({'path':self.path,'sub':claims(bearer).get('sub') if bearer else None}) + if bearer and claims(bearer)['exp'] < time.time(): + raw=json.dumps({'error':{'message':'expired bearer','type':'authentication_error'}}).encode();self.send_response(401);self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(raw)));self.end_headers();self.wfile.write(raw);return + data={'id':'local-probe','object':'chat.completion','created':int(time.time()),'model':'hermes-test','choices':[{'index':0,'message':{'role':'assistant','content':'local-only'},'finish_reason':'stop'}],'usage':{'prompt_tokens':1,'completion_tokens':1,'total_tokens':2}} + raw=json.dumps(data).encode(); self.send_response(200);self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(raw)));self.end_headers();self.wfile.write(raw) + def log_message(self,*a): pass +server=ThreadingHTTPServer(('127.0.0.1',0),Handler); threading.Thread(target=server.serve_forever,daemon=True).start() +url=f'http://127.0.0.1:{server.server_port}/v1' +# Runtime override preserves loopback routing, without relaxing URL validation. +os.environ['NOUS_INFERENCE_BASE_URL']=url +os.environ['HERMES_SHARED_AUTH_DIR']=str(home/'shared') +from run_agent import AIAgent +from agent.turn_iteration_prep import prepare_iteration +import agent.client_lifecycle as lifecycle +import hermes_cli.auth as auth +print(json.dumps({'mode':MODE,'module':lifecycle.__file__,'has_new':hasattr(AIAgent,'_adopt_nous_key_before_expiry'),'home':str(home)}),flush=True) +if MODE=='main': + spec=importlib.util.spec_from_file_location('main_prep',Path(__file__).with_name('main-turn_iteration_prep.py')); mod=importlib.util.module_from_spec(spec);sys.modules[spec.name]=mod;spec.loader.exec_module(mod);prepare_iteration=mod.prepare_iteration + +def store(token): + exp=claims(token)['exp']; state={'portal_base_url':'https://portal.nousresearch.com','inference_base_url':'https://inference-api.nousresearch.com/v1','client_id':'hermes-cli','token_type':'Bearer','scope':'inference:invoke','access_token':token,'refresh_token':'fixture-refresh-never-send','expires_at':datetime.fromtimestamp(exp,timezone.utc).isoformat(),'expires_in':3600,'agent_key':token,'agent_key_expires_at':datetime.fromtimestamp(exp,timezone.utc).isoformat()} + (home/'hermes'/'auth.json').write_text(json.dumps({'version':1,'active_provider':'nous','providers':{'nous':state}}), encoding='utf-8') +results=[] +for case, own_sub, store_sub, ttl in [('same-account','account-A','account-A',30),('explicit-account','account-A','account-B',30),('far-from-expiry','account-A','account-B',3000)]: + own=jwt(own_sub,ttl);fresh=jwt(store_sub,3600);store(fresh) + agent=AIAgent(api_key=own,base_url=url,provider='nous',model='hermes-test',quiet_mode=True,skip_context_files=True,skip_memory=True,enabled_toolsets=[]) + messages=[{'role':'user','content':'local fixture'}]; before=json.dumps(messages) + prepare_iteration(agent,messages=messages,api_call_count=0) + client=agent._create_request_openai_client(reason='review_probe') + reply=client.chat.completions.create(model='hermes-test',messages=messages) + result={'case':case,'before_sub':own_sub,'after_sub':claims(agent.api_key)['sub'],'adopted_fresh':agent.api_key==fresh,'messages_unchanged':json.dumps(messages)==before,'wire':records[-1],'reply':reply.choices[0].message.content} + results.append(result); print(json.dumps(result),flush=True) + agent._close_request_openai_client(client,reason='probe_done') + agent.client.close() +# Contended peer-adoption: real auth-store locking and SDK wire, no resolver mocks. +from concurrent.futures import ThreadPoolExecutor +fresh=jwt('account-A',3600); store(fresh) +expired=jwt('account-A',-30) +barrier=threading.Barrier(12) +def worker(_): + a=AIAgent(api_key=expired,base_url=url,provider='nous',model='hermes-test',quiet_mode=True,skip_context_files=True,skip_memory=True,enabled_toolsets=[]) + barrier.wait(timeout=30) + messages=[{'role':'user','content':'concurrent local fixture'}] + prepare_iteration(a,messages=messages,api_call_count=0) + c=a._create_request_openai_client(reason='concurrent_review_probe') + try: + c.chat.completions.create(model='hermes-test',messages=messages) + status=200 + except Exception as e: + status=getattr(e,'status_code',type(e).__name__) + finally: + a._close_request_openai_client(c,reason='probe_done'); a.client.close() + return status +with ThreadPoolExecutor(max_workers=12) as executor: + statuses=list(executor.map(worker,range(12))) +print(json.dumps({'concurrent_statuses':statuses,'successes':statuses.count(200),'401s':statuses.count(401)}),flush=True) +# No fresh peer key: twelve agents contend for one real local refresh POST. +refresh_reply=jwt('account-A',3600) +if MODE != 'main': + import shutil + shutil.rmtree(home/'shared',ignore_errors=True) + expired=jwt('account-A',30); store(expired) + state_file=home/'hermes'/'auth.json' + state=json.loads(state_file.read_text(encoding='utf-8')); state['providers']['nous']['portal_base_url']=url.removesuffix('/v1');state_file.write_text(json.dumps(state), encoding='utf-8') + barrier=threading.Barrier(12);records.clear() + with ThreadPoolExecutor(max_workers=12) as executor: + mint_statuses=list(executor.map(worker,range(12))) + posts=sum(r.get('refresh',False) for r in records) + print(json.dumps({'mint_statuses':mint_statuses,'refresh_posts':posts}),flush=True) + assert posts == 1 and mint_statuses == [200]*12 +server.shutdown() +(Path(os.environ.get('PROBE_OUT', tempfile.gettempdir()))/('cred-probe-'+MODE+'.json')).write_text(json.dumps({'cases':results,'concurrent_statuses':statuses},indent=2), encoding='utf-8') +assert all(r['messages_unchanged'] for r in results) +# Regression contract (fixed head): the explicit account-A key must stay A; same-account adoption must still happen. +assert results[1]['after_sub']=='account-A', 'explicit-account key was replaced by the singleton account (the round-1 defect)' +assert results[0]['after_sub']=='account-A' and results[0]['adopted_fresh'], 'same-account fresh key must still be adopted' diff --git a/evals/postmortem/review_probes/deadline_probe.py b/evals/postmortem/review_probes/deadline_probe.py new file mode 100644 index 0000000000..630795b7ac --- /dev/null +++ b/evals/postmortem/review_probes/deadline_probe.py @@ -0,0 +1,140 @@ +"""#103486: nested delegate deadline through actual dispatch + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import os, sys, tempfile, pathlib, json, time, threading, socket, subprocess +from types import SimpleNamespace as NS +root = pathlib.Path(sys.argv[1]).resolve() +sys.path.insert(0, str(root)); os.chdir(root) +for k in list(os.environ): + if k.startswith('HERMES_') or k.endswith(('_API_KEY', '_TOKEN')): + os.environ.pop(k, None) +home = tempfile.TemporaryDirectory(prefix='deadline-probe-') +os.environ['HERMES_HOME'] = home.name +os.environ['HERMES_DISABLE_TELEMETRY'] = '1' +pathlib.Path(home.name, 'config.yaml').write_text('timeouts:\n tools:\n sequential_call: 0.3\ndelegation:\n max_summary_chars: 24000\n', encoding='utf-8') +# Any accidental provider, metadata, or telemetry request fails closed. +def no_connect(*args, **kwargs): + raise RuntimeError('OFFLINE PROBE: network forbidden') +socket.socket.connect = no_connect +socket.create_connection = no_connect +import agent.tool_executor as te +import agent.turn_usage as tu +import tools.delegate_tool_results as dr +import agent.usage_pricing as up +import agent.codex_runtime as cr +from run_agent import AIAgent +for mod in (te, tu, dr, up, cr): + assert pathlib.Path(mod.__file__).is_relative_to(root), mod.__file__ +print(json.dumps({'sha': subprocess.check_output(['git','rev-parse','HEAD'],text=True, encoding='utf-8', errors='replace').strip(), 'modules': {m.__name__:m.__file__ for m in (te,tu,dr,up,cr)}, 'home':home.name}), flush=True) +a = AIAgent(api_key='offline-fixture', base_url='http://127.0.0.1:9/v1', provider='openai-compat', model='offline-test', enabled_toolsets=[], quiet_mode=True, skip_context_files=True, skip_memory=True, save_trajectories=False) +a._delegate_depth = 1; a._delegate_role = 'orchestrator' +assert te._resolve_sequential_tool_timeout() == 0.3 +rows=[] +for name in ('delegate_task','terminal','execute_code'): + entered=threading.Event(); finished=threading.Event() + def work(args): + entered.set() + time.sleep(0.65) + finished.set() + return json.dumps({'marker':'LEAF_DONE_MARKER'}) + start=time.monotonic() + result=te._run_sequential_tool_execution_middleware(a,function_name=name,function_args={},effective_task_id='offline',tool_call_id='probe-'+name,execute=work) + elapsed=time.monotonic()-start + assert entered.is_set(), 'tool body never entered' + row={'case':name,'elapsed':round(elapsed,4),'result_type':type(result.result).__name__,'result':str(result.result),'completed_on_return':finished.is_set()} + assert finished.wait(2), 'fixture did not finish' + rows.append(row) +print(json.dumps({'tool_boundary':rows}),flush=True) +# Warm real middleware, then interrupt a non-cooperative delegated operation. +a._interrupt_requested=False +entered=threading.Event(); release=threading.Event() +def blocked(args): + entered.set(); release.wait(8); return 'late' +def interrupt(): + assert entered.wait(2) + time.sleep(0.1) + a.interrupt('offline cancellation probe') +t = threading.Thread(target=interrupt); t.start() +start=time.monotonic() +r=te._run_sequential_tool_execution_middleware(a,function_name='delegate_task',function_args={},effective_task_id='offline',tool_call_id='probe-interrupt',execute=blocked) +elapsed=time.monotonic()-start +release.set();t.join() +print(json.dumps({'interrupt':{'elapsed':round(elapsed,4),'result_type':type(r.result).__name__,'result':str(r.result)}}),flush=True) +# Feed provider-shaped data to the production writer, not fabricated _last_turn_usage. +# Price lookup is unrelated and potentially networked; only pricing is stubbed. +tu.estimate_usage_cost=lambda *args,**kwargs: NS(amount_usd=None,status='unknown',source='offline') +up.estimate_usage_cost=tu.estimate_usage_cost +class Compressor: + context_length=200000 + max_tokens=8000 + threshold_tokens=190000 + def update_from_response(self, usage): + self.last_prompt_tokens=usage.get('prompt_tokens',-1) +def parent(provider,mode='chat_completions',client=None): + b=NS(context_compressor=Compressor(),provider=provider,api_mode=mode,model='offline',base_url='',client=client,session_id='offline-usage',_session_db=None,quiet_mode=True,verbose_logging=False) + for key in ('api_calls','prompt_tokens','completion_tokens','total_tokens','input_tokens','output_tokens','cache_read_tokens','cache_write_tokens','reasoning_tokens','estimated_cost_usd'): + setattr(b,'session_'+key,0) + return b +def record(b,raw): + tu.record_response_usage(b,NS(usage=raw),messages=[{'role':'user','content':'fixture'}],api_call_count=1,api_duration=0,compression_attempts=0,max_compression_attempts=3) + return {'prompt':b._last_turn_usage['prompt_tokens'],'input':b._last_turn_usage['input_tokens'],'cache_read':b._last_turn_usage['cache_read_tokens'],'cache_write':b._last_turn_usage['cache_write_tokens'],'budget':dr._parent_summary_char_budget(b,1)} +chat={'prompt_tokens':30000,'completion_tokens':100,'prompt_tokens_details':{'cached_tokens':20000,'cache_write_tokens':5000}} +anth={'input_tokens':5000,'output_tokens':100,'cache_read_input_tokens':20000,'cache_creation_input_tokens':5000} +responses={'input_tokens':30000,'output_tokens':100,'input_tokens_details':{'cached_tokens':20000,'cache_write_tokens':5000}} +cases=[('openai', 'chat_completions',chat),('nous','chat_completions',chat),('openrouter','chat_completions',chat),('anthropic','anthropic_messages',anth),('minimax','anthropic_messages',anth),('minimax-cn','anthropic_messages',anth),('bedrock','chat_completions',{'prompt_tokens':30000,'completion_tokens':100,'cache_read_input_tokens':20000,'cache_creation_input_tokens':5000}),('google','chat_completions',chat),('deepseek','chat_completions',{'prompt_tokens':30000,'completion_tokens':100,'prompt_cache_hit_tokens':20000}),('moonshot','chat_completions',{'prompt_tokens':30000,'completion_tokens':100,'cached_tokens':20000}),('openai-codex','codex_responses',responses),('openai-compat','chat_completions',{'input_tokens':30000,'output_tokens':100})] +usage_rows=[] +for provider,mode,raw in cases: + b=parent(provider,mode) + fresh=record(b,raw) + assert fresh['prompt']==30000 + b.session_prompt_tokens=25000000 + long_budget=dr._parent_summary_char_budget(b,1) + usage_rows.append({'provider':provider,'mode':mode,**fresh,'long_lived_budget':long_budget}) +# Raw native adapter response -> production conversion -> writer -> budget. +from agent.bedrock_adapter import normalize_converse_response +from agent.gemini_native_adapter import _usage_from_metadata +native=[] +for provider, raw in [('bedrock',normalize_converse_response({'usage':{'inputTokens':5000,'cacheReadInputTokens':20000,'cacheWriteInputTokens':5000,'outputTokens':100}}).usage),('google',_usage_from_metadata({'promptTokenCount':30000,'cachedContentTokenCount':25000,'candidatesTokenCount':100,'totalTokenCount':30100}))]: + b=parent(provider); values=record(b,raw); assert values['prompt']==30000 + native.append({'provider':provider,**values}) +print(json.dumps({'provider_writers':usage_rows,'native_adapters':native}),flush=True) +# Exercise the actual MoA accounting deposit/consume path without model calls. +from agent.moa_loop import MoAClient +moa_client = MoAClient('offline') +moa_client.chat.completions._fold_pending_accounting(up.CanonicalUsage(input_tokens=240000), None) +b=parent('moa',client=moa_client) +moa=record(b,chat) +results=[{'task_index':0,'summary':'X'*5000}] +dr._apply_summary_budget(results,b) +print(json.dumps({'moa':{**moa,'actual_aggregator_prompt':30000,'anchor_prompt':b._usage_anchor.prompt_tokens if hasattr(b._usage_anchor,'prompt_tokens') else repr(b._usage_anchor),'truncated':results[0].get('summary_truncated',False)}}),flush=True) +b=parent('openai-codex','codex_app_server') +b._last_turn_usage=None +codex=cr._record_codex_app_server_usage(b,NS(token_usage_last={'inputTokens':170000,'cachedInputTokens':20000,'outputTokens':100},model_context_window=200000)) +print(json.dumps({'codex_app_server':{'returned_prompt':codex['prompt_tokens'],'last_turn_usage':b._last_turn_usage,'compressor_prompt':b.context_compressor.last_prompt_tokens,'budget':dr._parent_summary_char_budget(b,1)}}),flush=True) +# Missing usage and full-context defaults, including a current turn with no usage. +b=parent('openai-compat'); record(b,{'prompt_tokens':190000,'completion_tokens':100}) +print(json.dumps({'near_full_budget':dr._parent_summary_char_budget(b,1),'near_full_batch_budget':dr._parent_summary_char_budget(b,5)}),flush=True) +# A successful provider response without usage keeps no current usage after turn reset. +b._last_turn_usage=None +record_outcome=tu.record_response_usage(b,NS(usage=None),messages=[],api_call_count=1,api_duration=0,compression_attempts=0,max_compression_attempts=3) +missing=[{'task_index':0,'summary':'X'*20000}] +dr._apply_summary_budget(missing,b) +print(json.dumps({'missing_usage_next_turn':{'session_prompt':b.session_prompt_tokens,'compressor_prompt':b.context_compressor.last_prompt_tokens,'last_turn_usage':b._last_turn_usage,'budget':dr._parent_summary_char_budget(b,1),'summary_length':len(missing[0]['summary']),'truncated':missing[0].get('summary_truncated',False)}}),flush=True) +# Full sequential executor -> resolver -> real middleware -> result commit. +# Replace only the model/delegation body; no subagents or model requests. +a._interrupt_requested=False +full_entered=threading.Event() +def full_body(args): + full_entered.set(); time.sleep(0.65) + return json.dumps({'marker':'FULL_BOUNDARY_DONE'}) +a._dispatch_delegate_task=full_body +calls=NS(tool_calls=[NS(id='full-boundary',type='function',function=NS(name='delegate_task',arguments='{}'))]) +messages=[] +start=time.monotonic() +te.execute_tool_calls_sequential(a,calls,messages,'offline-full') +assert full_entered.is_set() +print(json.dumps({'full_sequential_boundary':{'elapsed':round(time.monotonic()-start,4),'messages':messages}}),flush=True) +print('PROBE_COMPLETE',flush=True) diff --git a/evals/postmortem/review_probes/finalizer_schedule_probe.py b/evals/postmortem/review_probes/finalizer_schedule_probe.py new file mode 100644 index 0000000000..596309560f --- /dev/null +++ b/evals/postmortem/review_probes/finalizer_schedule_probe.py @@ -0,0 +1,56 @@ +"""#103507: pytest plugin forcing the consumer-first schedule (-p finalizer_schedule_probe --finalizer-probe=consumer-first) + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import asyncio +import json +import pytest + + +def pytest_addoption(parser): + parser.addoption('--finalizer-probe', default='off') + + +@pytest.fixture(autouse=True) +def finalizer_schedule_probe(request, monkeypatch): + mode = request.config.getoption('--finalizer-probe') + if mode == 'off': + yield + return + from agent import relay_llm, chat_completion_helpers + print('PROBE_IMPORT', relay_llm.__file__, chat_completion_helpers.__file__) + original_provider = relay_llm.ManagedLlmStream._provider_stream + original_count = chat_completion_helpers._StreamingCall._count_chunk + stats = {'gated_terminal_chunks': 0, 'consumer_releases': 0} + + def terminal(chunk): + choices = chunk.get('choices') or [] + return (not choices and chunk.get('usage') is not None) or ( + bool(choices) and choices[0].get('finish_reason') == 'tool_calls') + + async def gated_provider(stream, *args): + async for chunk in original_provider(stream, *args): + if terminal(chunk): + stream._review_gate = asyncio.Event() + stats['gated_terminal_chunks'] += 1 + yield chunk + await asyncio.wait_for(stream._review_gate.wait(), timeout=10) + else: + yield chunk + + def count_and_release(call, diag, chunk): + stream = call.managed_stream_holder.get('stream') + gate = getattr(stream, '_review_gate', None) + raw = chunk.model_dump() if hasattr(chunk, 'model_dump') else vars(chunk) + if terminal(raw) and gate is not None and not gate.is_set(): + gate.set() + stats['consumer_releases'] += 1 + return original_count(call, diag, chunk) + + assert mode == 'consumer-first' + monkeypatch.setattr(relay_llm.ManagedLlmStream, '_provider_stream', gated_provider) + monkeypatch.setattr(chat_completion_helpers._StreamingCall, '_count_chunk', count_and_release) + yield + print('PROBE_STATS', json.dumps(stats, sort_keys=True)) diff --git a/evals/postmortem/review_probes/goal_repaste_probe.py b/evals/postmortem/review_probes/goal_repaste_probe.py new file mode 100644 index 0000000000..6f6c8c852a --- /dev/null +++ b/evals/postmortem/review_probes/goal_repaste_probe.py @@ -0,0 +1,129 @@ +"""#103553: /goal re-paste vs option-selecting fragment + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted here). +It reproduced a defect in the first version of the PR; the fixed head must pass it. Paths are taken +from the command line / environment, never hard-coded. Usage: see the argument parsing at the top of the file. +""" +import asyncio +import contextlib +import copy +import io +import json +import os +from pathlib import Path +import queue +import socket +import sqlite3 +import sys +import tempfile +import threading +from unittest.mock import patch + +repo, tag = sys.argv[1:3] +sys.path.insert(0, repo) +# Remove all inherited Hermes/config and credential env before real imports. +for key in list(os.environ): + if key.startswith('HERMES_') or key.endswith(('_API_KEY', '_TOKEN')): + os.environ.pop(key, None) +home = tempfile.TemporaryDirectory(prefix='review-goaldup-') +os.environ['HERMES_HOME'] = home.name +os.environ['NO_PROXY'] = '*' +os.environ['TZ'] = 'UTC' +# Fail closed: these probes must never invoke a provider or external network. +socket.socket.connect = lambda *a, **k: (_ for _ in ()).throw(RuntimeError('network prohibited in review probe')) +from cli import HermesCLI +from hermes_cli import cli_commands_mixin, goals +from gateway.slash_commands_goals import GatewayGoalCommandsMixin +from gateway.config import Platform +from gateway.platforms.base import MessageEvent, MessageType +from gateway.session import SessionSource +from tui_gateway import server + +assert str(Path(cli_commands_mixin.__file__).resolve()).startswith(repo) +assert str(Path(goals.__file__).resolve()).startswith(repo) +server._hermes_home = Path(home.name) +goals._DB_CACHE.clear() +goals._get_session_db() +source = sqlite3.connect('file:/tmp/rf/state_copy.db?mode=ro', uri=True) +source.row_factory = sqlite3.Row +original = source.execute('SELECT content FROM messages WHERE id=264820').fetchone()['content'] +repeated = source.execute('SELECT content FROM messages WHERE id=267045').fetchone()['content'] +assert original == repeated +history = [{'role': 'user', 'content': original}, {'role': 'assistant', 'content': 'The wave is running; I will pick this up when the workers report back.'}] +output = {'tag': tag, 'modules': [cli_commands_mixin.__file__, goals.__file__, server.__file__], 'source_equal': original == repeated, 'original_chars': len(original)} + +def make_cli(hist, sid): + c = HermesCLI.__new__(HermesCLI) + c.session_id = sid + c.agent = None + c.conversation_history = copy.deepcopy(hist) + c._pending_input = queue.Queue() + return c + +def cli_case(hist, goal, sid): + c = make_cli(hist, sid) + before = copy.deepcopy(c.conversation_history) + with contextlib.redirect_stdout(io.StringIO()): + assert c.process_command('/goal ' + goal) + prompt = c._pending_input.get_nowait() + state = goals.GoalManager(sid).state + assert c.conversation_history == before + assert c._pending_input.empty() + assert state.goal == goal + return {'prompt': prompt, 'prompt_chars': len(prompt), 'state_goal_preserved': state.goal == goal, + 'history_unchanged': c.conversation_history == before, + 'continuation_still_contains_full_goal': goal in goals.GoalManager(sid).next_continuation_prompt()} + +output['cli_literal_witness'] = cli_case(history, original, tag+'-cli') +output['cli_empty_control'] = cli_case([], 'Ship the release', tag+'-empty') +output['cli_unrelated_control'] = cli_case([{'role':'user','content':'Unrelated'}], original, tag+'-unrelated') +output['cli_block_content'] = cli_case([{'role':'user','content':[{'type':'text','text': original}]}], original, tag+'-block') +options = [{'role':'user','content':'We can ship the API or ship the UI. Wait for my choice.'}, {'role':'assistant','content':'Which one should I work on?'}] +output['selection_api'] = cli_case(options, 'ship the API', tag+'-api') +output['selection_ui'] = cli_case(options, 'ship the UI', tag+'-ui') +output['different_selected_goals_same_model_prompt'] = output['selection_api']['prompt'] == output['selection_ui']['prompt'] +# /goal draft invokes its only paid dependency as an explicit unavailable stub. +c = make_cli(history, tag+'-draft') +with patch('hermes_cli.goals.draft_contract', return_value=None), contextlib.redirect_stdout(io.StringIO()): + assert c.process_command('/goal draft ' + original) +output['draft_fallback_prompt'] = c._pending_input.get_nowait() + +# Actual registered TUI/Desktop command dispatcher; goal bypasses CLI slash worker. +sid = tag+'-tui' +server._sessions[sid] = {'session_key': sid, 'history': copy.deepcopy(history), 'history_lock': threading.Lock(), 'history_version':0, 'running':False, 'attached_images': [], 'cols': 120} +rpc = server._methods['slash.exec'](1, {'command':'goal '+original, 'session_id':sid}) +output['tui_slash_rpc'] = rpc +assert goals.GoalManager(sid).state.goal == original + +# Real gateway handler + real enqueue method. Capture only adapter transport FIFO seam. +class GatewayProbe(GatewayGoalCommandsMixin): + def __init__(self): + self.mgr = goals.GoalManager(tag+'-gateway') + self.events = [] + self.conversation_history = copy.deepcopy(history) + async def _get_goal_manager_for_event(self, event): + return self.mgr, None + def _adapter_and_key_for(self, event): + return object(), 'probe-key' + def _enqueue_fifo(self, key, event, adapter): + self.events.append(event) +gw = GatewayProbe() +event = MessageEvent(text='/goal '+original, message_type=MessageType.TEXT, + source=SessionSource(platform=Platform.DISCORD, chat_id='probe', chat_type='dm', user_id='probe'), message_id='goal-probe') +asyncio.run(gw._handle_goal_command(event)) +assert len(gw.events) == 1 +output['gateway_kick'] = {'prompt':gw.events[0].text, 'goal_preserved':gw.mgr.state.goal == original} +# Actual archived replay turn, including useful actions and terminal waits. +next_id = source.execute("SELECT min(id) FROM messages WHERE session_id=? AND id>? AND role='user'", ('20260902_073639_918cf3',267045)).fetchone()[0] +rows = source.execute('SELECT id,role,content,tool_calls,tool_name,timestamp FROM messages WHERE session_id=? AND id>=? AND id1 else os.getcwd()) # repo root under test +os.environ['HERMES_HOME']=tempfile.mkdtemp(dir=root) +os.environ['TERMINAL_ENV']='local' +from tools import file_tools as f +from tools.file_operations import WriteResult +print('module',f.__file__) +with tempfile.TemporaryDirectory(dir=root) as d: + p=pathlib.Path(d)/'file.txt' + for n in [1000,4000,10000,20000]: + old='same repetitive record\n'*n;p.write_text(old, encoding='utf-8');new=old+'end\n' + start=time.monotonic(); hint=f._whole_file_rewrite_hint('default',str(p),new) + print('repeat',n,len(old),round(time.monotonic()-start,3),bool(hint),flush=True) + old=''.join(f'{i} arbitrary unique data for test\n' for i in range(1500));p.write_text(old, encoding='utf-8') + new=old.replace('750 arbitrary','750 edited') + start=time.monotonic();r=json.loads(f.write_file_tool(str(p),new,task_id='review')) + print('live-write',round(time.monotonic()-start,3),r,p.read_text(encoding='utf-8')==new,flush=True) + class RemoteOps: + env=object() + def write_file(self,path,content): + self.written=(path,content) + return WriteResult(bytes_written=len(content),verified=True) + ops=RemoteOps() + p.write_text(old, encoding='utf-8') + with patch.object(f,'_get_file_ops',return_value=ops): + print('remote_backend_is_host',f._file_ops_uses_host_paths(ops)) + r=json.loads(f.write_file_tool(str(p),new,task_id='remote-review')) + print('remote_empty_target_hint_from_host',r.get('hint'), 'host_unchanged',p.read_text(encoding='utf-8')==old,flush=True) + fifo=pathlib.Path(d)/'pipe.txt';os.mkfifo(fifo) + alarms=[] + def alarm(sig,frame): + alarms.append(time.monotonic()) + raise TimeoutError('host FIFO read blocked') + signal.signal(signal.SIGALRM,alarm) # windows-footgun: ok — POSIX-only FIFO hazard probe + with patch.object(f,'_get_file_ops',return_value=ops): + start=time.monotonic();signal.alarm(2) + try: print('remote_fifo',f.write_file_tool(str(fifo),'x'*20000,task_id='remote-fifo'), 'seconds',time.monotonic()-start,'read_alarm_fired',bool(alarms)) + finally:signal.alarm(0) + with patch.object(f,'_whole_file_rewrite_hint',return_value=None): + start=time.monotonic();print('remote_fifo_base_no_hint',f.write_file_tool(str(fifo),'x'*20000,task_id='remote-fifo'), 'seconds',time.monotonic()-start) diff --git a/evals/postmortem/review_probes/scanner_bypass_probe.py b/evals/postmortem/review_probes/scanner_bypass_probe.py new file mode 100644 index 0000000000..aaa47a1041 --- /dev/null +++ b/evals/postmortem/review_probes/scanner_bypass_probe.py @@ -0,0 +1,55 @@ +"""#103492: a hardline command hidden after a NEWLINE inside a double-quoted $(grep …) must stay blocked +by the public guard with no approval callback; a grep with a backtick operand must stay allowed. + +Independent-review probe (written by the /review subagent for tracking issue #103563, adapted). It +reproduced an approval bypass in the first version of the PR (approved, 0 callbacks). Runs the real +public guard (check_dangerous_command) in a temp HERMES_HOME with approvals.mode=manual and executes a +HARMLESS Bash witness (reboot shadowed by a function writing a marker) to prove reachability. + +Usage: python scanner_bypass_probe.py [] +""" +import importlib.util +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +ROOT=Path(sys.argv[1]).resolve() # repo root under test +sys.path.insert(0,str(ROOT)) +home=Path(tempfile.mkdtemp(prefix='review103492-confirm-')) +os.environ['HERMES_HOME']=str(home) +os.environ['HERMES_INTERACTIVE']='1' +(home/'config.yaml').write_text('approvals:\n mode: manual\n timeout: 1\n', encoding='utf-8') +from tools import approval_detection as d +from tools import approval as a +b=None +if len(sys.argv)>2: # optional: path to a baseline approval_detection.py (e.g. from main) to compare verdicts + spec=importlib.util.spec_from_file_location('baseline',sys.argv[2]); b=importlib.util.module_from_spec(spec); spec.loader.exec_module(b) +r='re'+'boot' +cmd='echo "$(grep -P \'safe\' /dev/null\n'+r+')"' +print(json.dumps({'module':d.__file__,'approval_module':a.__file__,'case':'newline command','command':cmd,'baseline':(b.detect_hardline_command(cmd) if b else None),'head':d.detect_hardline_command(cmd),'variants':list(d._command_detection_variants(cmd))})) +calls=[] +def deny(*args,**kwargs): + calls.append([args,kwargs]);return 'deny' +print(json.dumps({'pattern_guard':a.check_dangerous_command(cmd,'local',approval_callback=deny),'combined_guard':a.check_all_command_guards(cmd,'local',approval_callback=deny),'callback_calls':calls},default=str)) +# Shadow the dangerous command name with a harmless function; verify shell +# execution and marker contents, never invoke a host shutdown binary. +marker=home/'marker' +prefix=r+'(){ printf executed > "$MARKER"; }; ' +p=subprocess.run(['/bin/bash','--noprofile','--norc','-c',prefix+cmd],env={'PATH':'/usr/bin:/bin','HOME':str(home),'MARKER':str(marker)},capture_output=True,text=True, encoding='utf-8', errors='replace',timeout=5) +assert marker.read_text(encoding='utf-8')=='executed' +print(json.dumps({'case':'newline safe execution','exit':p.returncode,'marker':marker.read_text(encoding='utf-8'),'stdout':p.stdout,'stderr':p.stderr})) +# Prove the benign backtick argument really is well-formed and matches input. +f=home/'f';f.write_text('needle\n', encoding='utf-8') +cmd='grep -e `printf needle` '+str(f) +p=subprocess.run(['/bin/bash','--noprofile','--norc','-c',cmd],env={'PATH':'/usr/bin:/bin','HOME':str(home)},capture_output=True,text=True, encoding='utf-8', errors='replace',timeout=5) +assert p.returncode==0 and p.stdout=='needle\n' +print(json.dumps({'case':'backtick argument','command':cmd,'baseline':(b.detect_hardline_command(cmd) if b else None),'head':d.detect_hardline_command(cmd),'tokens':d._shell_tokens_with_spans(cmd,0),'exit':p.returncode,'stdout':p.stdout})) +# Reporter's exact spelling, fixture makes the sed address meaningful. +with tempfile.TemporaryDirectory() as tmp: + Path(tmp,'f').write_text('X\ny\nz\nw\n', encoding='utf-8') + cmd='sed -n "$(grep -n X f | cut -d: -f1),+3p" f' + p=subprocess.run(['/bin/bash','--noprofile','--norc','-c',cmd],cwd=tmp,env={'PATH':'/usr/bin:/bin','HOME':str(home)},capture_output=True,text=True, encoding='utf-8', errors='replace',timeout=5) + assert p.returncode==0 and p.stdout=='X\ny\nz\nw\n' + print(json.dumps({'case':'reported real fixture','baseline':(b.detect_hardline_command(cmd) if b else None),'head':d.detect_hardline_command(cmd),'exit':p.returncode,'stdout':p.stdout})) diff --git a/evals/postmortem/run.py b/evals/postmortem/run.py new file mode 100644 index 0000000000..5d078e3ab0 --- /dev/null +++ b/evals/postmortem/run.py @@ -0,0 +1,87 @@ +#!/usr/bin/env python3 +"""Run the post-mortem harness against one or two checkouts and print a comparison table. + + python -m evals.postmortem.run --repo /path/to/checkout # one ref: pass/fail per probe + python -m evals.postmortem.run --repo A --compare B # two refs: side by side + python -m evals.postmortem.run --repo A --live # also the probes that spend money + +Each probe is a standalone script run in a fresh interpreter with the target checkout on sys.path and a +temp HERMES_HOME (probes that need real credentials say so and are only run with --live). A probe +"passes" when its process exits 0 AND its stdout contains the expected marker documented in +PROBES below; the marker is the behaviour the corresponding PR fixed. Run the forensics lanes +separately (they need a state.db copy): see forensics/README section in ../README.md. +""" +from __future__ import annotations + +import argparse +import os +import subprocess +import sys +import tempfile +from pathlib import Path + +HERE = Path(__file__).resolve().parent + +# (script, args-template, expected stdout substring, live?, PR) +# Two probes pass on main as well: scanner_bypass_probe (main blocked the witness too, as "malformed") and +# notice_delivery_probe (main has no interim notice to mis-deliver). They guard against regressing INTO the +# round-1 defects, which is why they are here. +PROBES = [ + ("live_ab/hardline_scanner_matrix.py", ["{repo}"], "ALL OK", False, "#103492"), + ("review_probes/scanner_bypass_probe.py", ["{repo}"], '"hardline": true', False, "#103492"), + ("live_ab/subagent_context_cap.py", ["{repo}"], "trigger=200,000", False, "#103513"), + ("live_ab/batch_failure_notice.py", ["{repo}"], "TASK_FAILURE_NOTICE", False, "#103549"), + ("review_probes/notice_delivery_probe.py", ["{repo}"], "PROBE_COMPLETE", False, "#103549"), + ("review_probes/cache_estimator_probe.py", ["{repo}"], '"preflight_should_compress": true', False, "#103476"), + ("review_probes/rewrite_hint_probe.py", ["{repo}"], "remote_fifo", False, "#103551"), + ("live_ab/auth_stampede.py", ["{repo}", "12"], "server_401=0", False, "#103526"), + ("review_probes/credential_identity_probe.py", ["{repo}", "pr"], '"after_sub": "account-A"', False, "#103526"), + # live (real provider calls, cents each) + ("live_ab/goal_judge_wait.py", ["{repo}", "3"], "('wait', ", True, "#103534"), + ("live_ab/cache_prefix_wire.py", ["{repo}", "B"], "", True, "#103476"), +] + + +def run_probe(script: str, args: list[str], repo: str, timeout: int = 240) -> tuple[int, str]: + env = dict(os.environ) + env.setdefault("HERMES_HOME", tempfile.mkdtemp(prefix="pm-probe-")) + env["PYTHONPATH"] = repo + os.pathsep + env.get("PYTHONPATH", "") + cmd = [sys.executable, str(HERE / script), *[a.format(repo=repo) for a in args]] + try: + p = subprocess.run(cmd, cwd=repo, capture_output=True, text=True, encoding="utf-8", errors="replace", timeout=timeout, env=env) + return p.returncode, (p.stdout + "\n" + p.stderr) + except subprocess.TimeoutExpired: + return 124, "TIMEOUT" + + +def main(argv=None) -> int: + ap = argparse.ArgumentParser(description=(__doc__ or "").split("\n\n")[0]) + ap.add_argument("--repo", required=True) + ap.add_argument("--compare", default=None, help="second checkout to run side by side (e.g. main vs branch)") + ap.add_argument("--live", action="store_true", help="also run probes that make real provider calls") + ap.add_argument("--only", default=None, help="substring filter on script path or PR number") + a = ap.parse_args(argv) + repos = [a.repo] + ([a.compare] if a.compare else []) + rows = [] + for script, args, marker, live, pr in PROBES: + if live and not a.live: + continue + if a.only and a.only not in script and a.only not in pr: + continue + cells = [] + for repo in repos: + rc, out = run_probe(script, args, os.path.abspath(repo)) + ok = rc == 0 and (marker in out if marker else True) + cells.append("PASS" if ok else f"FAIL(rc={rc})") + rows.append((pr, script, *cells)) + width = max(len(r[1]) for r in rows) if rows else 20 + head = f"{'PR':<9} {'probe':<{width}} " + " ".join(f"{os.path.basename(os.path.normpath(r)):<14}" for r in repos) + print(head); print("-" * len(head)) + for r in rows: + print(f"{r[0]:<9} {r[1]:<{width}} " + " ".join(f"{c:<14}" for c in r[2:])) + failed = any("FAIL" in c for r in rows for c in r[2 + (1 if a.compare else 0):]) # only the LAST column gates + return 1 if failed else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/evals/postmortem/tests/__init__.py b/evals/postmortem/tests/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/evals/postmortem/tests/test_postmortem_harness.py b/evals/postmortem/tests/test_postmortem_harness.py new file mode 100644 index 0000000000..870e81cfae --- /dev/null +++ b/evals/postmortem/tests/test_postmortem_harness.py @@ -0,0 +1,88 @@ +"""The post-mortem forensics run end-to-end on a synthetic state.db and report the run population correctly. + +Guards the harness itself (it lives in evals/, outside the normal import graph): the root is discovered, +a compression-rollover child is excluded, pricing is fitted, and each lane writes its JSON without error. +""" +import json +import sqlite3 +import sys +import time +from pathlib import Path + +import pytest + +EVALS = Path(__file__).resolve().parents[2] + + +def _mk_db(path: Path) -> None: + c = sqlite3.connect(path) + c.executescript(""" + CREATE TABLE sessions(id TEXT PRIMARY KEY, parent_session_id TEXT, source TEXT, started_at REAL, ended_at REAL, + api_call_count INTEGER, input_tokens INTEGER, cache_read_tokens INTEGER, cache_write_tokens INTEGER, + output_tokens INTEGER, estimated_cost_usd REAL, system_prompt_hash TEXT); + CREATE TABLE system_prompts(hash TEXT PRIMARY KEY, prompt TEXT); + CREATE TABLE messages(id INTEGER PRIMARY KEY, session_id TEXT, role TEXT, content TEXT, tool_calls TEXT, + tool_name TEXT, reasoning TEXT, timestamp REAL); + CREATE TABLE state_meta(key TEXT PRIMARY KEY, value TEXT); + """) + t0 = time.time() - 4000 + c.execute("INSERT INTO system_prompts VALUES ('h1', ?)", ("x" * 35_000,)) + + def sess(sid, parent, source, start, end, calls, cr, cw, out): + cost = cr * 0.2e-6 + cw * 10e-6 + out * 40e-6 + c.execute("INSERT INTO sessions VALUES (?,?,?,?,?,?,?,?,?,?,?,?)", (sid, parent, source, start, end, calls, 0, cr, cw, out, cost, "h1")) + + sess("root", None, "cli", t0, t0 + 3600, 40, 2_000_000, 300_000, 40_000) + for i in range(12): # children at depth 1 + sess(f"c{i}", "root", "subagent", t0 + 10 + i, t0 + 600 + 60 * i, 20, 1_000_000, 200_000, 20_000) + sess("g0", "c0", "subagent", t0 + 20, t0 + 500, 10, 500_000, 100_000, 10_000) # depth 2 + sess("rollover", "root", "cli", t0 + 3601, t0 + 7200, 30, 9_000_000, 900_000, 90_000) # excluded + sess("unrelated", None, "telegram", t0, t0 + 100, 1, 1000, 100, 10) + # messages: a delegate_task timeout in c0, a truncated batch block in root, a hardline block, a nudge + c.execute("INSERT INTO messages(session_id,role,content,tool_name,timestamp) VALUES ('c0','tool',\"Error executing tool 'delegate_task': timed out after 420.0s\",'delegate_task',?)", (t0 + 100,)) + c.execute("INSERT INTO messages(session_id,role,content,tool_calls,timestamp) VALUES ('c0','assistant','','[{\"function\":{\"name\":\"terminal\",\"arguments\":\"{\\\\\"command\\\\\": \\\\\"sleep 600\\\\\"}\"}}]',?)", (t0 + 200,)) + c.execute("INSERT INTO messages(session_id,role,content,timestamp) VALUES ('root','user','[ASYNC DELEGATION BATCH COMPLETE — d]\\n--- ✓ TASK 1/2 ...\\n[SUMMARY TRUNCATED]\\n--- ✓ TASK 2/2 ...',?)", (t0 + 700,)) + c.execute("INSERT INTO messages(session_id,role,content,tool_name,timestamp) VALUES ('c1','tool','BLOCKED (hardline): command parser limit or malformed executable payload','terminal',?)", (t0 + 300,)) + c.execute("INSERT INTO messages(session_id,role,content,timestamp) VALUES ('root','assistant','Waiting on 3 batches; nothing to dispatch.',?)", (t0 + 800,)) + c.execute("INSERT INTO messages(session_id,role,content,timestamp) VALUES ('root','user','[Continuing toward your standing goal]\\nGoal: x',?)", (t0 + 810,)) + c.execute("INSERT INTO state_meta VALUES ('goal:root', ?)", (json.dumps({"status": "active", "waiting_on_session": "proc_x", "waiting_since": t0 + 900, "last_verdict": "wait", "turns_used": 3}),)) + c.commit(); c.close() + + +@pytest.fixture +def harness_path(monkeypatch): + monkeypatch.syspath_prepend(str(EVALS.parent)) + return EVALS + + +def test_run_population_excludes_rollover_and_unrelated_sessions(tmp_path, harness_path): + from evals.postmortem.forensics.common import Run + db = tmp_path / "state.db"; _mk_db(db) + run = Run.open(str(db), out=str(tmp_path / "out")) + assert run.root == "root" + assert set(run.in_run) == {"root", *{f"c{i}" for i in range(12)}, "g0"} + assert "rollover" not in run.in_run and "unrelated" not in run.in_run + s = run.summary() + assert s["by_depth"] == {0: 1, 1: 12, 2: 1} + assert abs(s["fitted_price_per_million"]["cache_write_tokens"] - 10.0) < 0.01 + assert abs(s["cost_usd"] - sum(run.cost(x) for x in run.in_run)) < 1e-6 + + +def test_every_lane_runs_and_writes_its_report(tmp_path, harness_path): + from evals.postmortem.forensics import delegation, goal_loop, tokens, tools + db = tmp_path / "state.db"; _mk_db(db) + out = tmp_path / "out" + for lane in (tokens, delegation, tools, goal_loop): + assert lane.main(["--db", str(db), "--out", str(out)]) == 0 + assert json.loads((out / "delegation.json").read_text(encoding="utf-8"))["observed"]["delegate_task_timeouts"] == 1 + assert json.loads((out / "tools.json").read_text(encoding="utf-8"))["observed"]["hardline_blocks_malformed_class"] == 1 + g = json.loads((out / "goal_loop.json").read_text(encoding="utf-8"))["observed"] + assert g["nudges"] == 1 and g["nudges_within_180s_of_a_waiting_turn"] == 1 + assert (out / "tokens.json").exists() + + +def test_runner_lists_a_probe_per_pr(harness_path): + from evals.postmortem import run as runner + prs = {p[4] for p in runner.PROBES} + assert {"#103492", "#103513", "#103549", "#103476", "#103551", "#103526", "#103534"} <= prs + assert sys.version_info >= (3, 10)