Paper2Agent_decoupleRpy / scripts /analyze_post_analysis_reads.py
Annie Voigt
feat(metrics): instrument n_post_analysis_reads to test the 808b91d claim
0e597a2
Raw
History Blame Contribute Delete
9.96 kB
#!/usr/bin/env python3
"""Retrospective before/after analysis of post-analysis CSV re-reads.
Answers the question `1ec068b` shipped without answering: **did returning
`top_table` from `dataset_compare_activity_by_group` actually stop the agent
re-opening `comparison_path` before writing its solution?**
The intent was a retrospective control arm: historical run records carry their
own `messages`, so `n_post_analysis_reads` could be recomputed with no re-running
and no reverting (both Spaces run the new code, so there is no old-code
deployment left to A/B against).
**That control arm does not exist in the data.** As of 2026-08-06 all 20 run
traces in `anne-voigt/decoupleRpy_results` persist `"messages": []`, because the
Gradio path drives `graph.stream()` directly and never fills
`WorkflowEngine.message_history` — the field `get_trace()` reads. Every
production trace was therefore written empty. `src/agent.py` now falls back to
the final graph state, so records written from here on carry the conversation;
the pre-fix baseline is unrecoverable and has to be rebuilt forward. This script
reports that explicitly (`no_messages`) rather than scoring an empty record as
zero re-reads.
The counting logic is imported from ``src/core/run_metrics.py`` — the SAME code
the live instrumentation uses. Do not re-implement it here: a divergence between
the two would silently invalidate the comparison.
Arm assignment is by the record's own ``execution_time`` versus ``--cutoff``
(the `1ec068b` prod deploy). It deliberately does NOT use which sink the run
landed in: dev runs currently write to the PROD sink despite
``LOG_SINK_HF_DATASET`` pointing at ``…_dev`` (open TODO). Splitting on time
sidesteps that bug rather than depending on it — harmless today because both
Spaces run identical code, but it is why the sink is not the discriminator.
NOTE: this file must NOT be named ``test_*.py`` — pytest collects on filename
alone and a script that exits at import kills the whole run.
Usage::
# local run_logs
python scripts/analyze_post_analysis_reads.py --source local --dir ./run_logs
# the prod HF dataset (needs decouplerpy_results_token)
python scripts/analyze_post_analysis_reads.py --source hf
# dump the per-run table too
python scripts/analyze_post_analysis_reads.py --source hf --per-run
"""
from __future__ import annotations
import argparse
import json
import os
import sys
from datetime import datetime
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(REPO_ROOT / "src"))
from core.run_metrics import compute_run_metrics # noqa: E402
# 1ec068b reached prod on 2026-08-05. Runs at or after this instant are "after".
DEFAULT_CUTOFF = "2026-08-05T00:00:00"
# The one tool the fix actually touched. Every other terminal tool still returns
# only paths, so the re-read tail is EXPECTED to persist there — which is why
# results are grouped by terminal tool and never averaged into one number.
FIXED_TOOL = "dataset_compare_activity_by_group"
# --------------------------------------------------------------------------- #
# Record loading
# --------------------------------------------------------------------------- #
def load_local(directory: str) -> list[tuple[str, dict]]:
out = []
for path in sorted(Path(directory).glob("**/*.json")):
try:
with open(path, encoding="utf-8") as f:
record = json.load(f)
if is_run_record(record):
out.append((path.stem, record))
except Exception as exc: # noqa: BLE001
print(f" ! skipping {path}: {exc}", file=sys.stderr)
return out
def is_run_record(record: dict) -> bool:
"""Filter out the upload/audit records that share the same sink."""
return isinstance(record, dict) and "upload_id" not in record and "messages" in record
def load_hf(repo_id: str) -> list[tuple[str, dict]]:
from huggingface_hub import HfApi, hf_hub_download
token = os.environ.get("decouplerpy_results_token")
api = HfApi(token=token)
files = [
f
for f in api.list_repo_files(repo_id=repo_id, repo_type="dataset")
if f.endswith("/trace.json")
]
print(f"Found {len(files)} trace records in {repo_id}")
out = []
for remote in files:
try:
local = hf_hub_download(
repo_id=repo_id, filename=remote, repo_type="dataset", token=token
)
with open(local, encoding="utf-8") as f:
record = json.load(f)
if is_run_record(record):
out.append((remote.split("/")[-2], record))
except Exception as exc: # noqa: BLE001
print(f" ! skipping {remote}: {exc}", file=sys.stderr)
return out
# --------------------------------------------------------------------------- #
# Arm assignment
# --------------------------------------------------------------------------- #
def _parse_time(record: dict, run_id: str) -> datetime | None:
raw = record.get("execution_time")
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S.%f", "%Y-%m-%dT%H:%M:%S"):
try:
return datetime.strptime(str(raw), fmt)
except (TypeError, ValueError):
continue
# run_ids are stamped YYYYmmdd_HHMMSS by the Gradio path.
try:
return datetime.strptime(run_id[:15], "%Y%m%d_%H%M%S")
except ValueError:
return None
def summarize(records: list[tuple[str, dict]], cutoff: datetime) -> list[dict]:
rows = []
for run_id, record in records:
metrics = compute_run_metrics(record.get("messages") or [])
when = _parse_time(record, run_id)
rows.append(
{
"run_id": run_id,
"time": when,
"arm": None if when is None else ("after" if when >= cutoff else "before"),
**metrics,
}
)
return rows
def fisher(before_hits, before_n, after_hits, after_n):
"""Two-sided Fisher exact p for (≥1 re-read) before vs after."""
try:
from scipy.stats import fisher_exact
except ImportError:
return None
table = [
[before_hits, before_n - before_hits],
[after_hits, after_n - after_hits],
]
return float(fisher_exact(table)[1])
def report(rows: list[dict], per_run: bool) -> None:
measurable = [r for r in rows if r["n_post_analysis_reads"] is not None]
unmeasurable = [r for r in rows if r["n_post_analysis_reads"] is None]
print(f"\nRuns loaded: {len(rows)}")
print(f" measurable: {len(measurable)} (a terminal analysis tool wrote a results CSV)")
print(f" not measurable: {len(unmeasurable)} (excluded — NOT counted as zero re-reads)")
reasons: dict[str, int] = {}
for r in unmeasurable:
key = r.get("unmeasurable_reason") or "unknown"
reasons[key] = reasons.get(key, 0) + 1
for reason, n in sorted(reasons.items(), key=lambda kv: -kv[1]):
note = ""
if reason == "no_messages":
note = " ← the record persisted an empty conversation; nothing to count"
print(f" {reason:<28} {n}{note}")
no_arm = [r for r in measurable if r["arm"] is None]
if no_arm:
print(f" undated: {len(no_arm)} (excluded from the arms)")
tools = sorted({r["terminal_tool"] or "(unknown)" for r in measurable})
print("\nBy terminal tool — fraction of runs with ≥1 post-analysis re-read")
print(f"{'terminal tool':<45} {'before':>14} {'after':>14} {'fisher p':>10}")
print("-" * 87)
for tool in tools:
arm_counts = {}
for arm in ("before", "after"):
sel = [
r
for r in measurable
if (r["terminal_tool"] or "(unknown)") == tool and r["arm"] == arm
]
hits = sum(1 for r in sel if r["n_post_analysis_reads"] >= 1)
arm_counts[arm] = (hits, len(sel))
(bh, bn), (ah, an) = arm_counts["before"], arm_counts["after"]
def fmt(h, n):
return f"{h}/{n}" + (f" ({h / n:.0%})" if n else " (n=0)")
p = fisher(bh, bn, ah, an) if bn and an else None
p_str = f"{p:.4f}" if p is not None else "—"
mark = " ← the fixed tool" if tool == FIXED_TOOL else ""
print(f"{tool:<45} {fmt(bh, bn):>14} {fmt(ah, an):>14} {p_str:>10}{mark}")
print(
"\nOnly " + FIXED_TOOL + " was changed by 1ec068b; a persisting tail on the other\n"
"tools is expected, not a refutation. Do not average across rows."
)
if per_run:
print("\nPer-run detail")
for r in sorted(measurable, key=lambda r: r["time"] or datetime.min):
print(
f" {r['run_id']:<24} {str(r['arm']):<7} "
f"reads={r['n_post_analysis_reads']} "
f"after_steps={r['n_steps_after_analysis']} "
f"tool={r['terminal_tool']}"
)
def main() -> int:
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument("--source", choices=("local", "hf"), default="local")
ap.add_argument("--dir", default="./run_logs", help="local source directory")
ap.add_argument(
"--repo",
default=os.environ.get("LOG_SINK_HF_DATASET", "anne-voigt/decoupleRpy_results"),
help="HF dataset repo id (hf source)",
)
ap.add_argument("--cutoff", default=DEFAULT_CUTOFF, help="1ec068b deploy instant (ISO)")
ap.add_argument("--per-run", action="store_true")
args = ap.parse_args()
cutoff = datetime.fromisoformat(args.cutoff)
records = load_local(args.dir) if args.source == "local" else load_hf(args.repo)
if not records:
print("No run records found — nothing to report.", file=sys.stderr)
return 1
report(summarize(records, cutoff), args.per_run)
return 0
if __name__ == "__main__":
raise SystemExit(main())