5f0ec034d9
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.
228 lines
9.5 KiB
Python
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()
|