From 6cc967a0fd53bbbc313f8b11bb19de2dcd4ca52a Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Wed, 23 Sep 2026 17:29:50 +0000 Subject: [PATCH] test(e2e): merge-order-safe live-gap xfails and deterministic delivery cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open. --- tests/conformance/persistence/README.md | 2 +- ...test_cell5_delivery_outbox_exactly_once.py | 19 +-- tests/e2e/core/delivery/_pending_fixes.py | 150 ++++++++++++++++++ .../delivery/test_messaging_exactly_once.py | 107 ++++++++----- 4 files changed, 229 insertions(+), 49 deletions(-) create mode 100644 tests/e2e/core/delivery/_pending_fixes.py diff --git a/tests/conformance/persistence/README.md b/tests/conformance/persistence/README.md index e9be19943f..1b959835f4 100644 --- a/tests/conformance/persistence/README.md +++ b/tests/conformance/persistence/README.md @@ -14,7 +14,7 @@ process mid-write. | 2 — `test_cell2_consume_once` | a parked handoff is claimed by exactly one of N racing processes | adapted from the tracking issue's spot-probe (8-process file-barrier race) | | 3 — `test_cell3_rotation_atomicity` | a compression rotation is visible entirely or not at all — never a compression-ended parent without a continuation (the #80337 orphan shape; recovery for the legacy population merged in #80487) | new in this suite | | 4 — fork determinism on edit/rewind | recovery yields exactly the chosen prefix after a fork | **stub** — interlocked with the rewind/archive redesign (#82956–#82959) | -| 5 — `test_cell5_delivery_outbox_exactly_once` | a crash between provider send and durable record must not double-deliver on catch-up (gateway reboot): per obligation ≤1 unmarked copy, ≥1 copy, every extra copy carries a recovery marker, ledger terminal after the boot sweep, later reboots send nothing; concurrent rebooters never both claim/send a row — effect exactly-once, distinct from cell 2's consume-once | **implemented** — real ledger + real `GatewayRunner` boot claim/redeliver halves, SIGKILL at each kill point; only the transport is faked (fsync'd journal). One strict-xfail **fire**: a boot killed inside a *plain* ('pending') redelivery resends it unmarked next boot. Cron-ticker catch-up not covered yet (#83197/#83557) | +| 5 — `test_cell5_delivery_outbox_exactly_once` | a crash between provider send and durable record must not double-deliver on catch-up (gateway reboot): per obligation ≤1 unmarked copy, ≥1 copy, every extra copy carries a recovery marker, ledger terminal after the boot sweep, later reboots send nothing; concurrent rebooters never both claim/send a row — effect exactly-once, distinct from cell 2's consume-once | **implemented** — real ledger + real `GatewayRunner` boot claim/redeliver halves, SIGKILL at each kill point; only the transport is faked (fsync'd journal). One **fire**, strict-xfail while it reproduces (fix: #120450): a boot killed inside a *plain* ('pending') redelivery resends it unmarked next boot. Cron-ticker catch-up not covered yet (#83197/#83557) | ## Method diff --git a/tests/conformance/persistence/test_cell5_delivery_outbox_exactly_once.py b/tests/conformance/persistence/test_cell5_delivery_outbox_exactly_once.py index 5e384d1cbd..57b0c8fb99 100644 --- a/tests/conformance/persistence/test_cell5_delivery_outbox_exactly_once.py +++ b/tests/conformance/persistence/test_cell5_delivery_outbox_exactly_once.py @@ -75,6 +75,7 @@ from tests.conformance.persistence._harness import ( reap, wait_for, ) +from tests.e2e.core.delivery._pending_fixes import expect_gap SESSION_KEY = "agent:main:telegram:dm:cell5" CHAT_ID = "4242" @@ -446,16 +447,16 @@ def test_crash_between_send_and_record_never_double_delivers(tmp_path, kill_poin assert integrity_ok(cell.db_path) -@pytest.mark.xfail( - strict=True, - reason=( - "FIRE (cell 5): sweep_recoverable claims a 'pending' row without flipping it to " - "'attempting', and _redeliver_claimed_obligations never marks it attempting before the " - "send — a boot killed inside that resend leaves it 'pending', so the next boot resends it " - "UNMARKED (2 unmarked copies). strict: flips red when fixed so the xfail is removed." - ), +PLAIN_REDELIVERY_GAP = ( + "FIRE (cell 5, fixed by #120450): sweep_recoverable claims a 'pending' row without flipping it " + "to 'attempting', and _redeliver_claimed_obligations never marks it attempting before the send " + "— a boot killed inside that resend leaves it 'pending', so the next boot resends it UNMARKED " + "(2 unmarked copies)." ) -def test_boot_killed_inside_plain_redelivery_never_double_delivers(tmp_path): + + +def test_boot_killed_inside_plain_redelivery_never_double_delivers(tmp_path, request): + expect_gap(request, 120450, PLAIN_REDELIVERY_GAP) # strict xfail only while the gap reproduces cell = Cell(tmp_path) ids = cell.seed_and_crash(["pending"] * N_OBLIGATIONS) cell.crashing_boot() diff --git a/tests/e2e/core/delivery/_pending_fixes.py b/tests/e2e/core/delivery/_pending_fixes.py new file mode 100644 index 0000000000..b8ce631bf3 --- /dev/null +++ b/tests/e2e/core/delivery/_pending_fixes.py @@ -0,0 +1,150 @@ +"""Merge-order-safe expected failures for live gaps whose fix is an open PR. + +A plain ``xfail(strict=True)`` turns main red the moment its fix merges (XPASS), and a +non-strict one guards nothing. Instead, each gap here has a PROBE: a few lines that reproduce +the defect's mechanism on the tree under test, in a throwaway interpreter with its own +``HOME``/``HERMES_HOME`` (no state leaks into the suite process). ``expect_gap`` applies a +strict xfail only while the probe still reproduces the defect; once the fix is in the tree the +cell runs as a plain test, so it must pass. Whichever lands first, suite or fix, main stays +green, and a probe that disagrees with the end-to-end cell still fails loudly (XPASS, or a +real failure) instead of hiding. + +When a fix has landed, delete its entry and every ``expect_gap`` call naming it. +""" + +from __future__ import annotations + +import functools +import json +import os +import subprocess +import sys +import tempfile +from pathlib import Path +from typing import Dict + +import pytest + +REPO_ROOT = Path(__file__).resolve().parents[4] + +_PRELUDE = "import json, os, sys\nsys.path.insert(0, os.getcwd())\n" + +# PR -> (extra env, script). A script prints ``open`` while the defect reproduces, ``fixed`` once +# it no longer does; anything else (including a crash) fails the cell that asked. +PROBES: Dict[int, tuple] = { + # The due gate compared same-zone wall clocks: 01:00 EST (fold=1) looked due at 01:01 EDT. + 120314: ({"HERMES_TIMEZONE": "America/New_York", "TZ": "UTC"}, r''' +from datetime import datetime, timezone +from zoneinfo import ZoneInfo +from cron import jobs +ny = ZoneInfo("America/New_York") +now = datetime(2026, 11, 1, 5, 1, tzinfo=timezone.utc).astimezone(ny) # 01:01 EDT +scheduled = datetime(2026, 11, 1, 6, 0, tzinfo=timezone.utc).astimezone(ny) # 01:00 EST, 59 min later +jobs._hermes_now = lambda: now +jobs.save_jobs([{"id": "p", "name": "p", "prompt": "p", "schedule": {"kind": "interval", "minutes": 60}, + "next_run_at": scheduled.isoformat(), "last_run_at": None, "enabled": True, + "state": "scheduled", "repeat": {"times": None, "completed": 0}, "deliver": "local"}]) +print("open" if jobs.get_due_jobs() else "fixed") +'''), + # No timezone configured: the next cron occurrence kept the base time's fixed UTC offset, so a + # 09:00 job in a DST process zone fired at 10:00 local the day after spring-forward. + 119970: ({"TZ": "America/New_York"}, r''' +from datetime import datetime +from cron import jobs +# what an unconfigured clock returns: the process zone's offset of the moment, as a fixed offset +jobs._hermes_now = lambda: datetime.fromisoformat("2026-03-07T09:00:30-05:00") +nxt = jobs.compute_next_run({"kind": "cron", "expr": "0 9 * * *"}) +print("fixed" if nxt == "2026-03-08T09:00:00-04:00" else "open") +'''), + # A failed first stream send disabled edits but left no message id, so the next tick sent a + # second first send: an uneditable partial preview stayed visible next to the final reply. + 120315: ({}, r''' +import asyncio +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock +from gateway.stream_consumer import GatewayStreamConsumer, StreamConsumerConfig + +async def main(): + delivered, results = [], iter([SimpleNamespace(success=False, error="timeout")]) + async def send(**kw): + r = next(results, None) or SimpleNamespace(success=True, message_id=f"m{len(delivered)}") + if r.success: + delivered.append(kw["content"]) + return r + adapter = MagicMock() + adapter.send = AsyncMock(side_effect=send) + adapter.edit_message = AsyncMock(return_value=SimpleNamespace(success=True)) + adapter.MAX_MESSAGE_LENGTH = 4096 + c = GatewayStreamConsumer(adapter, "chat", StreamConsumerConfig(edit_interval=0.01, + buffer_threshold=5, cursor="")) + c.on_delta("preview never landed ") + task = asyncio.create_task(c.run()) + await asyncio.sleep(0.08) + c.on_delta("and more streamed text ") + await asyncio.sleep(0.08) + c.on_delta("then the end.") + c.finish() + await asyncio.wait_for(task, timeout=10) + return delivered + +print("fixed" if asyncio.run(main()) == ["preview never landed and more streamed text then the end."] + else "open") +'''), + # Unclean startup ran the 120 s recency sweep: every recently active session was marked + # resume_pending and auto-resumed, i.e. answered a second time. + 120377: ({}, r''' +import inspect +from gateway.run import GatewayRunner +src = inspect.getsource(GatewayRunner._recover_unclean_sessions) +print("open" if "suspend_recently_active" in src else "fixed") +'''), + # Inbound de-duplication lived only on the adapter instance; the reconnect watcher builds a + # fresh adapter without it, so a replay after a reconnect was processed again. + 120444: ({}, r''' +import inspect +from gateway.run import GatewayRunner +src = inspect.getsource(GatewayRunner._reconnect_failed_platform) +print("fixed" if "dedup" in src.lower() else "open") +'''), + # The boot sweep claimed a 'pending' row without moving it to 'attempting', so a boot killed + # inside that plain resend left 'pending' behind and the next boot resent it UNMARKED. + 120450: ({}, r''' +import sqlite3 +from gateway import delivery_ledger as L +L.record_obligation(obligation_id="p", session_key="k", platform="telegram", chat_id="1", + thread_id=None, content="x") +with L._transaction() as conn: + conn.execute("UPDATE delivery_obligations SET owner_pid=NULL, owner_started_at=NULL") +assert [r["obligation_id"] for r in L.sweep_recoverable()] == ["p"] +with L._transaction() as conn: + state = conn.execute("SELECT state FROM delivery_obligations").fetchone()[0] +print("fixed" if state == "attempting" else "open") +'''), +} + + +@functools.lru_cache(maxsize=None) +def gap_open(pr: int) -> bool: + """True while the defect PR ``pr`` fixes still reproduces on this tree.""" + extra_env, script = PROBES[pr] + with tempfile.TemporaryDirectory(prefix=f"gap-{pr}-") as tmp: + home = Path(tmp) / "home" + (home / ".hermes").mkdir(parents=True) + env = {k: v for k, v in os.environ.items() + if not k.startswith(("PYTEST_", "HERMES_")) and not k.endswith("_API_KEY")} + env.update({"HOME": str(home), "HERMES_HOME": str(home / ".hermes"), + "PYTHONPATH": str(REPO_ROOT), **extra_env}) + proc = subprocess.run([sys.executable, "-c", _PRELUDE + script], cwd=str(REPO_ROOT), env=env, + stdin=subprocess.DEVNULL, capture_output=True, text=True, timeout=120) + verdict = proc.stdout.strip().splitlines()[-1:] if proc.returncode == 0 else [] + if verdict not in (["open"], ["fixed"]): + raise AssertionError(f"probe for #{pr} is broken (rc={proc.returncode}): " + f"{json.dumps(proc.stdout[-800:])} {proc.stderr[-2000:]}") + return verdict == ["open"] + + +def expect_gap(request, pr: int, reason: str) -> None: + """Strict xfail for this cell while #``pr``'s defect reproduces; a plain test once it doesn't.""" + assert f"#{pr}" in reason, f"reason for a #{pr} gap must name the PR: {reason!r}" + if gap_open(pr): + request.applymarker(pytest.mark.xfail(strict=True, reason=reason)) diff --git a/tests/e2e/core/delivery/test_messaging_exactly_once.py b/tests/e2e/core/delivery/test_messaging_exactly_once.py index 4427d4d172..e1e8bb1552 100644 --- a/tests/e2e/core/delivery/test_messaging_exactly_once.py +++ b/tests/e2e/core/delivery/test_messaging_exactly_once.py @@ -45,7 +45,8 @@ from tests.e2e.core.delivery._fake_platform import ( visible_copies, wait_until, ) -from tests.fakes.fake_llm_provider import FakeLLMServer, Text, ToolCall +from tests.e2e.core.delivery._pending_fixes import expect_gap, gap_open +from tests.fakes.fake_llm_provider import FakeLLMServer, StallMidStream, Text, ToolCall pytestmark = pytest.mark.skipif(sys.platform == "win32", reason="SIGKILL + process-group restart harness") @@ -221,9 +222,14 @@ FAULTS = { "tool_then_answer": (None, "any", "footer", lambda aid: [ToolCall("terminal", {"command": "echo hi"}), Director.answer(aid, "after the tool, done.")], "exact"), - "long_split": (None, "any", "footer", _say(LONG, chunk_chars=400), "exact"), - "long_streamed": (None, "any", "footer", _say(LONG, chunk_chars=120, delay_per_chunk=0.01), "exact"), - # Two 4.5k bursts 3 s apart (edit interval 0.8 s): the first send overflows (split, tail kept as + # The whole 8.1k reply in ONE stream chunk: the first send overflows every cap and is split + # once, with nothing streamed after the split. + "long_split": (None, "any", "footer", _say(LONG, chunk_chars=10_000), "exact"), + # 2.5k chunks 1 s apart (edit interval 0.05 s), so every chunk lands in its own consumer tick: + # the Discord-like first send already overflows (split, tail kept as the live preview) and the + # next chunks extend that tail; the Telegram-like one overflows on an edit and seals. + "long_streamed": (None, "any", "footer", _say(LONG, chunk_chars=2500, delay_per_chunk=1.0), "exact"), + # Two 4.5k bursts 3 s apart (edit interval 0.05 s): the first send overflows (split, tail kept as # the live preview), then the second burst leaves a remainder over the limit after sealing. "burst_overflow": (None, "any", "footer", _say(LONG, chunk_chars=4500, delay_per_chunk=3.0), "exact"), "slow_stream": (None, "any", "footer", _say("a slowly streamed answer " * 6, **SLOW), "exact"), @@ -238,25 +244,22 @@ FAULTS = { "one_copy"), "stream_ack_lost_final_edit": ("ack_lost", "any", "footer", _say("stream ack lost " * 8, **SLOW), "marked_dupes"), - "stream_timeout_first_send": ("timeout", "send", "header", _say("preview never landed " * 8, **SLOW), + # The first preview send times out; four more chunks follow, each in its own consumer tick. + "stream_timeout_first_send": ("timeout", "send", "header", + _say("preview never landed " * 8, chunk_chars=42, delay_per_chunk=0.5), "one_copy"), } -# Red until #120315 lands (overflow tail keeps the "(n/n)" indicator / seal remainder re-sent; -# uneditable partial preview after a failed first send): every streamed reply over the platform -# limit, plus the first-send timeout. Value = strict: burst_overflow and the 2000-char first-send -# timeout fail every run; the others depend on stream/edit-tick timing, so a strict mark would flake -# on the runs where the bug does not trigger. Delete the table with #120315. -UNLANDED_STREAM_FIX = { - ("burst_overflow", "fk_tg"): True, - ("burst_overflow", "fk_dc"): True, - ("stream_timeout_first_send", "fk_dc"): True, - ("stream_timeout_first_send", "fk_tg"): False, - ("long_split", "fk_tg"): False, - ("long_split", "fk_dc"): False, - ("long_streamed", "fk_tg"): False, - ("long_streamed", "fk_dc"): False, +# Streaming gaps #120315 fixes: an overflowing first send keeps the " (n/n)" indicator on the live +# tail that later chunks extend / a sealed remainder over the limit is re-sent; a failed first send +# leaves an uneditable partial preview next to the final. Every cell here fails every run while the +# gap is open (the chunk pacing above puts each chunk in its own consumer tick). +STREAM_OVERFLOW_GAP = "LIVE GAP (fixed by #120315): streamed overflow / failed first send" +STREAM_OVERFLOW_GAP_CELLS = { + ("burst_overflow", "fk_tg"), ("burst_overflow", "fk_dc"), + ("long_streamed", "fk_dc"), + ("stream_timeout_first_send", "fk_tg"), ("stream_timeout_first_send", "fk_dc"), } @@ -296,9 +299,8 @@ def test_delivery_fault_matrix(gw, director, platform, fault, request): assert faults_fired(gw, aid), f"injected {kind} never hit a platform call\n{dump(gw, platform, chat)}" if fault == "ack_lost" and STREAMING[platform]: request.applymarker(pytest.mark.xfail(strict=True, reason=STREAM_ACK_LOST_GAP)) - if (fault, platform) in UNLANDED_STREAM_FIX: - request.applymarker(pytest.mark.xfail(strict=UNLANDED_STREAM_FIX[(fault, platform)], - reason="fixed by #120315")) + if (fault, platform) in STREAM_OVERFLOW_GAP_CELLS: + expect_gap(request, 120315, STREAM_OVERFLOW_GAP) assert_exactly_once(gw, director, platform, chat, token, aid, marked_duplicates_allowed=expect == "marked_dupes", first_copy_may_be_marked=expect == "one_copy") @@ -401,7 +403,7 @@ def test_redelivered_inbound_id_processed_once(gw, director, when, request): if when == "after_turn": gw.inject(platform, text, message_id=mid, chat_id=chat) if when == "after_reconnect": - request.applymarker(pytest.mark.xfail(strict=True, reason=RECONNECT_DEDUP_GAP)) + expect_gap(request, 120444, RECONNECT_DEDUP_GAP) before = gw.rpc("reconnect", platform=platform)["id"] wait_until(lambda: gw.rpc("adapter_id", platform=platform)["id"] not in (before, None), "the reconnect watcher to install a fresh adapter", proc=gw.proc, log=gw.log) @@ -412,7 +414,7 @@ def test_redelivered_inbound_id_processed_once(gw, director, when, request): RECONNECT_DEDUP_GAP = ( - "LIVE GAP (#119848 family): inbound de-duplication state lives on the adapter instance " + "LIVE GAP (#119848 family, fixed by #120444): inbound de-duplication state lives on the adapter instance " "(MessageDeduplicator); the gateway's reconnect watcher builds a FRESH adapter, so a platform " "replaying a recent inbound id after the reconnect gets it processed and answered a second time.") @@ -425,15 +427,22 @@ CRASH_POINTS = { "completed_not_ledgered": ("pre_ledger", "hold"), # the final send/edit is in flight and the platform never saw it "send_in_flight": ("any", "hold"), - # the platform accepted the final send/edit; the ack died with the process + # the turn is persisted and the platform accepted the final send/edit; the ack died with the process "sent_ack_lost": ("any", "hold_after"), + # streaming: the platform accepted a preview edit that already shows the complete answer, while + # the provider stream (so the turn) is still open -> nothing persisted when the process dies + "stream_accepted_unpersisted": ("any", "hold_after"), } -# completed_not_ledgered: GatewayRunner._handle_message clears the durable active-turn marker in its -# finally BEFORE the base adapter records the delivery obligation, so a SIGKILL in that window leaves -# neither a resume marker nor a ledger row. Today the reply is still recovered - by the 120 s recency -# fallback re-running the turn (see RECENCY_FALLBACK_GAP). If that fallback is narrowed without -# closing the window, this scenario goes red: the persisted answer is then never delivered. +# completed_not_ledgered: until #120377, GatewayRunner._handle_message clears the durable active-turn +# marker in its finally BEFORE the base adapter records the delivery obligation, so a SIGKILL in that +# window leaves neither a resume marker nor a ledger row and only the 120 s recency fallback recovers +# the reply (by re-running the turn, see RECENCY_FALLBACK_GAP). #120377 removes that fallback and +# hands the persisted reply to the ledger instead; either way exactly one complete reply must show. +SENT_ACK_LOST_STREAM_GAP = ( + "LIVE GAP (fixed by #120377): streaming turn persisted, the platform accepted the final edit, " + "SIGKILL before the ack -> the unclean restart re-runs the already-answered turn (turn marker / " + "120 s recency fallback) and the model's second answer is sent UNMARKED next to the first.") STREAM_CRASH_AFTER_ACCEPT_GAP = ( "LIVE GAP: streaming turn, the platform already shows the complete final answer, SIGKILL before " "the turn is persisted -> restart auto-resumes and the model answers AGAIN; the second answer is " @@ -449,7 +458,8 @@ def _replies(gw: GatewayProcess, platform: str, chat: str, token: str, aid: str) CRASH_MATRIX = [("completed_not_ledgered", "fk_ne"), # streaming finals are sent before the handler returns ("send_in_flight", "fk_ne"), ("send_in_flight", "fk_tg"), - ("sent_ack_lost", "fk_ne"), ("sent_ack_lost", "fk_tg")] + ("sent_ack_lost", "fk_ne"), ("sent_ack_lost", "fk_tg"), + ("stream_accepted_unpersisted", "fk_tg")] @pytest.fixture @@ -475,15 +485,23 @@ def test_crash_between_completion_and_send(fresh_gw, director, point, platform, op, kind = CRASH_POINTS[point] streaming = STREAMING[platform] if point == "sent_ack_lost" and streaming: + expect_gap(request, 120377, SENT_ACK_LOST_STREAM_GAP) + if point == "stream_accepted_unpersisted": request.applymarker(pytest.mark.xfail(strict=True, reason=STREAM_CRASH_AFTER_ACCEPT_GAP)) token = f"{platform}.crash-{point}" chat = f"k-{platform}-{point}" aid = f"A-{token}" - director.script(token, Director.answer(aid, "answer computed before the crash.")) + answer = Director.answer(aid, "answer computed before the crash.") + if point == "stream_accepted_unpersisted": # the whole answer streams, then the stream stays open + answer = StallMidStream(text=answer.text, after_chars=len(answer.text), seconds=60) + director.script(token, answer) gw.fault(platform=platform, op=op, kind=kind, contains=footer(aid), chat_id=chat) gw.inject(platform, f"[in:{token}] question", message_id=f"in-{token}", chat_id=chat) wait_until(lambda: any(footer(aid) in h["content"] for h in gw.holds()), f"{token} to reach the crash point", proc=gw.proc, log=gw.log) + if point == "sent_ack_lost": # the stream's footer edit can land before the turn is persisted + wait_until(lambda: persisted_answers(gw.db_path, aid), f"{token} answer persisted before the kill", + proc=gw.proc, log=gw.log) gw.kill9() gw.start() @@ -518,10 +536,17 @@ def test_crash_between_completion_and_send(fresh_gw, director, point, platform, # Scenarios whose strict xfail documents a live gap: excluded from the whole-run audit below. -KNOWN_GAP_TOKENS = {"fk_tg.ack_lost", "fk_dc.ack_lost", "fk_tg.redeliver-after_reconnect"} +KNOWN_GAP_TOKENS = {"fk_tg.ack_lost", "fk_dc.ack_lost"} RESUME_EXPECTED: set = set() +def known_gap_tokens() -> set: + gaps = set(KNOWN_GAP_TOKENS) + if gap_open(120444): + gaps.add("fk_tg.redeliver-after_reconnect") + return gaps + + def test_zz_whole_run_audit(gw, director): """Across EVERY chat of the module (every scenario, every restart): no inbound ended with two unmarked complete replies (held answer, extra turn or auto-resume), and no already-answered @@ -529,8 +554,9 @@ def test_zz_whole_run_audit(gw, director): visible = gw.platform_view().visible() head = re.compile(r"<<((?:A-|R-)?[A-Za-z0-9_.-]+?)>>") problems = [] + gaps = known_gap_tokens() for token in sorted(director.turns): - if token in KNOWN_GAP_TOKENS or ".crash-" in token: # crash scenarios run on their own homes + if token in gaps or ".crash-" in token: # crash scenarios run on their own homes continue ids = {m.group(1) for v in visible for m in head.finditer(v.text) if m.group(1) in (f"A-{token}",) or m.group(1).startswith((f"R-{token}-", f"{token}-extra"))} @@ -549,7 +575,7 @@ def test_zz_whole_run_audit(gw, director): RECENCY_FALLBACK_GAP = ( - "LIVE GAP: on an unclean startup GatewayRunner._recover_unclean_sessions also runs the legacy " + "LIVE GAP (fixed by #120377): on an unclean startup GatewayRunner._recover_unclean_sessions also runs the legacy " "recency fallback SessionStore.suspend_recently_active(120), which marks EVERY session updated in " "the last 120 s resume_pending with reason 'restart_interrupted' - a reason in _AUTO_RESUME_REASONS - " "so _schedule_resume_pending_sessions synthesizes a turn for sessions whose turn had already " @@ -557,10 +583,10 @@ RECENCY_FALLBACK_GAP = ( "unsolicited second answer.") -@pytest.mark.xfail(strict=True, reason=RECENCY_FALLBACK_GAP) -def test_zzz_unclean_restart_reruns_nothing(gw, director): +def test_zzz_unclean_restart_reruns_nothing(gw, director, request): """Runs LAST on the module gateway: after every scenario above has been answered, an unclean restart with NOTHING in flight must not re-run or re-deliver anything.""" + expect_gap(request, 120377, RECENCY_FALLBACK_GAP) tok = "fk_tg.quiet-restart" director.script(tok, Director.answer(f"A-{tok}", "answered long before the crash.")) gw.inject("fk_tg", f"[in:{tok}] hi", message_id=f"in-{tok}", chat_id="quiet") @@ -569,7 +595,9 @@ def test_zzz_unclean_restart_reruns_nothing(gw, director): "ledger quiescent before the crash", proc=gw.proc, log=gw.log) assert_exactly_once(gw, director, "fk_tg", "quiet", tok, f"A-{tok}") before = {m.message_id for m in gw.platform_view().visible()} - turns = dict(director.turns) + # The Director is module-wide: crash cells on their own homes may legitimately resume their + # in-flight turn, so only what THIS restart adds counts. + turns, resumes = dict(director.turns), dict(director.resumes) gw.kill9() gw.start() # A buggy restart shows its first extra reply within a few seconds of boot; a correct one never @@ -584,7 +612,8 @@ def test_zzz_unclean_restart_reruns_nothing(gw, director): new = [m for m in gw.platform_view().visible() if m.message_id not in before] assert not new, f"an unclean restart with nothing in flight delivered {len(new)} new messages, e.g.:\n" + \ "\n".join(f" {m.chat_id}: {m.text[:100]!r}" for m in new[:8]) - assert director.turns == turns and not director.resumes, f"model re-run after restart: {director.resumes}" + assert director.turns == turns and director.resumes == resumes, ( + f"model re-run after restart: {director.resumes} (before the kill: {resumes})")