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.
331 lines
15 KiB
Python
331 lines
15 KiB
Python
#!/usr/bin/env python3
|
|
"""Standard bridge fan-out test — ONE primary vs MANY workers, concurrently, for
|
|
issue hunting through the running `fleetd` daemon, fully captured, with gap analysis.
|
|
|
|
Where conversation_test.py exercises a single worker over multiple turns, this drives
|
|
the path that only appears under fan-out: the primary spawns N workers, sends each a
|
|
distinct issue-hunting assignment on a slice of the codebase, fires them all at once,
|
|
and collects every reply concurrently. That stresses what a single worker never can —
|
|
|
|
• simultaneous delivery to many panes (the injector's per-worker, not global, writer),
|
|
• per-session rendezvous isolation (N blocked sends resolving independently),
|
|
• reply routing under concurrency (worker A's answer must never resolve worker B's send),
|
|
|
|
and, as the payload, whether a fleet of off-subscription workers can actually surface
|
|
real issues in the repo and report them back structurally via fleet_reply.
|
|
|
|
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 issue_hunt_test.py [--base URL] [--profile NAME] [--repo DIR]
|
|
[--out DIR] [--poll-timeout SECS] [--keep-workers]
|
|
[--assignments FILE]
|
|
|
|
--base bridge REST base URL (default http://127.0.0.1:8765)
|
|
--profile profile for every worker (default: the daemon's default)
|
|
--repo cwd handed to each worker (default: the bridge repo root)
|
|
--out output dir for the transcript (default: alongside this file)
|
|
--poll-timeout per-worker max wait, seconds (default 300)
|
|
--keep-workers do not stop spawned workers at the end
|
|
--assignments JSON file overriding the built-in assignment list
|
|
|
|
Exit code: 0 if every worker delivered AND produced a usable reply with no cross-talk;
|
|
1 otherwise (CI-usable). A per-worker and fleet-level gap report is printed to stdout.
|
|
"""
|
|
import argparse
|
|
import json
|
|
import pathlib
|
|
import sys
|
|
import threading
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from datetime import datetime
|
|
|
|
HERE = pathlib.Path(__file__).parent
|
|
REPO_ROOT = HERE.parent
|
|
|
|
# Each worker gets a distinct source file to hunt in, plus a `probe` — a token its reply
|
|
# should mention if it actually addressed ITS assignment (a soft cross-talk detector: a
|
|
# reply that references only another worker's file is a routing red flag). Targets are the
|
|
# hot files this project has been iterating on, so a real issue is plausible to find.
|
|
DEFAULT_ASSIGNMENTS = [
|
|
{"id": "completion", "probe": "CompletionResolver",
|
|
"target": "fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java"},
|
|
{"id": "worker", "probe": "WorkerService",
|
|
"target": "fleetd/src/main/java/dev/ltms/fleet/worker/WorkerService.java"},
|
|
{"id": "rendezvous", "probe": "Rendezvous",
|
|
"target": "fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java"},
|
|
]
|
|
|
|
PROMPT_TMPL = (
|
|
"You are one of several issue-hunting workers in the claude-bridge repo (it is your "
|
|
"current working directory). Your assignment: inspect the file `{target}` and find the "
|
|
"SINGLE most important real bug, correctness gap, or risk in it. Read the file before "
|
|
"answering. Reply via fleet_reply with EXACTLY these four lines:\n"
|
|
"1. {target}:<line>\n"
|
|
"2. issue: <one sentence>\n"
|
|
"3. fix: <one line>\n"
|
|
"4. severity: high|medium|low\n"
|
|
"Keep it under 90 words. If after reading you find nothing real, reply 'NO ISSUE' and one "
|
|
"line why. Do NOT hunt in any other file — only `{target}`."
|
|
)
|
|
|
|
|
|
def http(base, method, path, body=None, timeout=20):
|
|
"""JSON request. Tolerates an empty body (e.g. 204 No Content on DELETE) → returns {}."""
|
|
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:
|
|
raw = r.read().decode().strip()
|
|
return json.loads(raw) if raw else {}
|
|
|
|
|
|
def now():
|
|
return datetime.now().strftime("%H:%M:%S")
|
|
|
|
|
|
def spawn_worker(base, profile, repo, wid):
|
|
res = http(base, "POST", "/workers", {"profile": profile, "cwd": repo} if profile
|
|
else {"cwd": repo})
|
|
tid = res.get("terminalId") or res.get("sessionId")
|
|
if not tid:
|
|
raise RuntimeError(f"spawn failed for {wid}: {res}")
|
|
pane = res.get("paneId")
|
|
print(f"[{now()}] [{wid}] spawned {tid} (profile={profile or 'default'}, pane={pane})")
|
|
return tid, pane
|
|
|
|
|
|
def await_ready(base, tid, wid, timeout=150):
|
|
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()}] [{wid}] ready (status={st.get('status')})")
|
|
return True
|
|
time.sleep(3)
|
|
print(f"[{now()}] [{wid}] WARNING: never reported ready within {timeout}s — sending anyway")
|
|
return False
|
|
|
|
|
|
def fire(base, tid, prompt):
|
|
"""Fire one async send; return its ticket (delivery is confirmed by a ticket coming back)."""
|
|
sent = http(base, "POST", f"/sessions/{tid}/message", {"wait": False, "content": prompt})
|
|
return sent.get("ticket")
|
|
|
|
|
|
def poll(base, ticket, wid, t0, poll_timeout):
|
|
"""Poll a ticket to resolution; return (record fields) mirroring conversation_test."""
|
|
samples, reply, detail, source, phase = [], None, None, None, 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, detail, source = task.get("reply"), task.get("detail"), task.get("replySource")
|
|
break
|
|
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 {"phase": phase, "reply": reply, "detail": detail, "source": source,
|
|
"latency": round(time.time() - t0, 1), "transitions": " → ".join(trans)}
|
|
|
|
|
|
def grade(rec):
|
|
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 worker_lifecycle(base, profile, repo, poll_timeout, a, barrier):
|
|
"""Full per-worker path: spawn → ready → (barrier) → fire → poll. Spawn+ready run
|
|
concurrently across workers; the send waits on the shared `barrier` so every ready
|
|
worker fires within the same instant — the real simultaneous-delivery stress. A worker
|
|
that fails to spawn aborts the barrier so the rest don't block forever."""
|
|
wid = a["id"]
|
|
rec = {"id": wid, "target": a["target"], "probe": a["probe"], "spawned": False,
|
|
"tid": None, "pane": None, "ticket": None}
|
|
try:
|
|
tid, pane = spawn_worker(base, profile, repo, wid)
|
|
rec.update(tid=tid, pane=pane, spawned=True)
|
|
await_ready(base, tid, wid)
|
|
except Exception as e: # noqa: BLE001
|
|
barrier.abort() # release peers waiting on the barrier
|
|
rec.update(phase="failed", detail=f"spawn/ready error: {e}", reply=None,
|
|
source=None, latency=0.0, transitions="")
|
|
return rec
|
|
|
|
try:
|
|
barrier.wait(timeout=210) # all ready workers proceed together
|
|
except (threading.BrokenBarrierError, Exception): # noqa: BLE001
|
|
pass # a peer died or timed out — fire anyway rather than hang
|
|
t0 = time.time()
|
|
prompt = PROMPT_TMPL.format(target=a["target"])
|
|
try:
|
|
ticket = fire(base, tid, prompt)
|
|
rec["ticket"] = ticket
|
|
print(f"[{now()}] [{wid}] fired (ticket={ticket})")
|
|
rec.update(poll(base, ticket, wid, t0, poll_timeout))
|
|
except Exception as e: # noqa: BLE001
|
|
rec.update(phase="failed", detail=f"send error: {e}", reply=None,
|
|
source=None, latency=round(time.time() - t0, 1), transitions="")
|
|
return rec
|
|
|
|
|
|
def crosstalk_report(records):
|
|
"""Fleet-level isolation checks: distinct tickets, distinct non-empty replies, and each
|
|
reply addressing its OWN assigned file (probe token present). Returns (list_of_gaps)."""
|
|
gaps = []
|
|
tickets = [r.get("ticket") for r in records if r.get("ticket")]
|
|
if len(tickets) != len(set(tickets)):
|
|
gaps.append("ticket collision: two workers were handed the same ticket id")
|
|
replies = {r["id"]: (r.get("reply") or "").strip() for r in records}
|
|
# identical non-empty replies from distinct assignments ⇒ suspected reply misrouting
|
|
seen = {}
|
|
for wid, text in replies.items():
|
|
if text and text in seen:
|
|
gaps.append(f"identical reply from '{seen[text]}' and '{wid}' "
|
|
f"(distinct assignments should not yield byte-identical answers)")
|
|
elif text:
|
|
seen[text] = wid
|
|
# a reply that names ANOTHER worker's file but not its own ⇒ likely cross-routing
|
|
for r in records:
|
|
text = (r.get("reply") or "")
|
|
if not text.strip():
|
|
continue
|
|
own = r["probe"] in text or pathlib.Path(r["target"]).name in text
|
|
others = [o["probe"] for o in records if o["id"] != r["id"] and o["probe"] in text]
|
|
if not own and others:
|
|
gaps.append(f"worker '{r['id']}' (assigned {r['probe']}) replied about "
|
|
f"{others} but not its own file — possible cross-routing")
|
|
return gaps
|
|
|
|
|
|
def write_transcript(out_dir, records, meta):
|
|
path = out_dir / "issue_hunt_transcript.md"
|
|
with path.open("w") as f:
|
|
f.write(f"# Bridge fan-out issue-hunt — {datetime.now():%Y-%m-%d %H:%M}\n\n")
|
|
f.write(f"One primary vs **{len(records)} concurrent workers** "
|
|
f"(profile `{meta['profile']}`), each hunting a distinct file.\n\n")
|
|
for r in records:
|
|
g, note = grade(r)
|
|
f.write(f"### `{r['id']}` — {r['target']} (`{g}`, {r.get('latency')}s, "
|
|
f"phase={r.get('phase')}, source={r.get('source')})\n\n")
|
|
f.write(f"**ASSIGNMENT:** find the top issue in `{r['target']}`\n\n")
|
|
reply = r.get("reply")
|
|
f.write(f"**WORKER {r['id']}:** {reply if reply else '_(no reply)_ ' + str(r.get('detail'))}\n\n")
|
|
if r.get("transitions"):
|
|
f.write(f"_status: {r['transitions']}_\n")
|
|
if g != "OK":
|
|
f.write(f"\n> **GAP — {g}:** {note}\n")
|
|
f.write("\n")
|
|
if meta["gaps"]:
|
|
f.write("## Fleet-level gaps\n\n")
|
|
for gp in meta["gaps"]:
|
|
f.write(f"- {gp}\n")
|
|
return path
|
|
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description="Bridge fan-out issue-hunt test (1 primary, N workers)")
|
|
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("--out", default=str(HERE))
|
|
ap.add_argument("--poll-timeout", type=int, default=300)
|
|
ap.add_argument("--keep-workers", action="store_true")
|
|
ap.add_argument("--assignments", default=None)
|
|
args = ap.parse_args()
|
|
|
|
assignments = DEFAULT_ASSIGNMENTS
|
|
if args.assignments:
|
|
assignments = json.loads(pathlib.Path(args.assignments).read_text())
|
|
|
|
n = len(assignments)
|
|
barrier = threading.Barrier(n) # releases exactly when all n ready workers reach it
|
|
print(f"[{now()}] fan-out issue-hunt: 1 primary vs {n} workers "
|
|
f"(profile={args.profile or 'default'}, repo={args.repo})\n")
|
|
|
|
# Spawn + ready + fire + poll all workers concurrently; the barrier makes every send fire
|
|
# together once all are ready, so delivery pressure hits the daemon simultaneously.
|
|
records = []
|
|
with ThreadPoolExecutor(max_workers=n) as ex:
|
|
futures = [ex.submit(worker_lifecycle, args.base, args.profile, args.repo,
|
|
args.poll_timeout, a, barrier) for a in assignments]
|
|
for fut in futures:
|
|
records.append(fut.result())
|
|
|
|
records.sort(key=lambda r: [a["id"] for a in assignments].index(r["id"]))
|
|
|
|
print()
|
|
for r in records:
|
|
g, note = grade(r)
|
|
print(f"[{r['id']:11}] {g:8} {str(r.get('latency','?')):6}s "
|
|
f"phase={r.get('phase')} source={r.get('source')}")
|
|
print(f" status: {r.get('transitions') or '(none)'}")
|
|
print(f" reply: {((r.get('reply') or '(none) ' + str(r.get('detail'))).strip()[:200])}")
|
|
print(f" note: {note}\n")
|
|
|
|
gaps = crosstalk_report(records)
|
|
meta = {"profile": args.profile or "default", "gaps": gaps}
|
|
out_dir = pathlib.Path(args.out)
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
path = write_transcript(out_dir, records, meta)
|
|
|
|
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("=" * 72)
|
|
print(f"FAN-OUT ISSUE-HUNT SUMMARY — 1 primary vs {n} workers")
|
|
print(" channel: " + " ".join(f"{g}:{v}" for g, v in counts.items()))
|
|
print(f" transcript: {path}")
|
|
if gaps:
|
|
print(" FLEET GAPS:")
|
|
for gp in gaps:
|
|
print(f" ⚠ {gp}")
|
|
else:
|
|
print(" isolation: clean — distinct tickets, distinct replies, each on its own file")
|
|
delivered = all(r.get("ticket") for r in records)
|
|
replied = all(g in ("OK", "DEGRADED") for g in grades)
|
|
if delivered and replied and not gaps:
|
|
print(" RESULT: PASS — all workers delivered concurrently, replied, and stayed isolated.")
|
|
elif delivered and replied:
|
|
print(" RESULT: PASS (with notes) — all delivered & replied, but see FLEET GAPS.")
|
|
else:
|
|
print(" RESULT: FAIL — a worker did not deliver or did not reply (see GAP notes).")
|
|
print("=" * 72)
|
|
|
|
if not args.keep_workers:
|
|
for r in records:
|
|
if r.get("spawned") and r.get("pane"):
|
|
try:
|
|
http(args.base, "DELETE", f"/workers/{r['pane']}")
|
|
print(f"[{now()}] [{r['id']}] stopped (pane {r['pane']})")
|
|
except Exception as e: # noqa: BLE001
|
|
print(f"[{now()}] [{r['id']}] stop failed (ignore): {e}")
|
|
|
|
sys.exit(0 if (delivered and replied and not gaps) else 1)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|