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
54 changes: 36 additions & 18 deletions src/nanodot/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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


Expand Down Expand Up @@ -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()
Expand All @@ -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
Expand Down
21 changes: 21 additions & 0 deletions src/nanodot/core/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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()
Expand All @@ -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."""
Expand Down
42 changes: 42 additions & 0 deletions src/nanodot/core/teardown.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
"""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.

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]]] = []

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
9 changes: 6 additions & 3 deletions tests/test_public_mode.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
156 changes: 156 additions & 0 deletions tests/test_teardown.py
Original file line number Diff line number Diff line change
@@ -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() == []
Loading