CB-118: clip the completion baseline so the CB-115 guard survives >cap blocks
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
This commit is contained in:
@@ -0,0 +1,2 @@
|
||||
__pycache__/
|
||||
*.pyc
|
||||
@@ -0,0 +1,330 @@
|
||||
#!/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()
|
||||
Reference in New Issue
Block a user