diff --git a/tests/e2e/core/kanban/_helpers.py b/tests/e2e/core/kanban/_helpers.py index 81317b7a5f..9e82ceef90 100644 --- a/tests/e2e/core/kanban/_helpers.py +++ b/tests/e2e/core/kanban/_helpers.py @@ -42,12 +42,12 @@ def pid_alive(pid: Optional[int]) -> bool: if not pid: return False try: - os.kill(int(pid), 0) + os.kill(int(pid), 0) # windows-footgun: ok — Linux-gated (module skips off Linux) except (ProcessLookupError, PermissionError): return False # A zombie child of ours still answers kill(0); /proc tells the truth. try: - stat = Path(f"/proc/{int(pid)}/stat").read_text() + stat = Path(f"/proc/{int(pid)}/stat").read_text(encoding="utf-8") return stat.rsplit(")", 1)[1].split()[0] != "Z" except OSError: return False @@ -133,7 +133,7 @@ class Board: def kill_workers(self) -> None: for pid in list(self.spawned_pids): try: - os.kill(pid, signal.SIGKILL) + os.kill(pid, signal.SIGKILL) # windows-footgun: ok — Linux-gated (module skips off Linux) except (ProcessLookupError, PermissionError): pass diff --git a/tests/e2e/core/kanban/test_kanban_decompose_billing.py b/tests/e2e/core/kanban/test_kanban_decompose_billing.py index e9f1e2593a..e87e3cceb0 100644 --- a/tests/e2e/core/kanban/test_kanban_decompose_billing.py +++ b/tests/e2e/core/kanban/test_kanban_decompose_billing.py @@ -141,11 +141,11 @@ def gateway(b: Board) -> Iterator[subprocess.Popen]: try: yield proc finally: - for sig, grace in ((signal.SIGTERM, 20), (signal.SIGKILL, 10)): + for sig, grace in ((signal.SIGTERM, 20), (signal.SIGKILL, 10)): # windows-footgun: ok — POSIX-gated (skips on win32) if proc.poll() is not None: break try: - os.killpg(proc.pid, sig) + os.killpg(proc.pid, sig) # windows-footgun: ok — POSIX-gated (skips on win32) except ProcessLookupError: break try: @@ -156,7 +156,7 @@ def gateway(b: Board) -> Iterator[subprocess.Popen]: def gateway_tail(b: Board) -> str: p = b.root / "gateway.log" - return p.read_text(errors="replace")[-3000:] if p.exists() else "(no gateway log)" + return p.read_text(encoding="utf-8", errors="replace")[-3000:] if p.exists() else "(no gateway log)" def prove_ticks(b: Board, proc: subprocess.Popen, k: int) -> None: diff --git a/tests/e2e/core/kanban/test_kanban_dispatcher_restart.py b/tests/e2e/core/kanban/test_kanban_dispatcher_restart.py index 0dc66af54f..096d9d58b2 100644 --- a/tests/e2e/core/kanban/test_kanban_dispatcher_restart.py +++ b/tests/e2e/core/kanban/test_kanban_dispatcher_restart.py @@ -12,6 +12,9 @@ Invariants, read from kanban.db and the provider's request log: * every card ends ``done`` with exactly one ``completed`` run and one ``completed`` event; * each card's ``kanban_complete`` was billed exactly once (no duplicate worker ran a card); * no card ever had two runs open at the same time. + +A second scenario gives a live worker a claim TTL shorter than its provider call: the dispatcher +must extend that claim on every expired tick, never reclaim the card and spawn a duplicate. """ from __future__ import annotations @@ -74,7 +77,7 @@ def _kill_dispatcher_after(board: Board, kind: str | None) -> None: wait_until(lambda: _event_count(board, kind) > before or proc.poll() is not None, 60, f"dispatcher to write a {kind} event", interval=0.01) finally: - proc.send_signal(signal.SIGKILL) + proc.send_signal(signal.SIGKILL) # windows-footgun: ok — Linux-gated (module skips off Linux) proc.wait(timeout=30) for row in board.tasks(): if row["worker_pid"]: @@ -123,3 +126,64 @@ def test_dispatcher_sigkill_mid_tick_never_destroys_or_duplicates_cards(tmp_path assert outcomes["reclaimed"] >= 1, outcomes finally: board.kill_workers() + + +class SlowModel: + """The worker's first call hangs on the provider (no chunk, no tool, so no heartbeat) until + ``release``; every first call is counted so a duplicate worker on the card shows up as a bill.""" + + def __init__(self) -> None: + self.first_calls = 0 + self.hanging = threading.Event() + self.release = threading.Event() + self._lock = threading.Lock() + + def __call__(self, rec: dict): + if rec["body"]["messages"][-1].get("role") == "tool": + return Text("slow card closed") + with self._lock: + self.first_calls += 1 + self.hanging.set() + self.release.wait(120) + return ToolCall("kanban_complete", {"summary": "slow model finally answered"}) + + +def test_ttl_expiry_extends_a_live_hung_workers_claim_instead_of_respawning(tmp_path) -> None: + """A claim TTL shorter than one provider call: every tick past expiry must extend the live + worker's claim (``claim_extended``), never reclaim it and spawn a second worker beside it.""" + ttl = 2 + model = SlowModel() + with FakeLLMServer(model) as srv: + board = Board(tmp_path, srv.base_url, env_extra={"HERMES_KANBAN_CLAIM_TTL_SECONDS": str(ttl)}) + try: + tid = board.create("slow model card") + board.dispatch() + pid = int(board.task(tid)["worker_pid"]) + wait_until(model.hanging.is_set, 90, f"worker to hang on its provider call\n{board.diag(tid)}") + first_expiry = int(board.task(tid)["claim_expires"]) + # Tick until two TTL windows have passed and two ticks saw the claim expired, or the + # dispatcher gave the card away (then the assertions below name what went wrong). + def settled() -> bool: + board.dispatch() + if board.events(tid, "reclaimed") or len(board.events(tid, "spawned")) > 1: + return True + return (int(time.time()) > first_expiry + 2 * ttl + and len(board.events(tid, "claim_extended")) >= 2) + try: + wait_until(settled, 60, "two ticks past the claim TTL", interval=0.2) + except AssertionError as exc: + raise AssertionError(f"{exc}\n{board.diag(tid)}") from None + assert not board.events(tid, "reclaimed"), f"live worker's claim reclaimed\n{board.diag(tid)}" + task = board.task(tid) + assert pid_alive(pid) and task["worker_pid"] == pid, board.diag(tid) + assert task["status"] == "running" and int(task["claim_expires"]) > first_expiry, board.diag(tid) + assert [e["payload"]["pid"] for e in board.events(tid, "spawned")] == [pid], board.diag(tid) + assert all(e["payload"]["worker_pid"] == pid for e in board.events(tid, "claim_extended")) + assert model.first_calls == 1, f"a second worker billed the card\n{board.diag(tid)}" + model.release.set() + board.wait_worker_exit(tid, pid) + wait_until(lambda: board.task(tid)["status"] == "done", 30, f"card done\n{board.diag(tid)}") + assert [r["outcome"] for r in board.runs(tid)] == ["completed"], board.diag(tid) + finally: + model.release.set() + board.kill_workers() diff --git a/tests/e2e/core/kanban/test_kanban_rate_limit_review.py b/tests/e2e/core/kanban/test_kanban_rate_limit_review.py index 50331d5b6b..e51ec23e1a 100644 --- a/tests/e2e/core/kanban/test_kanban_rate_limit_review.py +++ b/tests/e2e/core/kanban/test_kanban_rate_limit_review.py @@ -3,9 +3,9 @@ card's later review handoff (#119070). Real processes end to end: every tick is a real ``hermes kanban dispatch`` process, every worker a real ``hermes chat -q`` process spawned by it, talking to the recording fake provider (the only -fake — it stands in for the vendor HTTP API). The fake plays two roles, told apart by what the -dispatcher put in the worker's system prompt: a review-lane worker is started with the bundled -``sdlc-review`` skill preloaded, an implementer is not. +fake — it stands in for the vendor HTTP API). The fake plays two roles, told apart structurally: +the scratch home carries its own ``sdlc-review`` skill whose body holds a unique marker, and the +review lane preloads that skill into the worker it spawns, so only a reviewer's prompt carries it. Flow under test (``HERMES_KANBAN_RATE_LIMIT_COOLDOWN_SECONDS=0``, ``agent.api_max_retries: 1``): @@ -20,6 +20,7 @@ result, and the request stream the fake provider recorded — never from log wor from __future__ import annotations +import json import sys import threading from dataclasses import dataclass, field @@ -37,9 +38,11 @@ pytestmark = [ pytest.mark.live_system_guard_bypass, ] -# The review lane preloads the bundled review skill; the dispatcher's own contract (kanban.md: -# "spawn the assigned profile with the bundled sdlc-review skill"). +# The review lane preloads the review skill by name (kanban.md: "spawn the assigned profile with +# the bundled sdlc-review skill"). A same-named user skill wins over the bundled copy, so seeding one +# with a marker in the scratch home tells a reviewer's prompt apart without reading prompt wording. REVIEW_SKILL = "sdlc-review" +REVIEW_MARK = "E2E_REVIEW_LANE_SKILL_5d21c9" RATE_LIMIT_EXIT_CODE = 75 # KANBAN_RATE_LIMIT_EXIT_CODE — documented worker exit contract MAX_TICKS = 6 # a healthy rate-limited card is done on tick 3; the rest prove "forever" @@ -62,13 +65,12 @@ def _known(name: str): # fake provider ----------------------------------------------------------------------------------- -def _system_text(body: dict) -> str: - parts = [] - for m in body.get("messages", []): - if m.get("role") == "system": - c = m.get("content") - parts.append(c if isinstance(c, str) else " ".join(p.get("text", "") for p in c or [])) - return "\n".join(parts) +def _seed_review_skill(board: Board) -> None: + skill = board.hermes_home / "skills" / "devops" / REVIEW_SKILL + skill.mkdir(parents=True, exist_ok=True) + (skill / "SKILL.md").write_text( + f"---\nname: {REVIEW_SKILL}\ndescription: Review Kanban handoffs (e2e stand-in).\n---\n\n" + f"# {REVIEW_SKILL}\n\nFollow {REVIEW_MARK} before approving.\n", encoding="utf-8") @dataclass @@ -95,7 +97,7 @@ class TwoRoleModel: def __call__(self, rec: dict[str, Any]) -> Any: body = rec["body"] msgs = body.get("messages", []) - role = "review" if f'"{REVIEW_SKILL}"' in _system_text(body) else "impl" + role = "review" if REVIEW_MARK in json.dumps(msgs) else "impl" with self._lock: last = self.attempts[-1] if self.attempts else None if last is None or (last.tick, last.role) != (self.tick, role): @@ -120,6 +122,13 @@ class TwoRoleModel: return [a for a in self.attempts if a.role == role] +def one_terminal_turn(att: Attempt, tool: str) -> bool: + """Billing bound for a finishing attempt: its first answer IS the terminal board call, followed by + at most one closing turn (dropping that closing call is an improvement, not a regression).""" + return (att.answers[:1] == [tool] and att.requests == len(att.answers) <= 2 + and all(a == "Text" for a in att.answers[1:])) + + # board driving ----------------------------------------------------------------------------------- @@ -164,6 +173,7 @@ def _flow(root: Path, rate_limited_attempts: int) -> Any: model = TwoRoleModel(rate_limited_attempts) with FakeLLMServer(model) as srv: board = Board(root, srv.base_url) + _seed_review_skill(board) tid = board.create(f"rate-limit review flow ({rate_limited_attempts} x 429)") try: ticks = drive(board, tid, model) @@ -199,11 +209,10 @@ def test_rate_limited_attempt_is_billed_once_and_requeued_without_a_failure(rate # Cooldown is 0: the reap and the respawn happen on the SAME tick right after the 429 worker died. assert [t["spawned"] for t in f.ticks[:2]] == [1, 1], diag # Provider-side billing: the 429 attempt is ONE request (api_max_retries: 1, no retry loop and - # no fallback re-send); the retry is exactly the handoff tool call plus its closing turn. + # no fallback re-send); the retry opens with the handoff call and at most one closing turn. impl = f.model.by_role("impl") - assert [(a.tick, a.requests, a.answers) for a in impl] == [ - (1, 1, ["Error"]), (2, 2, ["kanban_request_review", "Text"]), - ], diag + assert [(a.tick, a.requests, a.answers) for a in impl[:1]] == [(1, 1, ["Error"])], diag + assert [a.tick for a in impl] == [1, 2] and one_terminal_turn(impl[1], "kanban_request_review"), diag assert len(runs) - 2 == len(f.model.by_role("review")), diag @@ -216,7 +225,9 @@ def test_clean_handoff_spawns_the_reviewer_on_the_next_tick(tmp_path: Path) -> N assert f.run_outcomes() == ["review_requested", "completed"], diag assert [t["spawned"] for t in f.ticks] == [1, 1], diag assert not any(t["guarded"] for t in f.ticks), diag - assert [(a.role, a.tick, a.requests) for a in f.model.attempts] == [("impl", 1, 2), ("review", 2, 2)], diag + assert [(a.role, a.tick) for a in f.model.attempts] == [("impl", 1), ("review", 2)], diag + impl, review = f.model.attempts + assert one_terminal_turn(impl, "kanban_request_review") and one_terminal_turn(review, "kanban_complete"), diag assert b.task(tid)["consecutive_failures"] == 0, diag @@ -236,6 +247,6 @@ def test_rate_limited_then_review_handoff_reaches_the_reviewer(rate_limited_flow # Once fixed, the whole contract must hold, not just "something spawned". assert f.run_outcomes() == ["rate_limited", "review_requested", "completed"], diag # The reviewer starts on the tick right after the handoff (cooldown 0), billed one tool turn. - assert [(a.tick, a.requests, a.answers) for a in reviewers] == [(3, 2, ["kanban_complete", "Text"])], diag + assert [a.tick for a in reviewers] == [3] and one_terminal_turn(reviewers[0], "kanban_complete"), diag assert "blocker_auth" not in guarded, diag assert not b.events(tid, "gave_up") and b.task(tid)["consecutive_failures"] == 0, diag diff --git a/tests/e2e/core/kanban/test_kanban_worker_contract.py b/tests/e2e/core/kanban/test_kanban_worker_contract.py index a203911222..0f429241f7 100644 --- a/tests/e2e/core/kanban/test_kanban_worker_contract.py +++ b/tests/e2e/core/kanban/test_kanban_worker_contract.py @@ -11,6 +11,7 @@ Verdicts come from kanban.db rows, files on disk and the provider's request log. from __future__ import annotations +import json import re import shutil import sys @@ -33,6 +34,8 @@ KNOWN: dict[str, str] = { } _TASK_RE = re.compile(r"work kanban task (t_[0-9a-f]+)") +_EXIT_RE = re.compile(r"\[kanban-worker-exit\] rc=(-?\d+)") +_REFUSAL_KIND_RE = re.compile(r"block|refus|reject|artifact|violation") SKILL_MARK = "E2E_PINNED_SKILL_BODY_7f3a" @@ -70,10 +73,16 @@ def _artifact_location(board: Board, tid: str, where: str) -> Path: ]) def test_declared_artifact_is_attached_or_reported(tmp_path, where: str) -> None: board_ref: dict[str, Board] = {} + blocked: dict[str, bool] = {} payload = f"DELIVERABLE_{where.upper()}_c0ffee\n" def responder(rec: dict): - if rec["body"]["messages"][-1].get("role") == "tool": + last = rec["body"]["messages"][-1] + if last.get("role") == "tool": + # A refused completion is answered the way a careful worker would: park the card. + if '"error"' in str(last.get("content")) and not blocked.get("sent"): + blocked["sent"] = True + return ToolCall("kanban_block", {"reason": "artifact could not be delivered"}) return Text("delivered") tid = _task_id(rec) path = _artifact_location(board_ref["b"], tid, where) @@ -93,6 +102,8 @@ def test_declared_artifact_is_attached_or_reported(tmp_path, where: str) -> None raise KnownGap(f"card done with its declared {where}-workspace artifact never attached\n" f"{board.diag(tid)}") assert status == "done" or where == "outside", board.diag(tid) + if status != "done": + _assert_visible_refusal(board, srv, tid, _artifact_location(board, tid, where)) if attached: assert len(attached) == 1 and stored[0].read_text(encoding="utf-8") == payload, attached assert not stored[0].is_relative_to(board.hermes_home / "kanban" / "workspaces"), stored @@ -100,6 +111,23 @@ def test_declared_artifact_is_attached_or_reported(tmp_path, where: str) -> None board.kill_workers() +def _assert_visible_refusal(board: Board, srv: FakeLLMServer, tid: str, artifact: Path) -> None: + """A completion that did not land must say why, naming the artifact, somewhere a human or the + worker sees it; and the refusal must not turn into a crash loop.""" + board.dispatch("--max", "0") # reap the exited worker so its run is booked + name = artifact.name + tool_errors = [str(m.get("content")) for body in srv.main_requests() for m in body["messages"] + if m.get("role") == "tool" and '"error"' in str(m.get("content"))] + traces = [c for c in tool_errors if name in c] + traces += [e["kind"] for e in board.events(tid) + if _REFUSAL_KIND_RE.search(e["kind"]) and name in json.dumps(e["payload"])] + traces += [r["error"] for r in board.runs(tid) if name in (r["error"] or "")] + assert traces, f"card left {board.task(tid)['status']} with no refusal naming {name}\n{board.diag(tid)}" + runs = board.runs(tid) + assert len(runs) == 1 and runs[0]["outcome"] != "crashed", board.diag(tid) + assert not board.events(tid, "gave_up") and board.task(tid)["consecutive_failures"] <= 1, board.diag(tid) + + # pinned skills -------------------------------------------------------------------------------- @@ -127,7 +155,8 @@ def test_worker_with_resolvable_pinned_skill_sees_it_on_first_request(tmp_path) first = srv.main_requests()[0] assert SKILL_MARK in str(first["messages"]), "pinned skill body never reached the model" assert board.task(tid)["status"] == "done", board.diag(tid) - assert len(srv.main_requests()) == 2 + # The terminal call is the first answer; at most one closing turn follows it. + assert 1 <= len(srv.main_requests()) <= 2, len(srv.main_requests()) finally: board.kill_workers() @@ -143,9 +172,12 @@ def test_worker_with_unresolvable_pinned_skill_still_starts_its_session(tmp_path # The operator removes the skill after the pin was written. shutil.rmtree(board.hermes_home / "skills" / "e2e-archived") _run_one_card(board, tid) - wait_until(lambda: board.worker_log(tid), 10, "worker log") - if not srv.main_requests(): - raise KnownGap(f"worker died before its session; board:\n{board.diag(tid)}") + log = wait_until(lambda: board.worker_log(tid), 10, "worker log") + exits = [int(rc) for rc in _EXIT_RE.findall(log)] + # The bug's own signature: no model call at all, and the worker died of the stale pin + # (nonzero exit trailer, or its log names the missing skill). Anything else stays red. + if not srv.main_requests() and ((exits and exits[-1] != 0) or "e2e-archived" in log): + raise KnownGap(f"worker died before its session (exits={exits}); board:\n{board.diag(tid)}") assert SKILL_MARK not in str(srv.main_requests()[0]["messages"]), "removed skill still loaded" assert board.task(tid)["status"] == "done", board.diag(tid) finally: diff --git a/tests/e2e/core/kanban/test_kanban_worker_sigkill.py b/tests/e2e/core/kanban/test_kanban_worker_sigkill.py index 84c6ec83ae..981e124b77 100644 --- a/tests/e2e/core/kanban/test_kanban_worker_sigkill.py +++ b/tests/e2e/core/kanban/test_kanban_worker_sigkill.py @@ -111,7 +111,7 @@ def _drive(board: Board, director: Director) -> Scenario: wait_until(director.hanging.is_set, 90, f"attempt 2 to hang on its provider call\n{board.diag(tid)}") hb = board.task(tid)["last_heartbeat_at"] assert hb, f"attempt 2 never heartbeat\n{board.diag(tid)}" - os.kill(w2, signal.SIGKILL) + os.kill(w2, signal.SIGKILL) # windows-footgun: ok — Linux-gated (module skips off Linux) wait_until(lambda: not pid_alive(w2), 15, "SIGKILLed worker to disappear") # The fresh claim must start strictly after attempt 2's last heartbeat second. wait_until(lambda: int(time.time()) > int(hb) + 1, 5, "clock to pass the last heartbeat") @@ -158,9 +158,10 @@ def test_sigkilled_worker_is_reclaimed_and_the_retry_completes_once(scenario: Sc task = b.task(sc.tid) assert task["status"] == "done" and task["worker_pid"] is None and task["claim_lock"] is None assert [r for r in runs if r["outcome"] == "completed"][0]["summary"] == "ATTEMPT_THREE_DONE" - # Billing: every attempt paid only for its own turns, nothing ran after the card closed. + # Billing: every attempt paid only for its own turns, nothing ran after the card closed. Attempt 2 + # is exactly heartbeat + the hung call; attempt 3 opens with kanban_complete, then <= 1 closing turn. billed = sc.director.billed - assert sorted(billed) == [1, 2, 3] and billed[2] == 2 and billed[3] == 2, billed + assert sorted(billed) == [1, 2, 3] and billed[2] == 2 and 1 <= billed[3] <= 2, billed assert all(not t["spawned"] for t in sc.ticks_after_done), sc.ticks_after_done assert len(b.events(sc.tid, "completed")) == 1