feat(evals): post-mortem harness — forensics lanes + live A/B + review probes for the #102117 run fixes

evals/postmortem/ turns the one-off audit behind tracking issue #103563 into
something anyone with a Hermes state.db copy (and optionally rotated
agent.log*) can run on their own fan-out:

  forensics/   common.py discovers the run tree (root = most descendants,
               compression-rollover children excluded so cost buckets stay
               disjoint), fits pricing from estimated_cost_usd, and five lanes
               recompute the OBSERVED figures: tokens (buckets, depth/duration
               shares, context reconstruction, excess-cache-write proxy, cap
               replay), logcalls (per-call cache behaviour from agent.log with
               coverage printed first; strict and loose plateau definitions
               reported separately), delegation (timeouts, orphaned children,
               polling hours, batch-join withheld child-hours, truncated
               summaries), tools (hardline blocks, foreground refusals,
               whole-file rewrites), goal_loop (nudges, parked barrier), rework
               (public-surface drop at PR open + post-open commit inventory).
               Every figure is labeled OBSERVED or MODELED.
  live_ab/     the per-PR A/Bs (real code paths, fake providers, temp
               HERMES_HOME), paths from argv.
  review_probes/ the independent /review's probes, credited and adapted; each
               reproduced a round-1 defect and the fixed head must pass it.
  run.py       runs the offline probes against one or two checkouts and prints
               PASS/FAIL side by side (--live adds the ones that spend cents).
  tests/       synthetic-DB smoke test for the lanes and runner.

On the run's DB the lanes reproduce the tracking issue's population exactly
(1,394 sessions, 93,284 calls, $19,302.59; cache_write $11,159.76) and on
main vs an integration checkout of the 13 PRs the runner shows every probe
FAIL -> PASS (two guard-only probes pass on both, noted in run.py).

