Files
fleetd/e2e/conversation_sustained_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

171 lines
7.6 KiB
Python

#!/usr/bin/env python3
"""Sustained back-and-forth bridge test — ONE primary, ONE worker, many dependent turns
over a fixed wall-clock window (default 5 minutes), through the running `fleetd` daemon.
Where conversation_test.py proves a handful of turns work and issue_hunt_test.py proves
fan-out isolation, this proves the channel stays healthy under a *sustained, stateful*
conversation: a running-total game the worker must keep in its head across turns. Turn N's
prompt does NOT restate the total — the worker has to remember it from turn N-1 — so a
correct answer is evidence of genuine multi-turn continuity, not just per-turn liveness.
Each reply is machine-checked against the primary's own expected total; on a drift the
primary re-anchors (states the correct total once) and keeps going, and drift is reported.
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_sustained_test.py [--base URL] [--profile NAME] [--repo DIR]
[--duration SECS] [--turn-timeout SECS] [--keep-worker]
--base bridge REST base URL (default http://127.0.0.1:8765)
--profile worker profile (default: the daemon's default)
--repo cwd handed to the worker (default: the bridge repo root)
--duration wall-clock window, seconds (default 300 = 5 minutes)
--turn-timeout per-turn max wait, seconds (default 150)
--keep-worker do not stop the worker at the end
Exit code: 0 if every turn in the window resolved via a clean fleet_reply with no channel
break; 1 otherwise. A live per-turn log streams to stdout so the run can be watched.
"""
import argparse
import pathlib
import re
import sys
import time
# Reuse the exact REST primitives the other harnesses use (same daemon contract).
sys.path.insert(0, str(pathlib.Path(__file__).parent))
from issue_hunt_test import http, now, spawn_worker, await_ready, fire, poll # noqa: E402
HERE = pathlib.Path(__file__).parent
REPO_ROOT = HERE.parent
# The per-turn increments, cycled. Non-trivial and varied so the running total isn't a
# predictable multiple the worker could pattern-match without actually tracking it.
STEPS = [7, 3, 11, 5, 9, 4, 13, 6, 8, 2]
RULES = (
"Let's play a running-total game across several messages. The total starts at 0. "
"In each message I'll tell you to add a number; keep the running total yourself and "
"reply via fleet_reply with ONLY the current total as a plain integer — no words, no "
"punctuation, just the number. Do not restate the arithmetic. First move: add {step}."
)
NEXT = ("Add {step}. Reply via fleet_reply with only the new running total.")
REANCHOR = ("Let's re-sync — the running total is {total}. Now add {step}. Reply via "
"fleet_reply with only the new running total.")
def parse_int(reply):
"""Pull the worker's answer integer from its reply (last integer token wins)."""
if not reply:
return None
nums = re.findall(r"-?\d+", reply.replace(",", ""))
return int(nums[-1]) if nums else None
def one_turn(base, tid, prompt, turn_timeout):
"""Fire one prompt and block on its reply. Returns the poll record."""
t0 = time.time()
ticket = fire(base, tid, prompt)
rec = poll(base, ticket, "worker", t0, turn_timeout)
rec["ticket"] = ticket
return rec
def main():
ap = argparse.ArgumentParser(description="Sustained back-and-forth bridge test (1 primary, 1 worker)")
ap.add_argument("--base", default="http://127.0.0.1:8765")
ap.add_argument("--profile", default=None)
ap.add_argument("--repo", default=str(REPO_ROOT))
ap.add_argument("--duration", type=int, default=300)
ap.add_argument("--turn-timeout", type=int, default=150)
ap.add_argument("--keep-worker", action="store_true")
args = ap.parse_args()
mins = args.duration / 60
print(f"[{now()}] sustained back-and-forth: 1 primary <-> 1 worker for {args.duration}s "
f"(~{mins:.1f} min) profile={args.profile or 'default'} repo={args.repo}", flush=True)
tid, pane = spawn_worker(args.base, args.profile, args.repo, "worker")
await_ready(args.base, tid, "worker")
start = time.time()
expected = 0 # the primary's authoritative running total
reanchor = False # re-state the total next turn after a drift
turns, oks, drifts, breaks = 0, 0, 0, 0
latencies = []
print(f"[{now()}] --- conversation start (worker must keep the total in its head) ---\n", flush=True)
while time.time() - start < args.duration:
turns += 1
step = STEPS[(turns - 1) % len(STEPS)]
if turns == 1:
prompt = RULES.format(step=step)
elif reanchor:
prompt = REANCHOR.format(total=expected, step=step)
reanchor = False
else:
prompt = NEXT.format(step=step)
expected += step
el = round(time.time() - start)
print(f"[{now()}] turn {turns:>2} (t+{el}s) PRIMARY → add {step} (expect total {expected})", flush=True)
rec = one_turn(args.base, tid, prompt, args.turn_timeout)
got = parse_int(rec.get("reply"))
lat = rec.get("latency")
latencies.append(lat)
if rec.get("phase") != "done" or not (rec.get("reply") or "").strip():
breaks += 1
print(f"[{now()}] WORKER ✗ CHANNEL BREAK — phase={rec.get('phase')} "
f"source={rec.get('source')} detail={str(rec.get('detail'))[:100]} ({lat}s)\n", flush=True)
reanchor = True
continue
src = rec.get("source")
badge = "OK " if got == expected else "DRIFT"
if got == expected:
oks += 1
else:
drifts += 1
reanchor = True # re-sync the worker next turn
print(f"[{now()}] WORKER → {str(rec.get('reply')).strip()[:60]!r} = {got} "
f"[{badge}] via {src} ({lat}s)", flush=True)
if got != expected:
print(f"[{now()}] (expected {expected}; will re-anchor next turn)", flush=True)
print(flush=True)
dur = round(time.time() - start)
clean = sum(1 for lat in latencies if lat)
avg = round(sum(latencies) / len(latencies), 1) if latencies else 0
print("=" * 72)
print(f"SUSTAINED CONVERSATION SUMMARY — 1 primary <-> 1 worker over {dur}s (~{dur/60:.1f} min)")
print(f" turns: {turns}")
print(f" clean fleet_reply exchanges: {oks + drifts}/{turns} (channel breaks: {breaks})")
print(f" arithmetic correct (continuity held): {oks}/{turns} (drifts: {drifts})")
print(f" latency: avg {avg}s over {turns} turns")
ok = breaks == 0 and turns >= 2
if ok and drifts == 0:
print(" RESULT: PASS — every turn resolved via fleet_reply and the worker held the "
"running total across the whole window.")
elif ok:
print(f" RESULT: PASS (channel) — every turn resolved via fleet_reply for the full "
f"window; {drifts} arithmetic drift(s) (worker recovered after re-anchor).")
else:
print(" RESULT: FAIL — the channel broke on at least one turn (see CHANNEL BREAK above).")
print("=" * 72, flush=True)
if not args.keep_worker and pane:
try:
http(args.base, "DELETE", f"/workers/{pane}")
print(f"[{now()}] worker stopped (pane {pane})", flush=True)
except Exception as e: # noqa: BLE001
print(f"[{now()}] worker stop failed (ignore): {e}", flush=True)
sys.exit(0 if ok else 1)
if __name__ == "__main__":
main()