2a61fe69f1
captureBaseline stored the raw, unclipped last-assistant block while resolve() compares against clip(...) capped at MAX_SCRAPE_CHARS. For a block longer than 4000 chars the two capped representations never match even when the pane is unchanged, defeating the CB-115 misattribution guard and letting a stale completion resolve a rapid back-to-back send. Clip the baseline identically. Regression test: an unchanged >cap block stays suppressed. Surfaced by the fan-out issue-hunt E2E (1 primary -> 3 concurrent workers, e2e/issue_hunt_test.py, added here). The same hunt's WorkerService.stop() and Rendezvous.complete() findings were verified as false positives (locatePane is already guarded; the sender's finally-close already removes the waiter). Closes #2
331 lines
15 KiB
Python
331 lines
15 KiB
Python
#!/usr/bin/env python3
|
|
"""Standard bridge fan-out test — ONE primary vs MANY workers, concurrently, for
|
|
issue hunting through the running `bridged` daemon, fully captured, with gap analysis.
|
|
|
|
Where conversation_test.py exercises a single worker over multiple turns, this drives
|
|
the path that only appears under fan-out: the primary spawns N workers, sends each a
|
|
distinct issue-hunting assignment on a slice of the codebase, fires them all at once,
|
|
and collects every reply concurrently. That stresses what a single worker never can —
|
|
|
|
• simultaneous delivery to many panes (the injector's per-worker, not global, writer),
|
|
• per-session rendezvous isolation (N blocked sends resolving independently),
|
|
• reply routing under concurrency (worker A's answer must never resolve worker B's send),
|
|
|
|
and, as the payload, whether a fleet of off-subscription workers can actually surface
|
|
real issues in the repo and report them back structurally via bridge_reply.
|
|
|
|
It talks ONLY to the bridge's REST face on loopback — it never sets ANTHROPIC_BASE_URL
|
|
and never touches herdr directly, so it is subscription-safe by construction.
|
|
|
|
Usage:
|
|
python3 issue_hunt_test.py [--base URL] [--profile NAME] [--repo DIR]
|
|
[--out DIR] [--poll-timeout SECS] [--keep-workers]
|
|
[--assignments FILE]
|
|
|
|
--base bridge REST base URL (default http://127.0.0.1:8765)
|
|
--profile profile for every worker (default: the daemon's default)
|
|
--repo cwd handed to each worker (default: the bridge repo root)
|
|
--out output dir for the transcript (default: alongside this file)
|
|
--poll-timeout per-worker max wait, seconds (default 300)
|
|
--keep-workers do not stop spawned workers at the end
|
|
--assignments JSON file overriding the built-in assignment list
|
|
|
|
Exit code: 0 if every worker delivered AND produced a usable reply with no cross-talk;
|
|
1 otherwise (CI-usable). A per-worker and fleet-level gap report is printed to stdout.
|
|
"""
|
|
import argparse
|
|
import json
|
|
import pathlib
|
|
import sys
|
|
import threading
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from datetime import datetime
|
|
|
|
HERE = pathlib.Path(__file__).parent
|
|
REPO_ROOT = HERE.parent
|
|
|
|
# Each worker gets a distinct source file to hunt in, plus a `probe` — a token its reply
|
|
# should mention if it actually addressed ITS assignment (a soft cross-talk detector: a
|
|
# reply that references only another worker's file is a routing red flag). Targets are the
|
|
# hot files this project has been iterating on, so a real issue is plausible to find.
|
|
DEFAULT_ASSIGNMENTS = [
|
|
{"id": "completion", "probe": "CompletionResolver",
|
|
"target": "bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java"},
|
|
{"id": "worker", "probe": "WorkerService",
|
|
"target": "bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java"},
|
|
{"id": "rendezvous", "probe": "Rendezvous",
|
|
"target": "bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java"},
|
|
]
|
|
|
|
PROMPT_TMPL = (
|
|
"You are one of several issue-hunting workers in the claude-bridge repo (it is your "
|
|
"current working directory). Your assignment: inspect the file `{target}` and find the "
|
|
"SINGLE most important real bug, correctness gap, or risk in it. Read the file before "
|
|
"answering. Reply via bridge_reply with EXACTLY these four lines:\n"
|
|
"1. {target}:<line>\n"
|
|
"2. issue: <one sentence>\n"
|
|
"3. fix: <one line>\n"
|
|
"4. severity: high|medium|low\n"
|
|
"Keep it under 90 words. If after reading you find nothing real, reply 'NO ISSUE' and one "
|
|
"line why. Do NOT hunt in any other file — only `{target}`."
|
|
)
|
|
|
|
|
|
def http(base, method, path, body=None, timeout=20):
|
|
"""JSON request. Tolerates an empty body (e.g. 204 No Content on DELETE) → returns {}."""
|
|
data = json.dumps(body).encode() if body is not None else None
|
|
req = urllib.request.Request(base + path, data=data, method=method,
|
|
headers={"Content-Type": "application/json"})
|
|
with urllib.request.urlopen(req, timeout=timeout) as r:
|
|
raw = r.read().decode().strip()
|
|
return json.loads(raw) if raw else {}
|
|
|
|
|
|
def now():
|
|
return datetime.now().strftime("%H:%M:%S")
|
|
|
|
|
|
def spawn_worker(base, profile, repo, wid):
|
|
res = http(base, "POST", "/workers", {"profile": profile, "cwd": repo} if profile
|
|
else {"cwd": repo})
|
|
tid = res.get("terminalId") or res.get("sessionId")
|
|
if not tid:
|
|
raise RuntimeError(f"spawn failed for {wid}: {res}")
|
|
pane = res.get("paneId")
|
|
print(f"[{now()}] [{wid}] spawned {tid} (profile={profile or 'default'}, pane={pane})")
|
|
return tid, pane
|
|
|
|
|
|
def await_ready(base, tid, wid, timeout=150):
|
|
deadline = time.time() + timeout
|
|
while time.time() < deadline:
|
|
try:
|
|
st = http(base, "GET", f"/sessions/{tid}/status")
|
|
except urllib.error.URLError:
|
|
st = {}
|
|
if st.get("ready"):
|
|
print(f"[{now()}] [{wid}] ready (status={st.get('status')})")
|
|
return True
|
|
time.sleep(3)
|
|
print(f"[{now()}] [{wid}] WARNING: never reported ready within {timeout}s — sending anyway")
|
|
return False
|
|
|
|
|
|
def fire(base, tid, prompt):
|
|
"""Fire one async send; return its ticket (delivery is confirmed by a ticket coming back)."""
|
|
sent = http(base, "POST", f"/sessions/{tid}/message", {"wait": False, "content": prompt})
|
|
return sent.get("ticket")
|
|
|
|
|
|
def poll(base, ticket, wid, t0, poll_timeout):
|
|
"""Poll a ticket to resolution; return (record fields) mirroring conversation_test."""
|
|
samples, reply, detail, source, phase = [], None, None, None, None
|
|
deadline = time.time() + poll_timeout
|
|
while time.time() < deadline:
|
|
time.sleep(3)
|
|
task = http(base, "GET", f"/tasks/{ticket}")
|
|
phase = task.get("phase")
|
|
live = (task.get("detail") or "").replace("worker ", "") if phase == "pending" else ""
|
|
samples.append((round(time.time() - t0, 1), phase, live))
|
|
if phase in ("done", "failed"):
|
|
reply, detail, source = task.get("reply"), task.get("detail"), task.get("replySource")
|
|
break
|
|
trans, last = [], None
|
|
for el, ph, live in samples:
|
|
tag = live if ph == "pending" else ph
|
|
if tag != last:
|
|
trans.append(f"{tag}@{el}s")
|
|
last = tag
|
|
return {"phase": phase, "reply": reply, "detail": detail, "source": source,
|
|
"latency": round(time.time() - t0, 1), "transitions": " → ".join(trans)}
|
|
|
|
|
|
def grade(rec):
|
|
phase, source, reply = rec["phase"], rec["source"], rec["reply"]
|
|
has_reply = bool(reply and reply.strip())
|
|
if phase == "done" and source == "reply" and has_reply:
|
|
return "OK", "clean explicit bridge_reply"
|
|
if phase == "done" and has_reply:
|
|
return "DEGRADED", f"resolved via {source} (worker did not call bridge_reply)"
|
|
if phase == "done" and not has_reply:
|
|
return "EMPTY", "turn completed but reply was empty"
|
|
if phase == "failed":
|
|
return "FAILED", f"worker turn failed: {(rec['detail'] or '').strip()[:120]}"
|
|
return "WEDGE", "never resolved within the poll window (delivery wedge or lost turn)"
|
|
|
|
|
|
def worker_lifecycle(base, profile, repo, poll_timeout, a, barrier):
|
|
"""Full per-worker path: spawn → ready → (barrier) → fire → poll. Spawn+ready run
|
|
concurrently across workers; the send waits on the shared `barrier` so every ready
|
|
worker fires within the same instant — the real simultaneous-delivery stress. A worker
|
|
that fails to spawn aborts the barrier so the rest don't block forever."""
|
|
wid = a["id"]
|
|
rec = {"id": wid, "target": a["target"], "probe": a["probe"], "spawned": False,
|
|
"tid": None, "pane": None, "ticket": None}
|
|
try:
|
|
tid, pane = spawn_worker(base, profile, repo, wid)
|
|
rec.update(tid=tid, pane=pane, spawned=True)
|
|
await_ready(base, tid, wid)
|
|
except Exception as e: # noqa: BLE001
|
|
barrier.abort() # release peers waiting on the barrier
|
|
rec.update(phase="failed", detail=f"spawn/ready error: {e}", reply=None,
|
|
source=None, latency=0.0, transitions="")
|
|
return rec
|
|
|
|
try:
|
|
barrier.wait(timeout=210) # all ready workers proceed together
|
|
except (threading.BrokenBarrierError, Exception): # noqa: BLE001
|
|
pass # a peer died or timed out — fire anyway rather than hang
|
|
t0 = time.time()
|
|
prompt = PROMPT_TMPL.format(target=a["target"])
|
|
try:
|
|
ticket = fire(base, tid, prompt)
|
|
rec["ticket"] = ticket
|
|
print(f"[{now()}] [{wid}] fired (ticket={ticket})")
|
|
rec.update(poll(base, ticket, wid, t0, poll_timeout))
|
|
except Exception as e: # noqa: BLE001
|
|
rec.update(phase="failed", detail=f"send error: {e}", reply=None,
|
|
source=None, latency=round(time.time() - t0, 1), transitions="")
|
|
return rec
|
|
|
|
|
|
def crosstalk_report(records):
|
|
"""Fleet-level isolation checks: distinct tickets, distinct non-empty replies, and each
|
|
reply addressing its OWN assigned file (probe token present). Returns (list_of_gaps)."""
|
|
gaps = []
|
|
tickets = [r.get("ticket") for r in records if r.get("ticket")]
|
|
if len(tickets) != len(set(tickets)):
|
|
gaps.append("ticket collision: two workers were handed the same ticket id")
|
|
replies = {r["id"]: (r.get("reply") or "").strip() for r in records}
|
|
# identical non-empty replies from distinct assignments ⇒ suspected reply misrouting
|
|
seen = {}
|
|
for wid, text in replies.items():
|
|
if text and text in seen:
|
|
gaps.append(f"identical reply from '{seen[text]}' and '{wid}' "
|
|
f"(distinct assignments should not yield byte-identical answers)")
|
|
elif text:
|
|
seen[text] = wid
|
|
# a reply that names ANOTHER worker's file but not its own ⇒ likely cross-routing
|
|
for r in records:
|
|
text = (r.get("reply") or "")
|
|
if not text.strip():
|
|
continue
|
|
own = r["probe"] in text or pathlib.Path(r["target"]).name in text
|
|
others = [o["probe"] for o in records if o["id"] != r["id"] and o["probe"] in text]
|
|
if not own and others:
|
|
gaps.append(f"worker '{r['id']}' (assigned {r['probe']}) replied about "
|
|
f"{others} but not its own file — possible cross-routing")
|
|
return gaps
|
|
|
|
|
|
def write_transcript(out_dir, records, meta):
|
|
path = out_dir / "issue_hunt_transcript.md"
|
|
with path.open("w") as f:
|
|
f.write(f"# Bridge fan-out issue-hunt — {datetime.now():%Y-%m-%d %H:%M}\n\n")
|
|
f.write(f"One primary vs **{len(records)} concurrent workers** "
|
|
f"(profile `{meta['profile']}`), each hunting a distinct file.\n\n")
|
|
for r in records:
|
|
g, note = grade(r)
|
|
f.write(f"### `{r['id']}` — {r['target']} (`{g}`, {r.get('latency')}s, "
|
|
f"phase={r.get('phase')}, source={r.get('source')})\n\n")
|
|
f.write(f"**ASSIGNMENT:** find the top issue in `{r['target']}`\n\n")
|
|
reply = r.get("reply")
|
|
f.write(f"**WORKER {r['id']}:** {reply if reply else '_(no reply)_ ' + str(r.get('detail'))}\n\n")
|
|
if r.get("transitions"):
|
|
f.write(f"_status: {r['transitions']}_\n")
|
|
if g != "OK":
|
|
f.write(f"\n> **GAP — {g}:** {note}\n")
|
|
f.write("\n")
|
|
if meta["gaps"]:
|
|
f.write("## Fleet-level gaps\n\n")
|
|
for gp in meta["gaps"]:
|
|
f.write(f"- {gp}\n")
|
|
return path
|
|
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description="Bridge fan-out issue-hunt test (1 primary, N workers)")
|
|
ap.add_argument("--base", default="http://127.0.0.1:8765")
|
|
ap.add_argument("--profile", default=None)
|
|
ap.add_argument("--repo", default=str(REPO_ROOT))
|
|
ap.add_argument("--out", default=str(HERE))
|
|
ap.add_argument("--poll-timeout", type=int, default=300)
|
|
ap.add_argument("--keep-workers", action="store_true")
|
|
ap.add_argument("--assignments", default=None)
|
|
args = ap.parse_args()
|
|
|
|
assignments = DEFAULT_ASSIGNMENTS
|
|
if args.assignments:
|
|
assignments = json.loads(pathlib.Path(args.assignments).read_text())
|
|
|
|
n = len(assignments)
|
|
barrier = threading.Barrier(n) # releases exactly when all n ready workers reach it
|
|
print(f"[{now()}] fan-out issue-hunt: 1 primary vs {n} workers "
|
|
f"(profile={args.profile or 'default'}, repo={args.repo})\n")
|
|
|
|
# Spawn + ready + fire + poll all workers concurrently; the barrier makes every send fire
|
|
# together once all are ready, so delivery pressure hits the daemon simultaneously.
|
|
records = []
|
|
with ThreadPoolExecutor(max_workers=n) as ex:
|
|
futures = [ex.submit(worker_lifecycle, args.base, args.profile, args.repo,
|
|
args.poll_timeout, a, barrier) for a in assignments]
|
|
for fut in futures:
|
|
records.append(fut.result())
|
|
|
|
records.sort(key=lambda r: [a["id"] for a in assignments].index(r["id"]))
|
|
|
|
print()
|
|
for r in records:
|
|
g, note = grade(r)
|
|
print(f"[{r['id']:11}] {g:8} {str(r.get('latency','?')):6}s "
|
|
f"phase={r.get('phase')} source={r.get('source')}")
|
|
print(f" status: {r.get('transitions') or '(none)'}")
|
|
print(f" reply: {((r.get('reply') or '(none) ' + str(r.get('detail'))).strip()[:200])}")
|
|
print(f" note: {note}\n")
|
|
|
|
gaps = crosstalk_report(records)
|
|
meta = {"profile": args.profile or "default", "gaps": gaps}
|
|
out_dir = pathlib.Path(args.out)
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
path = write_transcript(out_dir, records, meta)
|
|
|
|
grades = [grade(r)[0] for r in records]
|
|
counts = {g: grades.count(g) for g in ("OK", "DEGRADED", "EMPTY", "FAILED", "WEDGE") if grades.count(g)}
|
|
print("=" * 72)
|
|
print(f"FAN-OUT ISSUE-HUNT SUMMARY — 1 primary vs {n} workers")
|
|
print(" channel: " + " ".join(f"{g}:{v}" for g, v in counts.items()))
|
|
print(f" transcript: {path}")
|
|
if gaps:
|
|
print(" FLEET GAPS:")
|
|
for gp in gaps:
|
|
print(f" ⚠ {gp}")
|
|
else:
|
|
print(" isolation: clean — distinct tickets, distinct replies, each on its own file")
|
|
delivered = all(r.get("ticket") for r in records)
|
|
replied = all(g in ("OK", "DEGRADED") for g in grades)
|
|
if delivered and replied and not gaps:
|
|
print(" RESULT: PASS — all workers delivered concurrently, replied, and stayed isolated.")
|
|
elif delivered and replied:
|
|
print(" RESULT: PASS (with notes) — all delivered & replied, but see FLEET GAPS.")
|
|
else:
|
|
print(" RESULT: FAIL — a worker did not deliver or did not reply (see GAP notes).")
|
|
print("=" * 72)
|
|
|
|
if not args.keep_workers:
|
|
for r in records:
|
|
if r.get("spawned") and r.get("pane"):
|
|
try:
|
|
http(args.base, "DELETE", f"/workers/{r['pane']}")
|
|
print(f"[{now()}] [{r['id']}] stopped (pane {r['pane']})")
|
|
except Exception as e: # noqa: BLE001
|
|
print(f"[{now()}] [{r['id']}] stop failed (ignore): {e}")
|
|
|
|
sys.exit(0 if (delivered and replied and not gaps) else 1)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|