Files
fleetd/e2e/conversation_test.py
Dai Ha 2e138a199b CB-634: one shared "fleet" workspace + rename bridged -> fleetd cutover
Two changes ship together here.

1. One shared herdr workspace. The lead and every worker now live in one
   workspace called "fleet", so the operator sees one "session" with many
   windows, not two. Before, the lead sat in a "leads" workspace and workers
   in "bridged-workers", which read as two sessions. The lead is still told
   apart from workers by its exact tab label ("lead: <name>"), so putting them
   in one space is safe. LeadTabScanner keeps the exclude-by-label mechanism
   for split layouts; Fleetd now passes an empty exclude set.

2. Rename the daemon from "bridged" to "fleetd" (the binary, config, scripts,
   launchd/systemd units, module dir, and MCP mount).
   - Module dir bridged/ -> fleetd/; jar finalName -> fleetd.jar.
   - Log line, comments, docs, and CLAUDE.md updated to say fleetd.
   - Scripts renamed: redeploy-bridged.sh -> redeploy-fleetd.sh,
     bridged-launchd-wrapper.sh -> fleetd-launchd-wrapper.sh.
   - Deploy units renamed: dev.ltms.bridged.plist -> dev.ltms.fleetd.plist,
     bridged.service -> fleetd.service; launchd Label -> dev.ltms.fleetd.
   - Config default bridged.yaml -> fleetd.yaml; the legacy bridged.yaml is
     still read as a fallback, and still gitignored.
   - MCP: drop the deprecated bridge_* tool twins; only fleet_* remain. The
     server name is "fleet". The mount name in the local .mcp.json becomes
     "fleet" (gitignored, not in this commit).
   - Env var defaults BRIDGED_API_TOKEN -> FLEETD_API_TOKEN, fixture
     BRIDGED_WORKER_TOKEN -> FLEETD_WORKER_TOKEN.

Kept on purpose: the BRIDGED_MEMBER marker. Renaming it is a coupled change to
the credential-scrub security control (an operator secrets.sh may guard on it),
so it stays until that migration is done on its own.

Metrics were already fleet_* (CB-632); MetricNamesTest still guards that no
name says bridged_.

The canonical CLAUDE.md block and the wiki template stay byte-identical
(wiki working tree edited, committed to the wiki repo separately).

949 tests pass (mvn clean install). 4 fewer than before = the 4 removed
bridge_* alias tests.
2026-08-25 04:01:08 +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 `fleetd` 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
fleet_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 fleet_reply"
if phase == "done" and has_reply:
return "DEGRADED", f"resolved via {source} (worker did not call fleet_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 fleet_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()