diff --git a/docs/deep-dive.md b/docs/deep-dive.md index 2e790f9..dbfad0e 100644 --- a/docs/deep-dive.md +++ b/docs/deep-dive.md @@ -254,7 +254,7 @@ A stable system is not one that claims to have no edges — it is one whose edge - **`.env` and `grapharc.toml` follow the same discovery rule: the working directory, and nowhere else.** Neither searches parent directories — a run must not be governed by a file you did not know about, and must not be *billed* to one either. **This is a behaviour change:** the credential loader used to walk up to `/`, so a `.env` in an ancestor directory (a `$HOME` one on a shared box, a client project one above a demo checkout) was picked up silently. If you relied on that, move the file into the directory you run from, `export` the variable, or pass `env_file=` to name it explicitly. A real environment variable still beats any file. - **`grapharc run` has no budget unless you give it one.** Set any of `--max-tokens`, `--max-iterations`, `--max-seconds`, or `--max-concurrency`; without them each dimension is unlimited and the gate admits a topology of any worst-case cost. -**Verified this pass:** `pytest` → green, 2,189 selected and 13 deselected (the live ones); `ruff check .` clean; all eight `grapharc demo` stages green, plus the `trace` / `metrics` / `viz` / `replay` tour against a freshly recorded demo trace; the wheel builds and imports all submodules in a clean virtualenv with `[all]`, and `0.1.7` on PyPI is that wheel. The counts are a snapshot, not a property of the project — `pytest` re-derives them in one command, which is the only reason they are quoted, and `tests/test_deep_dive.py` fails this line rather than letting it drift. +**Verified this pass:** `pytest` → green, 2,190 selected and 13 deselected (the live ones); `ruff check .` clean; all eight `grapharc demo` stages green, plus the `trace` / `metrics` / `viz` / `replay` tour against a freshly recorded demo trace; the wheel builds and imports all submodules in a clean virtualenv with `[all]`, and `0.1.7` on PyPI is that wheel. The counts are a snapshot, not a property of the project — `pytest` re-derives them in one command, which is the only reason they are quoted, and `tests/test_deep_dive.py` fails this line rather than letting it drift. [ROADMAP.md](../ROADMAP.md) tracks what is built and what is not, item by item. diff --git a/grapharc/runtime/budget.py b/grapharc/runtime/budget.py index 67fdb76..ef36758 100644 --- a/grapharc/runtime/budget.py +++ b/grapharc/runtime/budget.py @@ -333,6 +333,13 @@ def snapshot(self) -> dict[str, float | int]: # node alive; the cost while a node is being torn down is one timer per 50ms. _REARM_SECONDS = 0.05 +#: How many times a mechanism-2 teardown may be interrupted before it stops +#: retrying and clears the flag outside the lock. Two is already generous -- +#: once `armed` is false `fire` returns without re-arming, so no further +#: exception is queued -- and the loop exists for the interrupt that lands +#: *inside* the teardown, not for a steady stream of them. +_DISARM_ATTEMPTS = 8 + # The longest delay both mechanisms can actually be armed with. `setitimer` # raises `OverflowError` past the platform's `time_t` (~2**31 seconds), and # `threading.Timer` accepts a larger value but crashes its own thread once the @@ -487,11 +494,39 @@ def rearm(delay: float) -> None: rearm(armable) def disarm() -> None: - with lock: - state["armed"] = False - state["timer"].cancel() - if state["fired"]: - _async_raise(thread_id, None) + # An interrupt can land *inside* this teardown, and used to leave a + # timer running for the life of the thread. `fire` queues the async + # exception while holding `lock`, so a guard already blocked on + # `lock` here is handed it the moment it acquires the lock -- at the + # next bytecode, which is before `armed` is cleared and before the + # re-armed timer is cancelled. Letting that propagate left a live + # 50ms timer whose `fire` still read `armed` as true, so it raised + # `NodeDeadlineExceeded` into this thread every 50ms, indefinitely, + # long after the run that armed it had finished. On a pooled thread + # that is an unattributable crash in whatever ran next -- the exact + # failure `test_no_interrupt_survives_the_node_that_earned_it` + # exists to rule out, arriving by a path it did not cover. + # + # So the teardown is retried rather than abandoned. Swallowing the + # interrupt costs nothing: the guard decides the outcome from + # `state["fired"]` once this returns, and raises on it. + for _ in range(_DISARM_ATTEMPTS): + try: + with lock: + state["armed"] = False + state["timer"].cancel() + if state["fired"]: + _async_raise(thread_id, None) + return + except NodeDeadlineExceeded: + # An interrupt landing here *is* the deadline firing. + state["fired"] = True + # Last resort. Clearing `armed` is the single store that stops + # `fire` re-arming, so it is done outside the lock rather than + # risking another interrupt on the way to it. + state["armed"] = False + state["timer"].cancel() + _async_raise(thread_id, None) try: try: diff --git a/tests/test_budget_enforcement.py b/tests/test_budget_enforcement.py index 1a29b87..7d50e5c 100644 --- a/tests/test_budget_enforcement.py +++ b/tests/test_budget_enforcement.py @@ -35,12 +35,14 @@ import threading import time from typing import Annotated +from unittest import mock import pytest from langchain_core.messages import AIMessage, HumanMessage from pydantic import BaseModel from grapharc.runtime.budget import ( + _REARM_SECONDS, Budget, BudgetExceeded, BudgetMeter, @@ -741,3 +743,107 @@ def spend(state: State) -> dict: assert metrics.errors == 1 # The cost report and the audit trail must never disagree. assert replay(trace, run_id).tokens == metrics.tokens + + +def test_an_interrupt_landing_inside_the_teardown_leaves_no_timer_running(): + """The teardown race, and the leak it left behind. + + `fire` queues its async exception *while holding the lock*, so a guard + already blocked on that same lock inside `disarm` is handed the exception + the moment it acquires it — at the next bytecode, which is before `armed` + is cleared and before the re-armed timer is cancelled. `disarm` then + propagated, leaving a live 50 ms timer whose `fire` still read `armed` as + true: it re-raised into this thread every 50 ms for the life of the thread, + long after the run that armed it had finished. On a pooled thread that is + an unattributable crash in whatever ran next. + + Forcing the interleaving needs the exception delivered at exactly that + point, so the lock is wrapped and raises on the guard thread's *second* + entry — the first is the initial arm, the second is `disarm`. It releases + before raising, because a real async exception lands inside the `with` body + and that block's exit releases the lock; raising from `__enter__` instead + would hold the lock forever and deadlock the very timer under test rather + than letting it spin. + + The harm is then measured the way a caller feels it: whether anything is + still interrupting this thread once the guard has been released. Without the + retry in `disarm` this records six further interrupts; with it, none. + """ + from grapharc.runtime import budget as budget_module + + worker = {} + calls: list[object] = [] + has_fired = threading.Event() + real_lock = threading.Lock + + # Recorded, not delivered: injecting into the test runner's own thread would + # surface as an unrelated crash somewhere later in the session. + monkey = mock.patch.object( + budget_module, "_async_raise", lambda thread_id, exc: calls.append(exc) + ) + + class InterruptingLock: + """A lock that delivers the deadline interrupt inside `disarm`.""" + + def __init__(self) -> None: + self._lock = real_lock() + self._guard_entries = 0 + + def __getattr__(self, name): # Condition and Event poke at locked() etc. + return getattr(self._lock, name) + + def acquire(self, *args, **kwargs): + return self._lock.acquire(*args, **kwargs) + + def release(self): + return self._lock.release() + + def __enter__(self): + self._lock.acquire() + if threading.get_ident() == worker.get("id"): + self._guard_entries += 1 + if self._guard_entries == 2 and has_fired.is_set(): + self._lock.release() + raise NodeDeadlineExceeded("delivered inside the teardown") + return self + + def __exit__(self, *exc_info): + self._lock.release() + return False + + class NotingTimer(threading.Timer): + """Records that `fire` has run at least once, so the interrupt is + delivered to a teardown that actually has a re-armed timer to lose.""" + + def run(self): + has_fired.set() + return super().run() + + def body(): + worker["id"] = threading.get_ident() + meter = BudgetMeter(Budget(max_seconds=0.1)) + try: + with deadline_guard(meter, what="worker"): + time.sleep(0.4) # past the deadline, so the timer fires and re-arms + except NodeDeadlineExceeded: + pass + + with ( + monkey, + mock.patch.object(budget_module.threading, "Lock", InterruptingLock), + mock.patch.object(budget_module.threading, "Timer", NotingTimer), + ): + thread = threading.Thread(target=body) + thread.start() + thread.join(timeout=10) + assert not thread.is_alive() + assert has_fired.is_set(), "the timer never fired; the race was not exercised" + + settled = len(calls) + time.sleep(6 * _REARM_SECONDS) + + assert len(calls) == settled, ( + f"{len(calls) - settled} interrupt(s) queued after the guard was " + "released: a re-armed timer outlived its teardown and will keep " + "raising into this thread" + )