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 `bridged` 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()
|