The trajectories are deliberately not shipped: the DB holds 51,956 home
paths, 5,341 e-mails, private IPs, chat ids and real-shaped credentials in
tool output. The lane reports and recomputed JSON are in a secret gist
linked from #103563.
This commit is contained in:
Teknium
2026-09-05 09:13:10 -07:00
parent 006b1beb00
commit a8ca904922
33 changed files with 2306 additions and 0 deletions
+117
View File
@@ -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 <hermes-checkout>
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 <session_id>] [--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 <merge-base> --open <sha-at-open> --head <merged-sha>
```
Each writes `postmortem_out/<lane>.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 <repo> 40` | #103526 | `401s=0` (main: 40) |
| `review_probes/credential_identity_probe.py <repo> 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 <repo> B` (live) | #103476 | 0 mutated prefixes across 6 calls |
| `live_ab/goal_judge_wait.py <repo> 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.
View File
+186
View File
@@ -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},
}
+115
View File
@@ -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())
+70
View File
@@ -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())
+136
View File
@@ -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 [<session_id>] agent.conversation_loop: API call #N: model=... in=<prompt> out=<out> total=... latency=..s cache=<hit>/<total>
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())
+76
View File
@@ -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 <merge-base> --open <sha-at-pr-open> --head <merged-sha>
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())
+122
View File
@@ -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())
+100
View File
@@ -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())
+79
View File
@@ -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 <repo_root> <n_agents>
"""
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")
@@ -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 <repo_root>"""
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)
@@ -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 <repo_root> <A|B>
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}))
@@ -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 <repo_root> <A|B> [--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)")
@@ -0,0 +1,19 @@
"""Live judge A/B on the run's real "waiting" response shape. Usage: python judge_ab.py <repo_root> [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)
@@ -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 <repo_root>
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")
@@ -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 <repo_root> (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)}))
@@ -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)
@@ -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: <repo_root>
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')
@@ -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)
@@ -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'
@@ -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)
@@ -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))
@@ -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 id<? ORDER BY id', ('20260902_073639_918cf3',267045,next_id)).fetchall()
output['archive_replay_turn'] = {'first_id':rows[0]['id'], 'last_id':rows[-1]['id'], 'elapsed_seconds':rows[-1]['timestamp']-rows[0]['timestamp'], 'assistant_messages':sum(r['role']=='assistant' for r in rows), 'tool_messages':sum(r['role']=='tool' for r in rows), 'first_assistant':next(r['content'] for r in rows if r['role']=='assistant'), 'tools':[{k:r[k] for k in ['id','tool_name','timestamp']} for r in rows if r['role']=='tool']}
try:
import tiktoken
enc=tiktoken.get_encoding('cl100k_base')
output['illustrative_cl100k_tokens']={'original':len(enc.encode(original)), 'kick':len(enc.encode(output['cli_literal_witness']['prompt']))}
except Exception as exc:
output['tokenizer_unavailable']=str(exc)
out=Path(os.environ.get('PROBE_OUT','.'))/f'goaldup-{tag}-probe.json'
out.write_text(json.dumps(output, indent=2), encoding='utf-8')
print(json.dumps({'artifact':str(out), 'tag':tag, 'cli_chars':output['cli_literal_witness']['prompt_chars'], 'tui_chars':len(rpc['result']['message']), 'gateway_chars':len(gw.events[0].text), 'ambiguous_selection':output['different_selected_goals_same_model_prompt'], 'archive':output['archive_replay_turn'], 'tokens':output.get('illustrative_cl100k_tokens')}, indent=2))
server._sessions.clear()
@@ -0,0 +1,78 @@
"""#103496/#103534: judge process scoping and delegation WAIT lifecycle
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,subprocess,time,queue,asyncio,contextlib
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import patch
sys.path.insert(0,os.getcwd())
home=tempfile.TemporaryDirectory(prefix='goals-probe-')
os.environ['HERMES_HOME']=home.name
os.environ['HERMES_TEST_MODE']='1'
from hermes_cli import goals
from tools import process_registry as pr,async_delegation as ad
from hermes_cli.cli_loops_mixin import CLILoopsMixin
from gateway.run_goals import GatewayGoalsMixin
from tui_gateway import server as pt
out={'module':goals.__file__,'sha':subprocess.check_output(['git','rev-parse','HEAD'],text=True, encoding='utf-8', errors='replace').strip()}
proc=subprocess.Popen([sys.executable,'-c','import time; time.sleep(300)'],stdin=subprocess.DEVNULL)
try:
reg=pr.process_registry
now=time.time()
records={'proc_own':pr.ProcessSession(id='proc_own',command='own sleeper',task_id='default',owner_task_id='root',pid=proc.pid,started_at=now), 'proc_child':pr.ProcessSession(id='proc_child',command='child sleeper',task_id='default',owner_task_id='sa-child',pid=proc.pid,started_at=now)}
deleg={'d':{'status':'running','parent_session_id':'root','session_key':'','origin_ui_session_id':''}}
with patch.object(reg,'_running',records),patch.object(reg,'_finished',{}),patch.object(ad,'_records',deleg):
out['registry']=reg.list_sessions()
out['all']=[r['session_id'] for r in goals.gather_background_processes()]
if hasattr(goals,'count_active_delegations'):out['active']=goals.count_active_delegations('root')
captured=[]
def judge(*a,**kw):
captured.append({'processes':[x['session_id'] for x in kw.get('background_processes',[]) or []], 'active':kw.get('active_delegations',0)})
return ('continue','test',False,None,False)
class CLI(CLILoopsMixin):
def _get_goal_manager(self): return self.mgr
c=CLI();c.session_id='root';c.agent=SimpleNamespace(session_id='root');c._pending_input=queue.Queue();c.conversation_history=[{'role':'assistant','content':'Waiting on workers'}]
class GW(GatewayGoalsMixin):
async def _post_turn_manager(self,*a):return self.mgr
async def _run_in_executor_with_context(self,fn):return fn()
g=GW()
with patch.object(goals,'judge_goal',side_effect=judge):
for label in ['cli','gateway','tui']:
m=goals.GoalManager(session_id='root');m.set('finish');c.mgr=g.mgr=m
if label=='cli':c._maybe_continue_goal_after_turn()
elif label=='gateway':asyncio.run(g._post_turn_goal_continuation(session_entry=SimpleNamespace(session_id='root'),source=None,final_response='Waiting on workers'))
else:
with patch.object(pt,'_active_goal_manager',return_value=m),patch.object(pt,'_plan_goal_compression_recovery',return_value=(None,None)),patch.object(pt,'_emit'):
pt._goal_followup_after_turn('tab',{'session_key':'root'}, {'final_response':'Waiting on workers'},'complete','Waiting on workers')
out[label]=captured[-1]
m=goals.GoalManager(session_id='expiry');m.set('finish');m.wait_on(proc.pid);m.state.waiting_since=now-1801;m._save()
out['old_pid_waiting']=m.is_waiting()
m.wait_on_session('proc_own');m.state.waiting_since=now-1801;m._save();out['old_session_waiting']=m.is_waiting()
m.wait_for_seconds(3600);m.state.waiting_since=now-1801;m._save();out['long_timer_still_waiting']=m.is_waiting()
m.stop_waiting();m.wait_for_seconds(1200,reason='delegates');deleg['d']['status']='completed'
with patch.object(goals,'judge_goal',side_effect=judge):
before=len(captured)
kw={'active_delegations':0} if hasattr(goals,'count_active_delegations') else {}
out['after_delegation_complete']=m.evaluate_after_turn('Workers finished; next step is integration',**kw)
out['completion_judge_calls']=len(captured)-before
out['persisted_wait_until']=goals.GoalManager(session_id='expiry').state.waiting_until
# Drive the actual new judge branch, then deliver a worker-result turn.
c._pending_input=queue.Queue(); c.mgr=goals.GoalManager(session_id='root');c.mgr.set('finish integration')
deleg['d']['status']='running'; calls=[]
def fake_aux(*args,**kwargs):
text=str(kwargs.get('messages') or args);calls.append(text)
content='{"verdict":"wait","wait_for_seconds":1200,"reason":"workers"}' if 'Active delegations:' in text else '{"verdict":"continue","reason":"integrate next"}'
return SimpleNamespace(choices=[SimpleNamespace(message=SimpleNamespace(content=content))])
with patch('agent.auxiliary_client.call_llm',side_effect=fake_aux):
c._maybe_continue_goal_after_turn();out['lifecycle_initial_pending']=c._pending_input.qsize()
while not c._pending_input.empty():c._pending_input.get_nowait()
deleg['d']['status']='completed';c.conversation_history=[{'role':'assistant','content':'Workers finished; I still need to integrate their branches.'}]
c._maybe_continue_goal_after_turn()
out['lifecycle_completion_pending']=c._pending_input.qsize();out['lifecycle_aux_calls']=len(calls);out['lifecycle_still_waiting']=c.mgr.is_waiting()
finally:
proc.terminate();proc.wait(timeout=10)
Path(sys.argv[1]).write_text(json.dumps(out,indent=2), encoding='utf-8')
print(json.dumps(out,indent=2))
@@ -0,0 +1,86 @@
"""#103549: interim notice must not claim/ack/dedup the batch's final result
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, socket, json, time, threading, queue, asyncio
from pathlib import Path
from types import SimpleNamespace
from collections import OrderedDict
root=Path(sys.argv[1]).resolve(); sys.path.insert(0,str(root))
for key in list(os.environ):
if key.startswith('HERMES_') or any(s in key for s in ('API_KEY','TOKEN','SECRET')):
os.environ.pop(key,None)
os.environ['HERMES_HOME']=tempfile.mkdtemp(prefix='notice-review-')
socket.socket.connect=lambda *a,**k: (_ for _ in ()).throw(RuntimeError('network forbidden'))
import tools.async_delegation as ad
import tools.delegate_tool_dispatch as dd
from tools.process_registry import process_registry as reg
from tools.process_registry_notifications import format_process_notification
from gateway.run_notifications import GatewayNotificationsMixin
assert Path(ad.__file__).is_relative_to(root)
assert Path(dd.__file__).is_relative_to(root)
print('SOURCE',ad.__file__,dd.__file__,flush=True)
class Sink(GatewayNotificationsMixin):
def __init__(self):
self._completion_delivery_lock=threading.Lock()
self._completion_deliveries_inflight=set()
self._completion_deliveries_delivered=OrderedDict()
self._completion_delivery_retention=100
self.received=[]
async def _classify_completion_target(self,sid): return 'deliver'
async def _inject_watch_notification(self,text,evt):
self.received.append(('notice' if evt.get('task_failure_notice') else 'final',evt.get('results')))
return True
def batch_run(did, gates):
tasks=[{'goal':f'worker task {i}'} for i in range(3)]
children=[(i,t,SimpleNamespace()) for i,t in enumerate(tasks)]
b=dd._Batch(tasks,children,SimpleNamespace(quiet_mode=True),{'model':'offline'},None,'leaf',3,did,[],[], '', '',None,None,time.monotonic())
def child(i,t,c):
assert gates[i].wait(15)
return {'task_index':i,'status':'error' if i<2 else 'completed','error':'offline failure' if i<2 else None,'summary':None if i<2 else 'FINISHED_SUCCESS','duration_seconds':0.1}
b.run_child=child
results=[]
kwargs={'honor_parent_interrupt':False}
if 'detached' in __import__('inspect').signature(dd._run_children_parallel).parameters: kwargs['detached']=True
dd._run_children_parallel(b,results,**kwargs)
return {'results':results,'total_duration_seconds':1}
def start(did):
gates=[threading.Event() for _ in range(3)]
h=ad.dispatch_async_delegation_batch(goals=['worker task 0','worker task 1','worker task 2'],context=None,toolsets=None,role='leaf',model='offline',session_key='owned',parent_session_id='parent',runner=lambda:batch_run(did,gates),delegation_id=did)
assert h['status']=='dispatched'
return gates
def get(timeout=2): return reg.completion_queue.get(timeout=timeout)
def deliver(sink,e): return asyncio.run(sink._deliver_completion_notification(format_process_notification(e),e))
results={}
if not hasattr(ad,'push_task_failure_notice'):
g=start('base-control');g[0].set()
try:get(.2);raise AssertionError('unexpected early event')
except queue.Empty:pass
g[1].set();g[2].set();e=get();sink=Sink();results['main_control']={'no_early_notice':True,'final_delivered':deliver(sink,e),'received':sink.received}
else:
g=start('gateway-before-final');sink=Sink();g[0].set();first=get();assert first['task_failure_notice']
results['gateway']={'first_notice_accepted':deliver(sink,first),'after_notice':ad.get_durable_delegation('gateway-before-final')}
g[1].set();second=get();results['gateway']['second_notice_accepted']=deliver(sink,second)
g[2].set();final=get();assert not final.get('task_failure_notice')
results['gateway']['final_accepted']=deliver(sink,final)
results['gateway']['received']=sink.received
results['gateway']['after_final']=ad.get_durable_delegation('gateway-before-final')
print('GATEWAY_RECEIVED', [k for k,_ in sink.received])
assert [kind for kind,_ in sink.received]==['notice','notice','final'], 'fixed head must deliver both notices AND the final'
# Parent busy: the queued first notice gets accepted only after final persisted.
g=start('busy-parent');g[0].set();notice=get();g[1].set();second=get();g[2].set();final=get()
claim=ad.claim_event_delivery(notice,'tui-poller');assert claim=='', 'an interim notice must be NON-durable (empty claim token)'
ad.complete_event_delivery(notice,claim)
final_claim=ad.claim_event_delivery(final,'tui-poller')
results['busy_parent']={'notice_claimed':bool(claim),'final_claim':final_claim,'row':ad.get_durable_delegation('busy-parent')}
assert final_claim, 'the final result must still be claimable after a notice was consumed'
# Different TUI UI-dedup keys do not fix shared durable claim identity.
from tui_gateway.session_notifications import _notification_event_dedup_key
results['busy_parent']['ui_keys_differ']=_notification_event_dedup_key(notice)!=_notification_event_dedup_key(final)
print(json.dumps(results,indent=2),flush=True)
print('PROBE_COMPLETE',flush=True)
@@ -0,0 +1,48 @@
"""#103551: remote backend, FIFO, and quadratic-time hazards of the rewrite hint
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 sys, os, pathlib, tempfile, json, time, signal
from unittest.mock import patch
root=pathlib.Path(tempfile.mkdtemp(prefix='hint-probe-'))
sys.path.insert(0, sys.argv[1] if len(sys.argv)>1 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)
@@ -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 <repo_root> [<baseline_approval_detection.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}))
+87
View File
@@ -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())
View File
@@ -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)