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.
171 lines
7.6 KiB
Python
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()
|