fix(cron): isolate ledger helpers from stale execution modules
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
47
cron/ledger.py
Normal file
47
cron/ledger.py
Normal file
@@ -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()
|
||||
@@ -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
|
||||
|
||||
|
||||
33
tests/cron/test_upgrade_module_skew.py
Normal file
33
tests/cron/test_upgrade_module_skew.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user