Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 13 additions & 6 deletions grapharc/runtime/budget.py
Original file line number Diff line number Diff line change
Expand Up @@ -367,8 +367,8 @@ def deadline_guard(meter: BudgetMeter, *, what: str) -> Iterator[None]:
What this does *not* guarantee:

- Mechanism 2 cannot interrupt a thread parked inside a C call: a
`time.sleep(60)` sleeps out its 60 seconds and raises on return. This is
not a fan-out-only weakness. Mechanism 1 needs `invoke()` to be on the
`time.sleep(60)` is not interrupted mid-call, but the guard still raises
on exit. This is not a fan-out-only weakness. Mechanism 1 needs `invoke()` to be on the
process's main thread, so *any* run driven from a worker thread — every
request handler in a threaded server, every `ThreadPoolExecutor` caller —
falls back to mechanism 2 for the whole run, nodes and fan-out alike.
Expand All @@ -382,9 +382,10 @@ def deadline_guard(meter: BudgetMeter, *, what: str) -> Iterator[None]:
- Like any asynchronous exception, the interrupt lands wherever the node
happened to be: it is as safe as Ctrl-C, no safer.

Short of that last case the ceiling is honoured at the node boundary: if the
deadline passed and the node swallowed the exception, this guard raises on
exit, so the node's writes never reach state.
Short of the never-returns case the ceiling is honoured at the node
boundary: if the deadline passed — whether the node swallowed the exception
or no interrupt was ever delivered — this guard raises on exit, so the
node's writes never reach state.
"""
remaining = meter.remaining_seconds()
if remaining is None:
Expand Down Expand Up @@ -486,5 +487,11 @@ def disarm() -> None:
disarm()
except NodeDeadlineExceeded as exc:
raise NodeDeadlineExceeded(detail()) from exc
if state["fired"]:

# Reached only when the node returned normally. It may have swallowed the
# interrupt, or the deadline may have passed without the timer firing —
# the timer thread needs the GIL, which a node inside a long C call
# withholds until it returns; either way its writes must not land.
left = meter.remaining_seconds()
if state["fired"] or (left is not None and left <= 0):
raise NodeDeadlineExceeded(detail())
46 changes: 46 additions & 0 deletions tests/test_budget_enforcement.py
Original file line number Diff line number Diff line change
Expand Up @@ -542,6 +542,52 @@ def body():
assert ran_for < 2.0, "the node asked for 5s and was not stopped"


def test_an_overrun_is_refused_at_exit_even_if_the_timer_never_fired(monkeypatch):
"""A node that holds the GIL through the deadline — a long C call, or plain
timer-scheduling latency — denies the timer thread its turn: `fire()` never
runs, the node returns normally, and an exit check that tests only
`state["fired"]` lets the overrun's writes land. The contract is the node
boundary, so the exit check itself must notice the spent deadline.

The timer is replaced with one that never fires, which makes this the
deterministic statement of that contract: asserting on whether a real
timer's async exception got delivered in time is a race (see
`_run_swallower` above), whereas the exit check runs unconditionally.
"""

class NeverFires:
"""`threading.Timer`'s surface as the guard uses it, minus the firing."""

def __init__(self, interval, function):
self.daemon = False

def start(self):
pass

def cancel(self):
pass

monkeypatch.setattr(threading, "Timer", NeverFires)
outcome: dict[str, object] = {}

def body(): # a worker thread uses mechanism 2, like any threaded server
meter = BudgetMeter(Budget(max_seconds=0.05))
try:
with deadline_guard(meter, what="node 'n'"):
time.sleep(0.2) # outlast the deadline; nothing interrupts it
outcome["raised"] = None
except NodeDeadlineExceeded as exc:
outcome["raised"] = exc

worker = threading.Thread(target=body)
worker.start()
worker.join(timeout=10)
assert not worker.is_alive()
assert isinstance(outcome["raised"], NodeDeadlineExceeded), (
"the node overran, no interrupt fired, and the guard let its writes land"
)


def test_a_node_that_finishes_in_time_is_left_alone():
def brisk(state):
time.sleep(0.05)
Expand Down
Loading