From 6c6381855215e75f06960b857974d153ff23ba53 Mon Sep 17 00:00:00 2001 From: Hongsheng Liu Date: Fri, 2 Oct 2026 19:03:19 +0800 Subject: [PATCH 1/2] Unwind runner resources in reverse order on shutdown (#37) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Registrations are effects: resources register a named disposer in creation order and shutdown unwinds them in reverse, under the lifetime lease — the next runner never meets a half-closed store. - core.teardown.Teardown: registry that pops disposers in reverse, runs every step even when one raises, reports failures by name only (never exception text, which can carry private data), idempotent. - TaskLoop.close: drain an in-flight summary bounded by the summary budget (a hung provider costs the wait, not the shutdown), then drop the per-task single-flight locks. Safe to call twice. - _wiring(teardown) registers task-store, activity-log, memory-store, task-loop closes; _run_runner unwinds them in a finally inside the RunnerLease for both serve and --once paths. Previously nothing closed the SQLite connections after serve returned. - runner_control is unchanged: its __exit__ already unwinds in correct reverse order and its behavior is pinned by its own tests. Cooperative-stop semantics are untouched: stop_runner still waits for real lock release and reports a pending stop on timeout. Closes #37 --- src/nanodot/cli.py | 54 ++++++++---- src/nanodot/core/runner.py | 21 +++++ src/nanodot/core/teardown.py | 37 +++++++++ tests/test_public_mode.py | 9 +- tests/test_teardown.py | 156 +++++++++++++++++++++++++++++++++++ 5 files changed, 256 insertions(+), 21 deletions(-) create mode 100644 src/nanodot/core/teardown.py create mode 100644 tests/test_teardown.py diff --git a/src/nanodot/cli.py b/src/nanodot/cli.py index e568afa..01424da 100644 --- a/src/nanodot/cli.py +++ b/src/nanodot/cli.py @@ -22,6 +22,7 @@ from nanodot import __version__ from nanodot.paths import data_home +from nanodot.core.teardown import Teardown from nanodot.core.tasks import ( DEFAULT_NOTIFICATION_CONDITIONS, DEFAULT_STOP_CONDITIONS, ) @@ -116,7 +117,7 @@ def build_parser() -> argparse.ArgumentParser: # -- shared wiring ----------------------------------------------------------- -def _wiring() -> tuple: +def _wiring(teardown: Teardown | None = None) -> tuple: from nanodot.core.activity import ActivityLog from nanodot.core.config import Config from nanodot.core.memory import MemoryStore @@ -144,6 +145,13 @@ def _wiring() -> tuple: store, fetcher, sink, activity, provider=configured_provider(), memory=memory, ) + if teardown is not None: + # Registrations are effects: creation order here, reverse unwind in + # Teardown.run — the loop drains before its stores close. + teardown.register("task-store", store.close) + teardown.register("activity-log", activity.close) + teardown.register("memory-store", memory.close) + teardown.register("task-loop", loop.close) return secrets, store, activity, sink, fetcher, loop @@ -549,6 +557,7 @@ def _run_runner(args: argparse.Namespace) -> int: from nanodot.native.runner_control import RunnerControlError, RunnerLease stop = threading.Event() + teardown = Teardown() def _sigint(_signum, _frame) -> None: stop.set() @@ -558,28 +567,37 @@ def _sigint(_signum, _frame) -> None: def prepare() -> None: nonlocal store, daemon - _, store, _, _, _, loop = _wiring() + _, store, _, _, _, loop = _wiring(teardown) daemon = RunnerDaemon(loop, store) try: with RunnerLease(_pidfile(), stop, prepare=prepare): - assert store is not None and daemon is not None - if args.once: - attempted = daemon.tick(stop=stop) - blockers = [t for t in store.list() if t.blocker] - if blockers: - for task in blockers: - print(f"blocked: {task.id} ({task.target}): {task.blocker}", - file=sys.stderr) - return 1 - print(f"ran {attempted} task(s)") + try: + assert store is not None and daemon is not None + if args.once: + attempted = daemon.tick(stop=stop) + blockers = [t for t in store.list() if t.blocker] + if blockers: + for task in blockers: + print(f"blocked: {task.id} ({task.target}): {task.blocker}", + file=sys.stderr) + return 1 + print(f"ran {attempted} task(s)") + return 0 + signal.signal(signal.SIGINT, _sigint) + signal.signal(signal.SIGTERM, _sigint) + print("nanodot runner started — Ctrl-C to stop", flush=True) + daemon.serve(stop) + print("nanodot runner stopped") return 0 - signal.signal(signal.SIGINT, _sigint) - signal.signal(signal.SIGTERM, _sigint) - print("nanodot runner started — Ctrl-C to stop", flush=True) - daemon.serve(stop) - print("nanodot runner stopped") - return 0 + finally: + # Unwind while still owning the lifetime lock, in reverse + # registration order: the next runner never meets a + # half-closed store, and a failing step never blocks the + # rest of the shutdown. + for name in teardown.run(): + print(f"warning: {name} did not shut down cleanly", + file=sys.stderr) except (RunnerControlError, OSError, ValueError) as error: print(f"error: {error}", file=sys.stderr) return 1 diff --git a/src/nanodot/core/runner.py b/src/nanodot/core/runner.py index b8e1bad..df68ce8 100644 --- a/src/nanodot/core/runner.py +++ b/src/nanodot/core/runner.py @@ -97,6 +97,15 @@ def work() -> None: return None return result[0] if result else None + def drain(self, seconds: float) -> None: + """Wait, bounded by ``seconds``, for an in-flight summary to finish. + + A provider that will not finish costs the wait, not the shutdown: + the drain gives up and the summary is abandoned mid-flight. + """ + if self._inflight.acquire(timeout=max(0.0, seconds)): + self._inflight.release() + class RunOutcome(str, Enum): OK = "ok" @@ -131,6 +140,7 @@ def __init__( self._activity = activity self._provider = provider self._memory = memory + self._summary_budget_seconds = summary_budget_seconds self._summaries = _SummaryBudget(summary_budget_seconds) self._locks: dict[str, threading.Lock] = {} self._locks_guard = threading.Lock() @@ -139,6 +149,17 @@ def _lock_for(self, task_id: str) -> threading.Lock: with self._locks_guard: return self._locks.setdefault(task_id, threading.Lock()) + def close(self) -> None: + """Release the loop's own resources: drain an in-flight summary + (bounded by the summary budget) and drop the per-task locks. + + Call after serve/tick has returned. Safe to call more than once; + a later run_once rebuilds whatever it needs. + """ + self._summaries.drain(self._summary_budget_seconds) + with self._locks_guard: + self._locks.clear() + def run_once(self, task: Task, now: float) -> RunOutcome: """One bounded check for one task. Safe to call concurrently: a second run of the same task is skipped, never overlapped.""" diff --git a/src/nanodot/core/teardown.py b/src/nanodot/core/teardown.py new file mode 100644 index 0000000..6927bab --- /dev/null +++ b/src/nanodot/core/teardown.py @@ -0,0 +1,37 @@ +"""Registrations are effects — the shutdown discipline. + +Resources register a named disposer in creation order; shutdown unwinds +in reverse. A raising disposer never blocks the unwind: every step runs, +and failures are reported by name only — never exception text, which can +carry private data. See docs/design/adapter-seam.md, admission +discipline 2. +""" + +from __future__ import annotations + +from collections.abc import Callable + + +class Teardown: + """Collect disposers as resources are built; unwind them on stop.""" + + def __init__(self) -> None: + self._disposers: list[tuple[str, Callable[[], None]]] = [] + + def register(self, name: str, disposer: Callable[[], None]) -> None: + self._disposers.append((name, disposer)) + + def run(self) -> list[str]: + """Unwind in reverse registration order, exactly once. + + Returns the names of disposers that raised; every disposer runs + regardless. A second call unwinds nothing. + """ + failed: list[str] = [] + while self._disposers: + name, disposer = self._disposers.pop() + try: + disposer() + except Exception: + failed.append(name) + return failed diff --git a/tests/test_public_mode.py b/tests/test_public_mode.py index ddfc06e..022743e 100644 --- a/tests/test_public_mode.py +++ b/tests/test_public_mode.py @@ -305,17 +305,20 @@ def test_runner_policy_changes_require_stopped_runner( def test_runner_loads_policy_under_lock_before_ready(home: Path) -> None: from nanodot.cli import _wiring + from nanodot.core.teardown import Teardown - def checked_wiring(): + def checked_wiring(teardown=None): assert not (home / "runner.pid").exists() with pytest.raises(RunnerAlreadyRunning): with configuration_lock(home / "runner.pid"): pytest.fail("runner read configuration before acquiring ownership") - return _wiring() + return _wiring(teardown) with mock.patch("nanodot.cli._wiring", side_effect=checked_wiring) as wiring: assert main(["runner", "--once"]) == 0 - wiring.assert_called_once_with() + wiring.assert_called_once() + (passed,), kwargs = wiring.call_args + assert isinstance(passed, Teardown) and not kwargs def test_background_start_does_not_report_ready_with_invalid_policy(home: Path, capsys) -> None: diff --git a/tests/test_teardown.py b/tests/test_teardown.py new file mode 100644 index 0000000..39572ed --- /dev/null +++ b/tests/test_teardown.py @@ -0,0 +1,156 @@ +"""Registrations-are-effects shutdown: reverse unwind, residue-free stop +(docs/design/adapter-seam.md, admission discipline 2 / issue #37).""" + +from __future__ import annotations + +import threading +import time +from pathlib import Path + +import pytest +from fakes import FAILURE, FakeGitHub, FakeSink + +from nanodot.cli import main +from nanodot.core.activity import ActivityLog +from nanodot.core.memory import MemoryStore +from nanodot.core.runner import RunOutcome, TaskLoop +from nanodot.core.tasks import PRTarget, Task, TaskStore +from nanodot.core.teardown import Teardown +from nanodot.ports.inference import TaskDraft + +TARGET = PRTarget.parse("thinkflowlab/nanodot#12") + + +# -- Teardown: the registry -------------------------------------------------- + + +def test_teardown_unwinds_in_reverse_order_and_runs_once() -> None: + order: list[str] = [] + teardown = Teardown() + teardown.register("first", lambda: order.append("first")) + teardown.register("second", lambda: order.append("second")) + teardown.register("third", lambda: order.append("third")) + + assert teardown.run() == [] + assert order == ["third", "second", "first"] + assert teardown.run() == [] # a second stop unwinds nothing + assert order == ["third", "second", "first"] + + +def test_teardown_runs_every_step_and_names_only_the_failures() -> None: + order: list[str] = [] + teardown = Teardown() + + def failing() -> None: + order.append("second") + raise RuntimeError("private detail that must not be reported verbatim") + + teardown.register("first", lambda: order.append("first")) + teardown.register("second", failing) + teardown.register("third", lambda: order.append("third")) + + assert teardown.run() == ["second"] + assert order == ["third", "second", "first"] # every step ran, in reverse + + +# -- TaskLoop.close: drain the summary, drop the locks ----------------------- + + +class _GatedProvider: + """Summarize blocks until released — a misbehaving, slow provider.""" + + def __init__(self) -> None: + self.entered = threading.Event() + self.release = threading.Event() + + def summarize(self, change) -> str: + self.entered.set() + self.release.wait(5.0) + time.sleep(0.1) # linger after release so a drain is observable + return "model summary" + + def parse_intent(self, text: str) -> TaskDraft: + return TaskDraft() + + +def _loop(home: Path, budget: float) -> tuple[TaskStore, TaskLoop, Task, _GatedProvider]: + store = TaskStore(path=home / "nanodot.db") + activity = ActivityLog(path=home / "nanodot.db") + github = FakeGitHub(TARGET) + github.set_pr("open", head_sha="s1") + github.add_check("ci", FAILURE, sha="s1") # notable → a summary is requested + provider = _GatedProvider() + loop = TaskLoop( + store, github, FakeSink(), activity, + provider=provider, summary_budget_seconds=budget, + ) + task = store.create( + Task(target=TARGET, purpose="watch", cadence_seconds=300, next_check_at=0.0) + ) + return store, loop, task, provider + + +def test_task_loop_close_is_bounded_even_with_a_hung_provider(home: Path) -> None: + store, loop, task, provider = _loop(home, budget=0.2) + outcomes: list[RunOutcome] = [] + worker = threading.Thread( + target=lambda: outcomes.append(loop.run_once(store.get(task.id), 0.0)) + ) + worker.start() + assert provider.entered.wait(2.0), "summary never reached the provider" + + started = time.monotonic() + loop.close() + assert time.monotonic() - started < 2.0, "close hung on the summary drain" + + provider.release.set() + worker.join(5.0) + assert outcomes == [RunOutcome.OK] # the run itself was never interrupted + + loop.close() # idempotent with nothing in flight + + +def test_task_loop_close_drains_a_finishing_summary(home: Path) -> None: + store, loop, task, provider = _loop(home, budget=5.0) + outcomes: list[RunOutcome] = [] + worker = threading.Thread( + target=lambda: outcomes.append(loop.run_once(store.get(task.id), 0.0)) + ) + worker.start() + assert provider.entered.wait(2.0) + provider.release.set() + + started = time.monotonic() + loop.close() + assert time.monotonic() - started >= 0.08, "close did not wait for the provider" + worker.join(5.0) + assert outcomes == [RunOutcome.OK] + + +# -- The real wiring: a full pass unwinds under the lease -------------------- + + +def test_runner_once_unwinds_every_registration_in_reverse( + home: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + order: list[str] = [] + monkeypatch.setattr(TaskStore, "close", lambda self: order.append("task-store")) + monkeypatch.setattr(ActivityLog, "close", lambda self: order.append("activity-log")) + monkeypatch.setattr(MemoryStore, "close", lambda self: order.append("memory-store")) + monkeypatch.setattr(TaskLoop, "close", lambda self: order.append("task-loop")) + + assert main(["runner", "--once"]) == 0 + assert order == ["task-loop", "memory-store", "activity-log", "task-store"] + + +def test_runner_once_disposes_cleanly_and_leaves_readable_data( + home: Path, capsys: pytest.CaptureFixture[str] +) -> None: + """Unpatched: every real disposer ran cleanly, and the next owner of + the data home reads the same files without trouble.""" + assert main(["runner", "--once"]) == 0 + assert "did not shut down cleanly" not in capsys.readouterr().err + + store = TaskStore(path=home / "nanodot.db") + assert store.list() == [] + assert ActivityLog(path=home / "nanodot.db").query() == [] From d21e8cfcc45007976661774a2f7c14b3bb8cabe3 Mon Sep 17 00:00:00 2001 From: Hongsheng Liu Date: Fri, 2 Oct 2026 19:22:35 +0800 Subject: [PATCH 2/2] Document the teardown registry's thread contract --- src/nanodot/core/teardown.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/src/nanodot/core/teardown.py b/src/nanodot/core/teardown.py index 6927bab..77721d6 100644 --- a/src/nanodot/core/teardown.py +++ b/src/nanodot/core/teardown.py @@ -13,7 +13,12 @@ class Teardown: - """Collect disposers as resources are built; unwind them on stop.""" + """Collect disposers as resources are built; unwind them on stop. + + Not thread-safe: register from the orchestrating thread that will + later call run — the same thread in the native runner, and the + contract an adapter host must keep. + """ def __init__(self) -> None: self._disposers: list[tuple[str, Callable[[], None]]] = []