#!/usr/bin/env python3 """Standard bridge fan-out test — ONE primary vs MANY workers, concurrently, for issue hunting through the running `fleetd` 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 fleet_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": "fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java"}, {"id": "worker", "probe": "WorkerService", "target": "fleetd/src/main/java/dev/ltms/fleet/worker/WorkerService.java"}, {"id": "rendezvous", "probe": "Rendezvous", "target": "fleetd/src/main/java/dev/ltms/fleet/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 fleet_reply with EXACTLY these four lines:\n" "1. {target}:\n" "2. issue: \n" "3. fix: \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 fleet_reply" if phase == "done" and has_reply: return "DEGRADED", f"resolved via {source} (worker did not call fleet_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()