diff --git a/cron/executions.py b/cron/executions.py index 14c1d40e2e..755e398415 100644 --- a/cron/executions.py +++ b/cron/executions.py @@ -2,7 +2,7 @@ The ledger records what is known about each attempt; it is not a retry queue. Interrupted attempts become ``unknown`` only after their exact owner process is proved gone. Terminal states are -immutable. Also hosts the SQLite ledger helpers shared with ``cron.incidents`` / ``cron.notepad``. +immutable. """ from __future__ import annotations @@ -14,8 +14,9 @@ import time import uuid from contextlib import contextmanager from pathlib import Path -from typing import Any, Callable, Dict, Iterator, List, Optional +from typing import Any, Dict, Iterator, List, Optional +from cron.ledger import ledger_transaction, open_ledger, prepare_ledger from hermes_constants import get_hermes_home from hermes_time import now as _hermes_now @@ -30,48 +31,6 @@ _lock = threading.RLock() _PROCESS_ID = uuid.uuid4().hex -# --- shared SQLite ledger plumbing -------------------------------------------------------------- - -def open_ledger(path: Path) -> sqlite3.Connection: - """Open a profile-local ledger DB, creating its ``cron/`` dir with the store's permissions.""" - from cron.jobs import _ensure_cron_dir - - _ensure_cron_dir(path.parent) - return sqlite3.connect(path, timeout=5) - - -def prepare_ledger( - conn: sqlite3.Connection, *, db_label: str, synchronous_full: bool = True -) -> None: - """Row factory + busy timeout + WAL (with fallback) + optional ``synchronous=FULL``.""" - from hermes_state_wal import apply_wal_with_fallback - - conn.row_factory = sqlite3.Row - conn.execute("PRAGMA busy_timeout=5000") - apply_wal_with_fallback(conn, db_label=db_label) - if synchronous_full: - conn.execute("PRAGMA synchronous=FULL") - - -@contextmanager -def ledger_transaction( - lock: threading.RLock, - connect: Callable[[], sqlite3.Connection], - initialize_schema: Callable[[sqlite3.Connection], None], -) -> Iterator[sqlite3.Connection]: - """Open a connection, commit/rollback on exit, always close. ``sqlite3.Connection``'s own - context manager does NOT close (leaks WAL/SHM fds until GC); schema init runs inside the - ``try`` so a PRAGMA/DDL failure after ``connect()`` still closes.""" - with lock: - conn = connect() - try: - initialize_schema(conn) - with conn: - yield conn - finally: - conn.close() - - # --- executions ledger -------------------------------------------------------------------------- def _connect() -> sqlite3.Connection: diff --git a/cron/incidents.py b/cron/incidents.py index 4e8a517f7c..02a9ae4396 100644 --- a/cron/incidents.py +++ b/cron/incidents.py @@ -19,7 +19,7 @@ from pathlib import Path from typing import Any, Dict, Iterator, List, Optional from cron import executions as _executions -from cron.executions import ledger_transaction, open_ledger, prepare_ledger +from cron.ledger import ledger_transaction, open_ledger, prepare_ledger from hermes_constants import get_hermes_home from hermes_time import now as _hermes_now diff --git a/cron/ledger.py b/cron/ledger.py new file mode 100644 index 0000000000..984da83cfc --- /dev/null +++ b/cron/ledger.py @@ -0,0 +1,47 @@ +"""SQLite connection and transaction helpers shared by cron ledgers.""" + +from __future__ import annotations + +import sqlite3 +import threading +from contextlib import contextmanager +from pathlib import Path +from typing import Callable, Iterator + + +def open_ledger(path: Path) -> sqlite3.Connection: + """Open a profile-local ledger DB, creating its cron directory securely.""" + from cron.jobs import _ensure_cron_dir + + _ensure_cron_dir(path.parent) + return sqlite3.connect(path, timeout=5) + + +def prepare_ledger( + conn: sqlite3.Connection, *, db_label: str, synchronous_full: bool = True +) -> None: + """Configure row access, busy timeout, WAL, and optional full synchronization.""" + from hermes_state_wal import apply_wal_with_fallback + + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA busy_timeout=5000") + apply_wal_with_fallback(conn, db_label=db_label) + if synchronous_full: + conn.execute("PRAGMA synchronous=FULL") + + +@contextmanager +def ledger_transaction( + lock: threading.RLock, + connect: Callable[[], sqlite3.Connection], + initialize_schema: Callable[[sqlite3.Connection], None], +) -> Iterator[sqlite3.Connection]: + """Initialize, transact on, and always close one ledger connection.""" + with lock: + conn = connect() + try: + initialize_schema(conn) + with conn: + yield conn + finally: + conn.close() diff --git a/cron/notepad.py b/cron/notepad.py index a84c37b363..ee740203b2 100644 --- a/cron/notepad.py +++ b/cron/notepad.py @@ -15,7 +15,7 @@ from contextlib import contextmanager from pathlib import Path from typing import Any, Dict, Iterator, List, Optional -from cron.executions import ledger_transaction, open_ledger, prepare_ledger +from cron.ledger import ledger_transaction, open_ledger, prepare_ledger from hermes_constants import get_hermes_home from hermes_time import now as _hermes_now diff --git a/tests/cron/test_upgrade_module_skew.py b/tests/cron/test_upgrade_module_skew.py new file mode 100644 index 0000000000..6e3828cc0e --- /dev/null +++ b/tests/cron/test_upgrade_module_skew.py @@ -0,0 +1,33 @@ +"""Cron imports remain usable when a daemon spans an on-disk upgrade.""" + +from __future__ import annotations + +import subprocess +import sys +from pathlib import Path + + +def test_lazy_cron_stores_do_not_require_new_symbols_on_cached_executions_module(): + repo_root = Path(__file__).resolve().parents[2] + script = """ +import sys +import cron.executions as executions + +for name in ("ledger_transaction", "open_ledger", "prepare_ledger"): + delattr(executions, name) +sys.modules.pop("cron.incidents", None) +sys.modules.pop("cron.notepad", None) + +import cron.incidents +import cron.notepad +""" + + result = subprocess.run( + [sys.executable, "-c", script], + cwd=repo_root, + capture_output=True, + text=True, + check=False, + ) + + assert result.returncode == 0, result.stderr