From 55ebd5b9492549a81afce1fb69f41c4a1655839d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 15:19:21 +0200 Subject: [PATCH] e2e: sustained back-and-forth conversation harness (5-min stateful continuity) --- e2e/conversation_sustained_test.py | 170 +++++++++++++++++++++++++++++ 1 file changed, 170 insertions(+) create mode 100644 e2e/conversation_sustained_test.py diff --git a/e2e/conversation_sustained_test.py b/e2e/conversation_sustained_test.py new file mode 100644 index 0000000..a7f69de --- /dev/null +++ b/e2e/conversation_sustained_test.py @@ -0,0 +1,170 @@ +#!/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()