Files
fleetd/e2e/conversation_test.py
Dai Ha 5f0ec034d9 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.
2026-07-16 08:27:50 +02:00

228 lines
9.5 KiB
Python

#!/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()