3ce76a5d69
The alias removal left bridge_* tool names in prose. Fix them: - README no longer claims the old bridge_* names still answer (they were removed). - pom + LeadTabScanner comments name fleet_* tools. - FleetMcp comment no longer mentions the removed deprecated twin. - docs/MCP-Contract.md and e2e swept bridge_* -> fleet_*; e2e ask files renamed. The historical mcp__bridge__* mount-name note in CLAUDE.md is kept on purpose. 949 tests pass.
278 lines
14 KiB
Python
278 lines
14 KiB
Python
#!/usr/bin/env python3
|
|
"""Live fleet_ask test — the REVERSE rendezvous (CB-205), watched end to end.
|
|
|
|
Every other harness drives the forward path: primary `fleet_send` → worker `fleet_reply`.
|
|
This drives the one that runs the other way. A worker is told to pause its delegated turn,
|
|
ask the primary a question via `fleet_ask`, and only finish once it has the answer — so the
|
|
turn round-trips primary→worker→primary→worker inside a SINGLE delegation.
|
|
|
|
The mechanics that only this path exercises:
|
|
|
|
• a worker's mid-turn question surfacing on the primary's *own* blocked send (Outcome.QUESTION),
|
|
• the `turnId` correlation that lets the primary answer the exact paused turn,
|
|
• the answer resuming that same turn and the worker's final `fleet_reply` landing on the
|
|
re-opened forward waiter (never a stale or cross-wired one).
|
|
|
|
It is two blocking REST calls, no polling:
|
|
|
|
1. POST /sessions/{id}/message {content: TASK} → blocks, returns 202 {status:"question",
|
|
question, turnId} (the worker asked)
|
|
2. POST /sessions/{id}/message {content: ANSWER, turnId} → blocks, returns 200 {reply, replySource}
|
|
(the worker resumed and replied)
|
|
|
|
Like the rest of the suite it talks ONLY to the bridge's REST face on loopback — it never sets
|
|
ANTHROPIC_BASE_URL and never touches herdr, so it is subscription-safe by construction.
|
|
|
|
Usage:
|
|
python3 fleet_ask_test.py [--base URL] [--profile NAME] [--repo DIR] [--out DIR]
|
|
[--send-timeout SECS] [--answer-timeout SECS] [--keep-worker]
|
|
|
|
--base bridge REST base URL (default http://127.0.0.1:8765)
|
|
--profile profile for the worker (default: the daemon's default)
|
|
--repo cwd handed to the worker (default: the bridge repo root)
|
|
--out output dir for the transcript (default: alongside this file)
|
|
--send-timeout max wait for the worker to ASK (default 110)
|
|
--answer-timeout max wait for the worker to REPLY (default 90)
|
|
--keep-worker do not stop the spawned worker at the end
|
|
|
|
Exit code: 0 if the worker asked, the answer resumed the turn, and the final reply reflected
|
|
the answer; 1 otherwise (CI-usable).
|
|
"""
|
|
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
|
|
REPO_ROOT = HERE.parent
|
|
|
|
# The color the primary will hand back when the worker asks. The final reply must reflect it,
|
|
# uppercased — proof the answer actually reached the resumed turn (not a value the worker could
|
|
# have guessed: it is told to ask, and only the primary knows which of red/blue is chosen).
|
|
ANSWER_COLOR = "blue"
|
|
|
|
# A task that CANNOT be completed without asking: the worker is not told which color to choose,
|
|
# only that the primary will name one when asked. So a correct final reply is only reachable by
|
|
# actually calling fleet_ask and using the answer.
|
|
TASK_PROMPT = (
|
|
"You are a bridge worker in a quick coordination game. You do NOT know which color to pick — "
|
|
"only the primary does. Do exactly this, in order:\n"
|
|
"1. Call the `fleet_ask` tool with EXACTLY this question: \"PICK A COLOR: red or blue?\"\n"
|
|
"2. The primary will answer with one color word. Take that color and uppercase it.\n"
|
|
"3. Call `fleet_reply` with EXACTLY one line: CHOSEN=<COLOR> (e.g. CHOSEN=GREEN if told green).\n"
|
|
"Do not guess a color. Do not call fleet_reply before fleet_ask has returned an answer. "
|
|
"Do nothing else — no file reads, no other tools."
|
|
)
|
|
|
|
|
|
def http(base, method, path, body=None, timeout=20):
|
|
"""JSON request. Tolerates an empty body (e.g. 204) → {}. Raises on non-2xx via urllib."""
|
|
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):
|
|
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: {res}")
|
|
pane = res.get("paneId")
|
|
print(f"[{now()}] spawned {tid} (profile={profile or 'default'}, pane={pane})")
|
|
return tid, pane
|
|
|
|
|
|
def await_ready(base, tid, 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()}] worker ready (status={st.get('status')})")
|
|
return True
|
|
time.sleep(3)
|
|
print(f"[{now()}] WARNING: worker never reported ready within {timeout}s — sending anyway")
|
|
return False
|
|
|
|
|
|
def post_message(base, tid, body, timeout):
|
|
"""One BLOCKING send. Returns (http_status, parsed_json). urllib raises on 4xx/5xx, so a
|
|
stale-turn 409 is surfaced here rather than swallowed."""
|
|
data = json.dumps(body).encode()
|
|
req = urllib.request.Request(f"{base}/sessions/{tid}/message", data=data, method="POST",
|
|
headers={"Content-Type": "application/json"})
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=timeout) as r:
|
|
raw = r.read().decode().strip()
|
|
return r.status, (json.loads(raw) if raw else {})
|
|
except urllib.error.HTTPError as e:
|
|
raw = e.read().decode().strip()
|
|
return e.code, (json.loads(raw) if raw else {})
|
|
|
|
|
|
def run(base, profile, repo, send_timeout, answer_timeout):
|
|
"""Drive the full reverse rendezvous. Returns a result record for grading + transcript."""
|
|
rec = {"spawned": False, "tid": None, "pane": None, "phase": "spawn",
|
|
"question": None, "turnId": None, "reply": None, "replySource": None,
|
|
"detail": None, "ask_latency": None, "answer_latency": None}
|
|
|
|
tid, pane = spawn_worker(base, profile, repo)
|
|
rec.update(tid=tid, pane=pane, spawned=True)
|
|
await_ready(base, tid)
|
|
|
|
# 1) Delegate the ask-forcing task. This blocks until the worker calls fleet_ask, at which
|
|
# point our own send unblocks carrying the question and the turnId to answer on.
|
|
print(f"[{now()}] delegating task (blocks until the worker asks; up to {send_timeout}s)…")
|
|
t0 = time.time()
|
|
rec["phase"] = "awaiting_question"
|
|
try:
|
|
code, resp = post_message(base, tid, {"content": TASK_PROMPT, "timeoutMs": send_timeout * 1000},
|
|
timeout=send_timeout + 15)
|
|
except Exception as e: # noqa: BLE001
|
|
rec.update(phase="send_error", detail=f"delegating send failed: {e}")
|
|
return rec
|
|
rec["ask_latency"] = round(time.time() - t0, 1)
|
|
|
|
if resp.get("status") != "question":
|
|
# The worker finished (or stalled) without asking — the whole point didn't happen.
|
|
rec.update(phase="no_question", detail=f"HTTP {code}: {json.dumps(resp)[:300]}",
|
|
reply=resp.get("reply"), replySource=resp.get("replySource"))
|
|
return rec
|
|
rec.update(phase="question", question=resp.get("question"), turnId=resp.get("turnId"))
|
|
print(f"[{now()}] worker ASKED ({rec['ask_latency']}s): {rec['question']!r} turnId={rec['turnId']}")
|
|
|
|
if not rec["turnId"]:
|
|
rec.update(phase="no_turnid", detail="question surfaced without a turnId to answer on")
|
|
return rec
|
|
|
|
# 2) Answer on that exact turn. This blocks again until the resumed worker calls fleet_reply.
|
|
print(f"[{now()}] answering '{ANSWER_COLOR}' on turn {rec['turnId']} (blocks until reply; up to {answer_timeout}s)…")
|
|
t1 = time.time()
|
|
rec["phase"] = "awaiting_reply"
|
|
try:
|
|
code, resp = post_message(base, tid,
|
|
{"content": ANSWER_COLOR, "turnId": rec["turnId"],
|
|
"timeoutMs": answer_timeout * 1000},
|
|
timeout=answer_timeout + 15)
|
|
except Exception as e: # noqa: BLE001
|
|
rec.update(phase="answer_error", detail=f"answer send failed: {e}")
|
|
return rec
|
|
rec["answer_latency"] = round(time.time() - t1, 1)
|
|
|
|
if code == 409 or resp.get("error") == "stale_turn":
|
|
rec.update(phase="stale_turn", detail=f"HTTP {code}: {json.dumps(resp)[:300]}")
|
|
return rec
|
|
if resp.get("reply") is None:
|
|
rec.update(phase="no_reply", detail=f"HTTP {code}: {json.dumps(resp)[:300]}")
|
|
return rec
|
|
rec.update(phase="replied", reply=resp.get("reply"), replySource=resp.get("replySource"))
|
|
print(f"[{now()}] worker RESUMED and replied ({rec['answer_latency']}s): {rec['reply']!r} "
|
|
f"(source={rec['replySource']})")
|
|
return rec
|
|
|
|
|
|
def grade(rec):
|
|
"""PASS only if the worker asked, the turn resumed, and the reply reflects the answer."""
|
|
if rec["phase"] == "no_question":
|
|
return "NO_ASK", "the worker finished/stalled without ever calling fleet_ask"
|
|
if rec["phase"] in ("send_error", "answer_error", "spawn"):
|
|
return "ERROR", rec.get("detail") or "transport error before the round-trip completed"
|
|
if rec["phase"] == "no_turnid":
|
|
return "NO_TURNID", "the question surfaced without a turnId — the primary could not answer"
|
|
if rec["phase"] == "stale_turn":
|
|
return "STALE", "answering the turn was rejected as stale (it lapsed or was already answered)"
|
|
if rec["phase"] in ("awaiting_reply", "no_reply"):
|
|
return "NO_RESUME", "the worker asked but never resumed to a final reply within the window"
|
|
if rec["phase"] == "replied":
|
|
reflected = ANSWER_COLOR.upper() in (rec["reply"] or "").upper()
|
|
if reflected and rec["replySource"] == "reply":
|
|
return "OK", "asked, resumed the same turn, and the reply reflected the primary's answer"
|
|
if reflected:
|
|
return "DEGRADED", f"reply reflected the answer but resolved via {rec['replySource']} " \
|
|
"(worker did not call fleet_reply cleanly)"
|
|
return "WRONG_ANSWER", f"the worker replied but did not reflect '{ANSWER_COLOR}' — " \
|
|
f"the answer may not have reached the resumed turn: {rec['reply']!r}"
|
|
return "WEDGE", f"unexpected terminal phase {rec['phase']}: {rec.get('detail')}"
|
|
|
|
|
|
def write_transcript(out_dir, rec, meta):
|
|
path = out_dir / "fleet_ask_transcript.md"
|
|
g, note = grade(rec)
|
|
with path.open("w") as f:
|
|
f.write(f"# Live fleet_ask — reverse rendezvous — {datetime.now():%Y-%m-%d %H:%M}\n\n")
|
|
f.write(f"One worker paused its delegated turn to ask the primary, then resumed with the "
|
|
f"answer (profile `{meta['profile']}`). Result: **`{g}`**.\n\n")
|
|
f.write("## Round-trip\n\n")
|
|
f.write(f"1. **primary → worker** (delegation): the ask-forcing task.\n")
|
|
f.write(f"2. **worker → primary** (`fleet_ask`, {rec.get('ask_latency')}s): "
|
|
f"{rec.get('question')!r} — surfaced on the primary's blocked send as a "
|
|
f"`question` with `turnId={rec.get('turnId')}`.\n")
|
|
f.write(f"3. **primary → worker** (answer on that turn): `{ANSWER_COLOR}`.\n")
|
|
f.write(f"4. **worker → primary** (`fleet_reply`, {rec.get('answer_latency')}s, "
|
|
f"source={rec.get('replySource')}): {rec.get('reply')!r}\n\n")
|
|
f.write(f"> **{g}:** {note}\n")
|
|
if rec.get("detail"):
|
|
f.write(f">\n> detail: {rec['detail']}\n")
|
|
return path
|
|
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description="Live fleet_ask reverse-rendezvous test (CB-205)")
|
|
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("--send-timeout", type=int, default=110, help="max wait for the worker to ASK")
|
|
ap.add_argument("--answer-timeout", type=int, default=90, help="max wait for the worker to REPLY")
|
|
ap.add_argument("--keep-worker", action="store_true")
|
|
args = ap.parse_args()
|
|
|
|
print(f"[{now()}] live fleet_ask: 1 primary, 1 worker "
|
|
f"(profile={args.profile or 'default'}, repo={args.repo})\n")
|
|
|
|
rec = {"spawned": False, "pane": None}
|
|
try:
|
|
rec = run(args.base, args.profile, args.repo, args.send_timeout, args.answer_timeout)
|
|
finally:
|
|
if not args.keep_worker and rec.get("spawned") and rec.get("pane"):
|
|
try:
|
|
http(args.base, "DELETE", f"/workers/{rec['pane']}")
|
|
print(f"[{now()}] stopped worker (pane {rec['pane']})")
|
|
except Exception as e: # noqa: BLE001
|
|
print(f"[{now()}] stop failed (ignore): {e}")
|
|
|
|
g, note = grade(rec)
|
|
out_dir = pathlib.Path(args.out)
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
path = write_transcript(out_dir, rec, {"profile": args.profile or "default"})
|
|
|
|
print()
|
|
print("=" * 72)
|
|
print("LIVE fleet_ask SUMMARY — reverse rendezvous (CB-205)")
|
|
print(f" asked: {rec.get('question')!r} (turnId={rec.get('turnId')}, {rec.get('ask_latency')}s)")
|
|
print(f" answered: {ANSWER_COLOR!r}")
|
|
print(f" replied: {rec.get('reply')!r} (source={rec.get('replySource')}, {rec.get('answer_latency')}s)")
|
|
print(f" transcript: {path}")
|
|
print(f" RESULT: {g} — {note}")
|
|
print("=" * 72)
|
|
|
|
sys.exit(0 if g == "OK" else 1)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|