Files
fleetd/e2e/conversation_sustained_test.py
T

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 bridge_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 bridge_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 bridge_reply with only the new running total.")
REANCHOR = ("Let's re-sync — the running total is {total}. Now add {step}. Reply via "
"bridge_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 bridge_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 bridge_reply and the worker held the "
"running total across the whole window.")
elif ok:
print(f" RESULT: PASS (channel) — every turn resolved via bridge_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()