diff --git a/e2e/README.md b/e2e/README.md new file mode 100644 index 0000000..21dcb84 --- /dev/null +++ b/e2e/README.md @@ -0,0 +1,77 @@ +# Bridge conversation test (`e2e/`) + +A standard, repeatable **live** end-to-end test of the two-way channel: it drives a real +multi-turn conversation between a primary and an off-subscription worker **through the +running `bridged` daemon**, captures the full transcript, and grades the channel. + +This is the committed form of the ad-hoc channel test that discovered the CB-115 gaps +(herdr `unknown` misclassification wedging delivery, dirty completion scrapes, and workers +never calling `bridge_reply` in conversation). Run it after any change to the injector, +status handling, completion/failure paths, or the worker reply charter. + +## What it exercises + +Each turn goes through the whole gateway exactly as a primary Opus session would — async +fire-and-poll (`POST /sessions/{id}/message {"wait":false}` → `GET /tasks/{ticket}`), so it +also validates the path that beats the caller's MCP timeout. It never sets +`ANTHROPIC_BASE_URL` and never talks to herdr directly, so it is **subscription-safe by +construction** — it only calls the bridge's loopback REST face. + +```mermaid +sequenceDiagram + participant T as conversation_test.py + participant B as bridged (REST) + participant W as worker (off-sub) + T->>B: POST /workers (spawn) + T->>B: GET /sessions/{id}/status (await ready) + loop each turn + T->>B: POST /sessions/{id}/message {wait:false} + B-->>T: ticket + B->>W: inject prompt (status-gated) + W-->>B: bridge_reply + T->>B: GET /tasks/{ticket} (poll) + B-->>T: done + reply + end + T->>B: DELETE /workers/{pane} (stop) +``` + +## Prerequisites + +- `bridged` is running (default REST on `http://127.0.0.1:8765`) with at least one worker + profile configured and its backend reachable. +- herdr is up (the daemon needs it). +- Python 3 (standard library only — no pip installs). + +## Run + +```bash +# spawn the default-profile worker, run the built-in 5-turn conversation, grade, clean up +python3 e2e/conversation_test.py + +# pick a profile / reuse a live worker / use your own prompts +python3 e2e/conversation_test.py --profile ollama +python3 e2e/conversation_test.py --tid term_abc123 --keep-worker +python3 e2e/conversation_test.py --prompts my_prompts.txt --out /tmp/run1 +``` + +A prompts file is one prompt per line; blank lines and `#` comments are ignored. + +## Output & grading + +- Writes `transcript.md` (in `--out`, default `e2e/`) — every turn's prompt, worker reply, + latency, resolution source, and observed status transitions, with an inline `> **GAP**` + note on any non-clean turn. +- Prints a per-turn line and an overall summary, and **exits non-zero** if any turn wedged, + failed, or returned empty — so it is CI-usable. + +Per-turn grade: + +| Grade | Meaning | +|------------|---------------------------------------------------------------------| +| `OK` | delivered and resolved by an explicit `bridge_reply` (`source=reply`) | +| `DEGRADED` | delivered and answered, but resolved via completion-scrape fallback | +| `EMPTY` | turn completed but the reply was empty | +| `FAILED` | the worker's turn ended in failure (`phase=failed`) | +| `WEDGE` | never resolved within the poll window (delivery wedge / lost turn) | + +`PASS` requires every turn to be `OK` or `DEGRADED`; a clean run is every turn `OK`. diff --git a/e2e/conversation_test.py b/e2e/conversation_test.py new file mode 100644 index 0000000..9c7647d --- /dev/null +++ b/e2e/conversation_test.py @@ -0,0 +1,227 @@ +#!/usr/bin/env python3 +"""Standard bridge conversation test — a multi-turn primary↔worker exchange through +the running `bridged` daemon, fully captured, with automatic gap analysis. + +This is the repeatable form of the ad-hoc channel test that surfaced the CB-115 gaps +(herdr `unknown` misclassification, dirty completion scrape, workers not calling +bridge_reply). It drives a real off-subscription worker over the live gateway exactly +as a primary Opus session would (async fire-and-poll), records every turn, and grades +the channel. + +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 conversation_test.py [--base URL] [--profile NAME] [--tid TERMINAL_ID] + [--prompts FILE] [--out DIR] [--poll-timeout SECS] + [--keep-worker] + + --base bridge REST base URL (default http://127.0.0.1:8765) + --profile worker profile to spawn (default: the daemon's default) + --tid reuse an existing worker (skips spawn + readiness wait) + --prompts newline-separated prompt file (default: the built-in script) + --out output dir for transcript.md (default: alongside this file) + --poll-timeout per-turn max wait, seconds (default 300) + --keep-worker do not stop a spawned worker at the end + +Exit code: 0 if every turn delivered AND produced a usable reply; 1 otherwise (so it +is CI-usable). A per-turn and overall gap report is printed to stdout. +""" +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 + +# A default conversation: short, varied turns (a question, a follow-up that needs the +# prior context, a tiny reasoning task, a meta-question, a close) — enough to exercise +# multi-turn delivery + reply on the channel without being a real coding workload. +DEFAULT_PROMPTS = [ + "Hi! Quick check that our channel works. In one sentence, what are you and what model are you running?", + "Thanks. Now a small task: what is 17 * 23? Show just the number.", + "Good. Remembering that result, is it a prime number? Answer yes or no with a one-line reason.", + "Switching topic: name one thing that would make this bridge conversation feel more reliable to you as the worker.", + "That's all — please acknowledge and we'll wrap up.", +] + + +def http(base, method, path, body=None, timeout=20): + 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: + return json.loads(r.read().decode()) + + +def now(): + return datetime.now().strftime("%H:%M:%S") + + +def spawn_worker(base, profile): + q = f"?profile={profile}" if profile else "" + res = http(base, "POST", f"/workers{q}") + tid = res.get("terminalId") or res.get("sessionId") + if not tid: + sys.exit(f"spawn failed: {res}") + pane = res.get("paneId") + print(f"[{now()}] spawned worker {tid} (profile={profile or 'default'}, pane={pane})") + return tid, pane + + +def await_ready(base, tid, timeout=120): + print(f"[{now()}] waiting for worker readiness (bridge MCP connect)…") + 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 — running anyway") + return False + + +def run_turn(base, tid, turn, prompt, poll_timeout): + """Fire one prompt async, poll the ticket to resolution, return a structured record.""" + t0 = time.time() + sent = http(base, "POST", f"/sessions/{tid}/message", {"wait": False, "content": prompt}) + ticket = sent.get("ticket") + samples = [] # (elapsed, phase, live_status) + reply = detail = source = phase = 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 = task.get("reply") + detail = task.get("detail") + source = task.get("replySource") + break + latency = round(time.time() - t0, 1) + + # compress status samples into a transition string + 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 { + "turn": turn, "prompt": prompt, "ticket": ticket, "phase": phase, + "source": source, "reply": reply, "detail": detail, "latency": latency, + "transitions": " → ".join(trans), "time": now(), + } + + +def grade(rec): + """Classify a turn's outcome. Returns (grade, note).""" + 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 write_transcript(out_dir, records): + path = out_dir / "transcript.md" + with path.open("w") as f: + f.write(f"# Bridge conversation test — {datetime.now():%Y-%m-%d %H:%M}\n\n") + for rec in records: + g, note = grade(rec) + f.write(f"### Turn {rec['turn']} — {rec['time']} " + f"(`{g}`, {rec['latency']}s, phase={rec['phase']}, source={rec['source']})\n\n") + f.write(f"**PRIMARY:** {rec['prompt']}\n\n") + f.write(f"**WORKER:** {rec['reply'] if rec['reply'] else '_(no reply)_ ' + str(rec['detail'])}\n\n") + f.write(f"_status: {rec['transitions']}_\n") + if g != "OK": + f.write(f"\n> **GAP — {g}:** {note}\n") + f.write("\n") + return path + + +def main(): + ap = argparse.ArgumentParser(description="Standard bridge conversation test") + ap.add_argument("--base", default="http://127.0.0.1:8765") + ap.add_argument("--profile", default=None) + ap.add_argument("--tid", default=None) + ap.add_argument("--prompts", default=None) + ap.add_argument("--out", default=str(HERE)) + ap.add_argument("--poll-timeout", type=int, default=300) + ap.add_argument("--keep-worker", action="store_true") + args = ap.parse_args() + + prompts = DEFAULT_PROMPTS + if args.prompts: + prompts = [ln.strip() for ln in pathlib.Path(args.prompts).read_text().splitlines() + if ln.strip() and not ln.startswith("#")] + + spawned = False + tid = args.tid + pane = None + if not tid: + tid, pane = spawn_worker(args.base, args.profile) + spawned = True + await_ready(args.base, tid) + + print(f"[{now()}] running {len(prompts)}-turn conversation on {tid}\n") + records = [] + for i, prompt in enumerate(prompts, 1): + rec = run_turn(args.base, tid, i, prompt, args.poll_timeout) + g, note = grade(rec) + records.append(rec) + print(f"[turn {i}] {g:8} {rec['latency']:6}s phase={rec['phase']} source={rec['source']}") + print(f" status: {rec['transitions']}") + print(f" reply: {(rec['reply'] or '(none) ' + str(rec['detail'])).strip()[:200]}") + print(f" note: {note}\n") + + out_dir = pathlib.Path(args.out) + out_dir.mkdir(parents=True, exist_ok=True) + path = write_transcript(out_dir, records) + + # ---- gap report ------------------------------------------------------- + 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("=" * 68) + print(f"CONVERSATION TEST SUMMARY — {len(records)} turns") + print(" " + " ".join(f"{g}:{n}" for g, n in counts.items())) + print(f" transcript: {path}") + ok = all(g in ("OK", "DEGRADED") for g in grades) + reply_clean = all(g == "OK" for g in grades) + if reply_clean: + print(" RESULT: PASS — every turn delivered and got a clean bridge_reply.") + elif ok: + print(" RESULT: PASS (with notes) — every turn delivered & replied, but some via fallback.") + else: + print(" RESULT: FAIL — one or more turns wedged, failed, or returned empty (see GAP notes).") + print("=" * 68) + + if spawned and pane and not args.keep_worker: + try: + http(args.base, "DELETE", f"/workers/{pane}") + print(f"[{now()}] stopped worker {tid} (pane {pane})") + except Exception as e: # noqa: BLE001 - best-effort cleanup + print(f"[{now()}] worker stop failed (ignore): {e}") + + sys.exit(0 if ok else 1) + + +if __name__ == "__main__": + main()