From 426855e378eb7103e70dc8624151cc19977a112a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 16:08:21 +0200 Subject: [PATCH] e2e: live bridge_ask reverse-rendezvous harness (CB-205) Drives the reverse path end to end over REST loopback: a worker is delegated a task it cannot finish without asking, calls bridge_ask mid-turn, and the primary answers on the surfaced turnId so the worker resumes the SAME turn. Two blocking sends, no polling. Verified live: worker asked in ~9s, resumed and replied CHOSEN=BLUE via clean bridge_reply after the primary answered. Subscription-safe by construction (REST face only; never sets ANTHROPIC_BASE_URL). --- e2e/bridge_ask_test.py | 277 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 277 insertions(+) create mode 100644 e2e/bridge_ask_test.py diff --git a/e2e/bridge_ask_test.py b/e2e/bridge_ask_test.py new file mode 100644 index 0000000..20bbabd --- /dev/null +++ b/e2e/bridge_ask_test.py @@ -0,0 +1,277 @@ +#!/usr/bin/env python3 +"""Live bridge_ask test — the REVERSE rendezvous (CB-205), watched end to end. + +Every other harness drives the forward path: primary `bridge_send` → worker `bridge_reply`. +This drives the one that runs the other way. A worker is told to pause its delegated turn, +ask the primary a question via `bridge_ask`, and only finish once it has the answer — so the +turn round-trips primary→worker→primary→worker inside a SINGLE delegation. + +The mechanics that only this path exercises: + + • a worker's mid-turn question surfacing on the primary's *own* blocked send (Outcome.QUESTION), + • the `turnId` correlation that lets the primary answer the exact paused turn, + • the answer resuming that same turn and the worker's final `bridge_reply` landing on the + re-opened forward waiter (never a stale or cross-wired one). + +It is two blocking REST calls, no polling: + + 1. POST /sessions/{id}/message {content: TASK} → blocks, returns 202 {status:"question", + question, turnId} (the worker asked) + 2. POST /sessions/{id}/message {content: ANSWER, turnId} → blocks, returns 200 {reply, replySource} + (the worker resumed and replied) + +Like the rest of the suite it talks ONLY to the bridge's REST face on loopback — it never sets +ANTHROPIC_BASE_URL and never touches herdr, so it is subscription-safe by construction. + +Usage: + python3 bridge_ask_test.py [--base URL] [--profile NAME] [--repo DIR] [--out DIR] + [--send-timeout SECS] [--answer-timeout SECS] [--keep-worker] + + --base bridge REST base URL (default http://127.0.0.1:8765) + --profile profile for the worker (default: the daemon's default) + --repo cwd handed to the worker (default: the bridge repo root) + --out output dir for the transcript (default: alongside this file) + --send-timeout max wait for the worker to ASK (default 110) + --answer-timeout max wait for the worker to REPLY (default 90) + --keep-worker do not stop the spawned worker at the end + +Exit code: 0 if the worker asked, the answer resumed the turn, and the final reply reflected +the answer; 1 otherwise (CI-usable). +""" +import argparse +import json +import pathlib +import sys +import time +import urllib.error +import urllib.request +from datetime import datetime + +HERE = pathlib.Path(__file__).parent +REPO_ROOT = HERE.parent + +# The color the primary will hand back when the worker asks. The final reply must reflect it, +# uppercased — proof the answer actually reached the resumed turn (not a value the worker could +# have guessed: it is told to ask, and only the primary knows which of red/blue is chosen). +ANSWER_COLOR = "blue" + +# A task that CANNOT be completed without asking: the worker is not told which color to choose, +# only that the primary will name one when asked. So a correct final reply is only reachable by +# actually calling bridge_ask and using the answer. +TASK_PROMPT = ( + "You are a bridge worker in a quick coordination game. You do NOT know which color to pick — " + "only the primary does. Do exactly this, in order:\n" + "1. Call the `bridge_ask` tool with EXACTLY this question: \"PICK A COLOR: red or blue?\"\n" + "2. The primary will answer with one color word. Take that color and uppercase it.\n" + "3. Call `bridge_reply` with EXACTLY one line: CHOSEN= (e.g. CHOSEN=GREEN if told green).\n" + "Do not guess a color. Do not call bridge_reply before bridge_ask has returned an answer. " + "Do nothing else — no file reads, no other tools." +) + + +def http(base, method, path, body=None, timeout=20): + """JSON request. Tolerates an empty body (e.g. 204) → {}. Raises on non-2xx via urllib.""" + 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): + 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: {res}") + pane = res.get("paneId") + print(f"[{now()}] spawned {tid} (profile={profile or 'default'}, pane={pane})") + return tid, pane + + +def await_ready(base, tid, 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()}] worker ready (status={st.get('status')})") + return True + time.sleep(3) + print(f"[{now()}] WARNING: worker never reported ready within {timeout}s — sending anyway") + return False + + +def post_message(base, tid, body, timeout): + """One BLOCKING send. Returns (http_status, parsed_json). urllib raises on 4xx/5xx, so a + stale-turn 409 is surfaced here rather than swallowed.""" + data = json.dumps(body).encode() + req = urllib.request.Request(f"{base}/sessions/{tid}/message", data=data, method="POST", + headers={"Content-Type": "application/json"}) + try: + with urllib.request.urlopen(req, timeout=timeout) as r: + raw = r.read().decode().strip() + return r.status, (json.loads(raw) if raw else {}) + except urllib.error.HTTPError as e: + raw = e.read().decode().strip() + return e.code, (json.loads(raw) if raw else {}) + + +def run(base, profile, repo, send_timeout, answer_timeout): + """Drive the full reverse rendezvous. Returns a result record for grading + transcript.""" + rec = {"spawned": False, "tid": None, "pane": None, "phase": "spawn", + "question": None, "turnId": None, "reply": None, "replySource": None, + "detail": None, "ask_latency": None, "answer_latency": None} + + tid, pane = spawn_worker(base, profile, repo) + rec.update(tid=tid, pane=pane, spawned=True) + await_ready(base, tid) + + # 1) Delegate the ask-forcing task. This blocks until the worker calls bridge_ask, at which + # point our own send unblocks carrying the question and the turnId to answer on. + print(f"[{now()}] delegating task (blocks until the worker asks; up to {send_timeout}s)…") + t0 = time.time() + rec["phase"] = "awaiting_question" + try: + code, resp = post_message(base, tid, {"content": TASK_PROMPT, "timeoutMs": send_timeout * 1000}, + timeout=send_timeout + 15) + except Exception as e: # noqa: BLE001 + rec.update(phase="send_error", detail=f"delegating send failed: {e}") + return rec + rec["ask_latency"] = round(time.time() - t0, 1) + + if resp.get("status") != "question": + # The worker finished (or stalled) without asking — the whole point didn't happen. + rec.update(phase="no_question", detail=f"HTTP {code}: {json.dumps(resp)[:300]}", + reply=resp.get("reply"), replySource=resp.get("replySource")) + return rec + rec.update(phase="question", question=resp.get("question"), turnId=resp.get("turnId")) + print(f"[{now()}] worker ASKED ({rec['ask_latency']}s): {rec['question']!r} turnId={rec['turnId']}") + + if not rec["turnId"]: + rec.update(phase="no_turnid", detail="question surfaced without a turnId to answer on") + return rec + + # 2) Answer on that exact turn. This blocks again until the resumed worker calls bridge_reply. + print(f"[{now()}] answering '{ANSWER_COLOR}' on turn {rec['turnId']} (blocks until reply; up to {answer_timeout}s)…") + t1 = time.time() + rec["phase"] = "awaiting_reply" + try: + code, resp = post_message(base, tid, + {"content": ANSWER_COLOR, "turnId": rec["turnId"], + "timeoutMs": answer_timeout * 1000}, + timeout=answer_timeout + 15) + except Exception as e: # noqa: BLE001 + rec.update(phase="answer_error", detail=f"answer send failed: {e}") + return rec + rec["answer_latency"] = round(time.time() - t1, 1) + + if code == 409 or resp.get("error") == "stale_turn": + rec.update(phase="stale_turn", detail=f"HTTP {code}: {json.dumps(resp)[:300]}") + return rec + if resp.get("reply") is None: + rec.update(phase="no_reply", detail=f"HTTP {code}: {json.dumps(resp)[:300]}") + return rec + rec.update(phase="replied", reply=resp.get("reply"), replySource=resp.get("replySource")) + print(f"[{now()}] worker RESUMED and replied ({rec['answer_latency']}s): {rec['reply']!r} " + f"(source={rec['replySource']})") + return rec + + +def grade(rec): + """PASS only if the worker asked, the turn resumed, and the reply reflects the answer.""" + if rec["phase"] == "no_question": + return "NO_ASK", "the worker finished/stalled without ever calling bridge_ask" + if rec["phase"] in ("send_error", "answer_error", "spawn"): + return "ERROR", rec.get("detail") or "transport error before the round-trip completed" + if rec["phase"] == "no_turnid": + return "NO_TURNID", "the question surfaced without a turnId — the primary could not answer" + if rec["phase"] == "stale_turn": + return "STALE", "answering the turn was rejected as stale (it lapsed or was already answered)" + if rec["phase"] in ("awaiting_reply", "no_reply"): + return "NO_RESUME", "the worker asked but never resumed to a final reply within the window" + if rec["phase"] == "replied": + reflected = ANSWER_COLOR.upper() in (rec["reply"] or "").upper() + if reflected and rec["replySource"] == "reply": + return "OK", "asked, resumed the same turn, and the reply reflected the primary's answer" + if reflected: + return "DEGRADED", f"reply reflected the answer but resolved via {rec['replySource']} " \ + "(worker did not call bridge_reply cleanly)" + return "WRONG_ANSWER", f"the worker replied but did not reflect '{ANSWER_COLOR}' — " \ + f"the answer may not have reached the resumed turn: {rec['reply']!r}" + return "WEDGE", f"unexpected terminal phase {rec['phase']}: {rec.get('detail')}" + + +def write_transcript(out_dir, rec, meta): + path = out_dir / "bridge_ask_transcript.md" + g, note = grade(rec) + with path.open("w") as f: + f.write(f"# Live bridge_ask — reverse rendezvous — {datetime.now():%Y-%m-%d %H:%M}\n\n") + f.write(f"One worker paused its delegated turn to ask the primary, then resumed with the " + f"answer (profile `{meta['profile']}`). Result: **`{g}`**.\n\n") + f.write("## Round-trip\n\n") + f.write(f"1. **primary → worker** (delegation): the ask-forcing task.\n") + f.write(f"2. **worker → primary** (`bridge_ask`, {rec.get('ask_latency')}s): " + f"{rec.get('question')!r} — surfaced on the primary's blocked send as a " + f"`question` with `turnId={rec.get('turnId')}`.\n") + f.write(f"3. **primary → worker** (answer on that turn): `{ANSWER_COLOR}`.\n") + f.write(f"4. **worker → primary** (`bridge_reply`, {rec.get('answer_latency')}s, " + f"source={rec.get('replySource')}): {rec.get('reply')!r}\n\n") + f.write(f"> **{g}:** {note}\n") + if rec.get("detail"): + f.write(f">\n> detail: {rec['detail']}\n") + return path + + +def main(): + ap = argparse.ArgumentParser(description="Live bridge_ask reverse-rendezvous test (CB-205)") + 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("--send-timeout", type=int, default=110, help="max wait for the worker to ASK") + ap.add_argument("--answer-timeout", type=int, default=90, help="max wait for the worker to REPLY") + ap.add_argument("--keep-worker", action="store_true") + args = ap.parse_args() + + print(f"[{now()}] live bridge_ask: 1 primary, 1 worker " + f"(profile={args.profile or 'default'}, repo={args.repo})\n") + + rec = {"spawned": False, "pane": None} + try: + rec = run(args.base, args.profile, args.repo, args.send_timeout, args.answer_timeout) + finally: + if not args.keep_worker and rec.get("spawned") and rec.get("pane"): + try: + http(args.base, "DELETE", f"/workers/{rec['pane']}") + print(f"[{now()}] stopped worker (pane {rec['pane']})") + except Exception as e: # noqa: BLE001 + print(f"[{now()}] stop failed (ignore): {e}") + + g, note = grade(rec) + out_dir = pathlib.Path(args.out) + out_dir.mkdir(parents=True, exist_ok=True) + path = write_transcript(out_dir, rec, {"profile": args.profile or "default"}) + + print() + print("=" * 72) + print("LIVE bridge_ask SUMMARY — reverse rendezvous (CB-205)") + print(f" asked: {rec.get('question')!r} (turnId={rec.get('turnId')}, {rec.get('ask_latency')}s)") + print(f" answered: {ANSWER_COLOR!r}") + print(f" replied: {rec.get('reply')!r} (source={rec.get('replySource')}, {rec.get('answer_latency')}s)") + print(f" transcript: {path}") + print(f" RESULT: {g} — {note}") + print("=" * 72) + + sys.exit(0 if g == "OK" else 1) + + +if __name__ == "__main__": + main()