2e138a199b
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.
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 `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()
|