e2e: standard bridge conversation test harness
A repeatable multi-turn primary↔worker conversation driven entirely through the bridge's loopback REST face (async fire-and-poll) — never sets ANTHROPIC_BASE_URL and never touches herdr, so it is subscription-safe by construction. Records every turn to a transcript, grades each (OK / DEGRADED / EMPTY / FAILED / WEDGE), and exits non-zero if any turn fails to deliver-and-reply, so it is CI-usable. This is the repeatable form of the ad-hoc channel test that surfaced the CB-115 and CB-116 gaps.
This commit is contained in:
@@ -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`.
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user