From 2a61fe69f1db6c54d2968cb130278b9c354048d2 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 09:11:20 +0200 Subject: [PATCH] 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 --- .../bridged/inject/CompletionResolver.java | 11 +- .../inject/CompletionResolverTest.java | 25 ++ e2e/.gitignore | 2 + e2e/issue_hunt_test.py | 330 ++++++++++++++++++ 4 files changed, 367 insertions(+), 1 deletion(-) create mode 100644 e2e/.gitignore create mode 100644 e2e/issue_hunt_test.py diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 8c3e0b1..2c815b4 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -97,7 +97,11 @@ public final class CompletionResolver implements TurnListener { } String baseline; try { - baseline = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); + // Clip to the same cap resolve() applies to the tail (line ~134): the CB-115 misattribution + // guard compares baseline.equals(tail), so both sides must be the same capped representation. + // An unclipped baseline vs a clipped tail would never match for a >MAX_SCRAPE_CHARS block, + // defeating the guard and letting a stale completion resolve the send. + baseline = clip(lastAssistantBlock(agents.read(target, SCRAPE_SOURCE))); } catch (RuntimeException e) { baseline = null; // fail open: no baseline ⇒ no suppression log.debug("delivery baseline for {} failed: {}", target, e.getMessage()); @@ -105,6 +109,11 @@ public final class CompletionResolver implements TurnListener { inFlight.put(target, new InFlight(waiter, baseline)); } + /** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */ + InFlight inFlight(String target) { + return inFlight.get(target); + } + @Override public void onTurnComplete(String target) { // Read the in-flight turn on the poller thread — before any next-turn delivery can overwrite diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index f919e63..17a27f8 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -156,6 +156,31 @@ class CompletionResolverTest { assertEquals("No, 391 = 17 × 23.", waiter.getNow(null).text()); } + @Test + void suppressesAnUnchangedCompletionEvenWhenTheBlockExceedsTheScrapeCap() { + // The fan-out issue-hunt finding: captureBaseline once stored the RAW (unclipped) assistant + // block while resolve compares against a clip()'d tail. For a block longer than MAX_SCRAPE_CHARS + // the two capped representations differ even when the pane never changed, so the CB-115 + // byte-identical guard failed to fire and a stale completion could resolve the send. Both sides + // must clip identically; here an unchanged >cap block on rapid back-to-back turns stays suppressed. + String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(longBlock); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + var waiter = rendezvous.open("term_a"); // a send is blocked on this turn + resolver.captureBaseline("term_a"); // baseline is the clipped >cap block + var turn = resolver.inFlight("term_a"); + assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(), + "the delivery baseline is clipped to the same cap resolve() applies to the tail"); + + resolver.resolve("term_a", turn); // scrape unchanged → clipped tail == baseline → suppress + + assertFalse(waiter.isDone(), + "an unchanged >cap block must still be recognised as stale and suppressed"); + assertTrue(rendezvous.isWaiting("term_a"), "the send stays waiting for a real reply"); + } + @Test void resolvesWhenThereIsNoBaseline() { // No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves. diff --git a/e2e/.gitignore b/e2e/.gitignore new file mode 100644 index 0000000..7a60b85 --- /dev/null +++ b/e2e/.gitignore @@ -0,0 +1,2 @@ +__pycache__/ +*.pyc diff --git a/e2e/issue_hunt_test.py b/e2e/issue_hunt_test.py new file mode 100644 index 0000000..9b706cc --- /dev/null +++ b/e2e/issue_hunt_test.py @@ -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}:\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 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()