From d301af3c59b4d36c874212324c5e68b027c5c00d Mon Sep 17 00:00:00 2001 From: poisdahl <4091911+poisdahl@users.noreply.github.com> Date: Sat, 8 Aug 2026 10:59:52 +0200 Subject: [PATCH 1/5] fix(cron): make degraded jobs saves conflict-safe --- cron/jobs.py | 576 ++++++++++++++------- hermes_cli/backup.py | 4 + tests/cron/test_jobs_conflict_safe_save.py | 353 +++++++++++++ tests/cron/test_jobs_shrink_merge_80624.py | 2 +- tests/hermes_cli/test_backup.py | 15 +- 5 files changed, 753 insertions(+), 197 deletions(-) create mode 100644 tests/cron/test_jobs_conflict_safe_save.py diff --git a/cron/jobs.py b/cron/jobs.py index 377b34c35487a..a6c81aeff7cd3 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -9,9 +9,11 @@ import copy from contextvars import ContextVar from dataclasses import dataclass +import errno import json import logging import shutil +import stat import tempfile import threading import time @@ -112,6 +114,13 @@ def _ensure_croniter() -> bool: # legitimate critical section (field updates only) while keeping the ticker's # worst-case stall well under one status-alarm threshold. _JOBS_LOCK_TIMEOUT_SECONDS = 30.0 +_JOBS_COMMIT_LOCK_TIMEOUT_SECONDS = 5.0 +_JOBS_GENERATION_MAX_ATTEMPTS = 3 +_UNSUPPORTED_LOCK_ERRNOS = frozenset( + value + for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP") + if (value := getattr(errno, name, None)) is not None +) OUTPUT_DIR = CRON_DIR / "output" ONESHOT_GRACE_SECONDS = 120 @@ -123,6 +132,14 @@ class _CronStorePaths: output_dir: Path +@dataclass +class _LoadedJobsSnapshot: + owner: List[Dict[str, Any]] + base: List[Dict[str, Any]] + stamp: Optional[Tuple[int, int, int]] + stamp_trusted: bool = True + + _cron_store_override: ContextVar[Optional[_CronStorePaths]] = ContextVar( "cron_store_override", default=None, @@ -267,6 +284,126 @@ def _jobs_lock_file() -> Path: return _current_cron_store().cron_dir / ".jobs.lock" +def _jobs_commit_lock_file() -> Path: + return _current_cron_store().cron_dir / ".jobs.commit.lock" + + +def _prepare_commit_lock_file(handle, owner_source: Optional[os.stat_result]) -> None: + """Keep the shared lock writable by the cron-directory owner.""" + try: + os.fchmod(handle.fileno(), 0o600) + except (AttributeError, OSError, NotImplementedError): + pass + + try: + geteuid = getattr(os, "geteuid", None) + if ( + owner_source is not None + and os.name == "posix" + and geteuid is not None + and geteuid() == 0 + ): + os.fchown(handle.fileno(), owner_source.st_uid, owner_source.st_gid) + except (AttributeError, OSError) as exc: + logger.warning( + "Could not restore ownership of %s to the cron directory owner: %s", + _jobs_commit_lock_file(), + exc, + ) + + +@contextlib.contextmanager +def _jobs_commit_lock(): + """Serialize the final read, reconcile, and publication of jobs.json.""" + if fcntl is None and msvcrt is None: + yield + return + + ensure_dirs() + lock_path = _jobs_commit_lock_file() + try: + owner_source = os.stat(lock_path.parent) + except OSError: + owner_source = None + flags = os.O_RDWR | os.O_CREAT | getattr(os, "O_BINARY", 0) + if os.name == "posix": + flags |= getattr(os, "O_NOFOLLOW", 0) + fd = None + try: + fd = os.open(lock_path, flags, 0o600) + if not stat.S_ISREG(os.fstat(fd).st_mode): + raise RuntimeError("Cron jobs commit lock is not a regular file") + handle = os.fdopen(fd, "r+b") + fd = None + except OSError as exc: + raise RuntimeError("Unable to open the cron jobs commit lock") from exc + finally: + if fd is not None: + os.close(fd) + try: + _prepare_commit_lock_file(handle, owner_source) + except BaseException: + handle.close() + raise + acquired = False + try: + if fcntl is None and msvcrt is not None: + handle.seek(0, os.SEEK_END) + if handle.tell() == 0: + handle.write(b" ") + handle.flush() + handle.seek(0) + + deadline = time.monotonic() + _JOBS_COMMIT_LOCK_TIMEOUT_SECONDS + backend_error = None + while True: + try: + if fcntl is not None: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + else: + handle.seek(0) + getattr(msvcrt, "locking")( + handle.fileno(), getattr(msvcrt, "LK_NBLCK"), 1 + ) + acquired = True + break + except (BlockingIOError, OSError) as exc: + if getattr(exc, "errno", None) in _UNSUPPORTED_LOCK_ERRNOS: + backend_error = exc + break + if time.monotonic() >= deadline: + raise RuntimeError( + "Timed out waiting for the cron jobs commit lock; " + "refusing to publish a potentially stale snapshot" + ) + time.sleep(0.01) + + if backend_error is not None: + logger.warning( + "Cron jobs commit lock is unavailable (%s); proceeding without " + "cross-process publication locking", + backend_error, + ) + yield + return + yield + finally: + try: + if acquired: + try: + if fcntl is not None: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + else: + handle.seek(0) + getattr(msvcrt, "locking")( + handle.fileno(), getattr(msvcrt, "LK_UNLCK"), 1 + ) + except OSError: + pass + finally: + handle.close() + + @contextlib.contextmanager def _jobs_lock(): """Serialize a load_jobs→modify→save_jobs critical section. @@ -298,12 +435,10 @@ def _jobs_lock(): with _jobs_file_lock: _jobs_lock_state.depth = 1 - # Stamp of jobs.json as of this section's load_jobs() (#80703's - # fast-path, credit @JoaoMarcos44): lets _save_jobs_unlocked skip the - # shrink-merge parse when the file provably hasn't changed since this - # section read it. Reset on entry/exit so stale stamps from unlocked - # loads or prior sections can never suppress a needed merge. - _jobs_lock_state.load_stamp = None + # Keep a bounded base snapshot for each list returned by load_jobs() + # during this outer mutation section. Object identity keeps nested + # load/modify/save callers from overwriting one another's provenance. + _jobs_lock_state.load_snapshots = [] lock_fd = None try: try: @@ -369,7 +504,7 @@ def _jobs_lock(): lock_fd.close() finally: _jobs_lock_state.depth = 0 - _jobs_lock_state.load_stamp = None + _jobs_lock_state.load_snapshots = [] # Fields on a cron job that must never change after creation. ``id`` is used # as a filesystem path component under ``OUTPUT_DIR``; allowing it to be @@ -1079,13 +1214,13 @@ def load_jobs() -> List[Dict[str, Any]]: """Load all jobs from storage.""" jobs_file = _current_cron_store().jobs_file ensure_dirs() - # Stamp BEFORE reading (fail-safe direction — see _record_load_stamp): - # a sibling write racing this load leaves the stamp older than disk, so - # the save-path merge runs instead of being wrongly skipped. + # Stamp BEFORE reading (fail-safe direction): a sibling write racing this + # load leaves the stamp older than disk, so the save path reconciles it. pre_read_stamp = _jobs_file_stamp(jobs_file) if not jobs_file.exists(): - _record_load_stamp(None) - return [] + jobs = [] + _record_loaded_jobs(jobs, None) + return jobs try: data, _strict_retry = _parse_jobs_file(jobs_file) @@ -1102,19 +1237,19 @@ def load_jobs() -> List[Dict[str, Any]]: # down the whole cron subsystem. if isinstance(data, dict): jobs = data.get("jobs", []) + _record_loaded_jobs(jobs, pre_read_stamp) if _strict_retry and jobs: # Hit control-character corruption — rewrite with proper escaping. save_jobs(jobs) logger.warning("Auto-repaired jobs.json (had invalid control characters)") - _record_load_stamp(pre_read_stamp) return jobs if isinstance(data, list): + _record_loaded_jobs(data, pre_read_stamp) # Bare array — likely saved/edited outside save_jobs(). Wrap it back # into the expected {"jobs": [...]} structure. if data: save_jobs(data) logger.warning("Auto-repaired jobs.json (bare list wrapped as dict)") - _record_load_stamp(pre_read_stamp) return data raise RuntimeError( @@ -1127,7 +1262,7 @@ def _peek_jobs_unlocked() -> Optional[List[Dict[str, Any]]]: Caller must hold ``_jobs_lock()``. Returns ``[]`` when the file is missing, ``None`` when the payload is unreadable/corrupt (caller should - not attempt a shrink-merge against an unknown baseline). Never calls + not attempt reconciliation against an unknown baseline). Never calls ``save_jobs`` — the repair-free property is what keeps the save path re-entrancy-safe (a repairing read here would recurse through ``_save_jobs_unlocked``). @@ -1151,13 +1286,12 @@ def _jobs_file_stamp(jobs_file: Path) -> Optional[Tuple[int, int, int]]: """Cheap change-detection stamp for jobs.json: ``(mtime_ns, size, ino)``. ``None`` means the file is missing/unstatable. Used as a fast-path gate - in front of the shrink-merge so the healthy no-race save costs one - ``stat()`` instead of a full read+parse (the ``advance_next_runs`` - batching exists because this path is hot — see its docstring). - ``st_ino`` is included because every legitimate writer goes through - mkstemp+rename (new inode), so even a same-size write inside one mtime - quantum on a coarse-clock filesystem (ext4 jiffies, network mounts) - cannot false-match. + before reconciliation so the healthy no-race save uses stat checks instead + of a full read+parse (the ``advance_next_runs`` batching exists because + this path is hot — see its docstring). ``st_ino`` is included because + every legitimate writer goes through mkstemp+rename (new inode), so even a + same-size write inside one mtime quantum on a coarse-clock filesystem + (ext4 jiffies, network mounts) cannot false-match. """ try: st = jobs_file.stat() @@ -1166,84 +1300,170 @@ def _jobs_file_stamp(jobs_file: Path) -> Optional[Tuple[int, int, int]]: return None -def _record_load_stamp(stamp: Optional[Tuple[int, int, int]]) -> None: - """Remember jobs.json's stamp for the enclosing _jobs_lock() section. - - No-op outside a critical section. Lets the save path skip the - shrink-merge parse when the file provably hasn't changed since this - section loaded it (#80703's fast-path). The caller must capture the - stamp BEFORE reading the file: a sibling landing mid-read then leaves - the recorded stamp OLDER than disk — a mismatch, so the merge runs - (fail-safe direction). Stamping after the read would let that sibling's - write be certified as "seen" without being in the loaded payload, - wrongly suppressing the recovery. - """ +def _record_loaded_jobs( + owner: List[Dict[str, Any]], + stamp: Optional[Tuple[int, int, int]], +) -> None: if not getattr(_jobs_lock_state, "depth", 0): return - _jobs_lock_state.load_stamp = stamp + _jobs_lock_state.load_snapshots.append( + _LoadedJobsSnapshot(owner, copy.deepcopy(owner), stamp) + ) -def _merge_unexpected_disk_jobs( - jobs: List[Dict[str, Any]], - *, - removed_ids: Optional[Collection[str]] = None, -) -> List[Dict[str, Any]]: - """Return *jobs* plus any on-disk jobs missing from the save payload (#80624). - - Under ``_jobs_lock()``'s degraded flock-timeout path (#60703), two - processes can both believe they own the store. A writer that loaded an - older/smaller snapshot then calls ``save_jobs`` and would otherwise - clobber concurrent creates (the filed ``no_agent`` watchdog pattern: - CLI/tool create succeeds, then a gateway tick/remove rewrites - ``jobs.json`` empty or without the new id). - - Intentional deletes pass ``removed_ids``. Any other id present on disk - but absent from *jobs* is treated as a concurrent create and merged - back before the atomic write. The caller's list is never mutated — a - new list is returned when anything was recovered. - - Fast path: when the enclosing critical section recorded a load stamp - and the file's ``(mtime_ns, size)`` still matches, nothing can have - changed underneath us, so the read+parse is skipped entirely — one - ``stat()`` on the healthy no-race save. - """ - stamp = getattr(_jobs_lock_state, "load_stamp", None) - if stamp is not None and _jobs_file_stamp(_current_cron_store().jobs_file) == stamp: - return jobs +def _loaded_jobs_snapshot( + owner: List[Dict[str, Any]], +) -> Optional[_LoadedJobsSnapshot]: + for snapshot in reversed(getattr(_jobs_lock_state, "load_snapshots", ())): + if snapshot.owner is owner: + return snapshot + return None + + +_MISSING_JOB = object() - disk_jobs = _peek_jobs_unlocked() - if disk_jobs is None: - return jobs - intended_remove = {str(i) for i in (removed_ids or ()) if i} - new_ids: Set[str] = set() +def _valid_job_id(job: Any) -> Optional[str]: + if not isinstance(job, dict) or not job.get("id"): + return None + return str(job["id"]) + + +def _index_unique_jobs( + jobs: List[Dict[str, Any]], +) -> Optional[Tuple[List[str], Dict[str, Dict[str, Any]]]]: + order = [] + by_id = {} for job in jobs: - if isinstance(job, dict) and job.get("id"): - new_ids.add(str(job["id"])) + job_id = _valid_job_id(job) + if job_id is None or job_id in by_id: + return None + order.append(job_id) + by_id[job_id] = job + return order, by_id - recovered: List[Dict[str, Any]] = [] - for disk_job in disk_jobs: - if not isinstance(disk_job, dict): - continue - disk_id = disk_job.get("id") - if not disk_id: - continue - disk_id = str(disk_id) - if disk_id in new_ids or disk_id in intended_remove: - continue - recovered.append(disk_job) - new_ids.add(disk_id) +def _warn_recovered_job_ids(job_ids: Collection[str]) -> None: + recovered = list(job_ids) if not recovered: - return jobs + return logger.warning( "Preserved %d cron job(s) present on disk but missing from the " "in-memory save payload (concurrent create under degraded lock " "or stale writer) (#80624): %s", len(recovered), - [j.get("id") for j in recovered], + recovered, ) - return jobs + recovered + + +def _merge_jobs_three_way( + base: List[Dict[str, Any]], + desired: List[Dict[str, Any]], + current: List[Dict[str, Any]], + *, + removed_ids: Optional[Collection[str]] = None, +) -> List[Dict[str, Any]]: + """Merge desired and current snapshots against one loaded base.""" + intended_remove = {str(job_id) for job_id in (removed_ids or ()) if job_id} + base_ids = { + job_id for job in base if (job_id := _valid_job_id(job)) is not None + } + desired_ids = { + job_id for job in desired if (job_id := _valid_job_id(job)) is not None + } + external_removals = intended_remove - base_ids - desired_ids + if external_removals: + current = [ + job for job in current if _valid_job_id(job) not in external_removals + ] + recovered_ids = list( + dict.fromkeys( + job_id + for job in current + if (job_id := _valid_job_id(job)) is not None + and job_id not in base_ids + and job_id not in desired_ids + ) + ) + + if desired == base: + _warn_recovered_job_ids(recovered_ids) + return copy.deepcopy(current) + if current == base: + return copy.deepcopy(desired) + if desired == current: + return copy.deepcopy(desired) + + base_index = _index_unique_jobs(base) + desired_index = _index_unique_jobs(desired) + current_index = _index_unique_jobs(current) + if base_index is None or desired_index is None or current_index is None: + raise RuntimeError( + "Refusing to merge concurrent cron changes with missing or duplicate ids" + ) + base_order, base_by_id = base_index + desired_order, desired_by_id = desired_index + current_order, current_by_id = current_index + + chosen = {} + for job_id in dict.fromkeys(base_order + desired_order + current_order): + base_job = base_by_id.get(job_id, _MISSING_JOB) + desired_job = desired_by_id.get(job_id, _MISSING_JOB) + current_job = current_by_id.get(job_id, _MISSING_JOB) + if desired_job == base_job: + merged = current_job + elif current_job == base_job: + merged = desired_job + elif desired_job == current_job: + merged = desired_job + else: + raise RuntimeError( + "Refusing to save conflicting concurrent cron job change for " + f"id {job_id!r}" + ) + chosen[job_id] = merged + + result = [] + emitted = set() + for job in base: + job_id = _valid_job_id(job) + assert job_id is not None + emitted.add(job_id) + merged = chosen[job_id] + if merged is not _MISSING_JOB: + result.append(copy.deepcopy(merged)) + for job_id in desired_order + current_order: + if job_id in emitted: + continue + emitted.add(job_id) + merged = chosen[job_id] + if merged is not _MISSING_JOB: + result.append(copy.deepcopy(merged)) + _warn_recovered_job_ids(recovered_ids) + return result + + +def _merge_jobs_without_base( + desired: List[Dict[str, Any]], + current: List[Dict[str, Any]], + removed_ids: Optional[Collection[str]], +) -> List[Dict[str, Any]]: + """Preserve on-disk jobs for callers without a loaded base.""" + if desired == current: + return copy.deepcopy(desired) + intended_remove = {str(job_id) for job_id in (removed_ids or ()) if job_id} + desired_ids = { + job_id for job in desired if (job_id := _valid_job_id(job)) is not None + } + result = copy.deepcopy(desired) + recovered = [] + for job in current: + job_id = _valid_job_id(job) + if job_id is not None and job_id not in desired_ids and job_id not in intended_remove: + recovered.append(copy.deepcopy(job)) + desired_ids.add(job_id) + _warn_recovered_job_ids([job["id"] for job in recovered]) + return result + recovered def _save_jobs_unlocked( @@ -1255,127 +1475,95 @@ def _save_jobs_unlocked( """Save all jobs to storage. Caller must hold _jobs_lock(). ``removed_ids`` lists job ids this mutation intentionally deleted. - ``replace=True`` skips the shrink-merge guard (tests / disaster recovery + ``replace=True`` skips the reconciliation guard (tests / disaster recovery that mean to rewrite the store wholesale). """ jobs_file = _current_cron_store().jobs_file ensure_dirs() - # Snapshot the current owner BEFORE the atomic replace so a privileged - # writer (root CLI in Docker) can hand ownership back to the gateway user - # afterwards instead of locking its ticker out (#68483). When the file is - # being created for the first time, inherit the cron dir's owner — in the - # Docker image that is the PUID/PGID gateway user who must be able to - # read the store on the next tick. - try: - _stat_before = os.stat(jobs_file) - except OSError: - try: - _stat_before = os.stat(jobs_file.parent) - except OSError: - _stat_before = None - - # Shrink-merge + rewrite loop (#80624): under the degraded flock-timeout - # path another process can create a job between our load and our write. - # Merge unexpected disk ids into the payload, stage the write, then - # re-peek; if new ids appeared, merge again and restage before replace. - # The merge itself fast-paths to a single stat() when the enclosing - # section's load stamp still matches (see _merge_unexpected_disk_jobs). - tmp_path = None - try: - for _attempt in range(5): - if not replace: - jobs = _merge_unexpected_disk_jobs(jobs, removed_ids=removed_ids) + desired = copy.deepcopy(jobs) + loaded = _loaded_jobs_snapshot(jobs) + base = copy.deepcopy(loaded.base) if loaded is not None else None + base_stamp = loaded.stamp if loaded is not None else None + + with _jobs_commit_lock(): + for _attempt in range(_JOBS_GENERATION_MAX_ATTEMPTS): + if replace: + observed_stamp = None + reconciled = copy.deepcopy(desired) + elif ( + base is not None + and loaded is not None + and loaded.stamp_trusted + and base_stamp is not None + and _jobs_file_stamp(jobs_file) == base_stamp + ): + observed_stamp = base_stamp + reconciled = copy.deepcopy(desired) + else: + observed_stamp = _jobs_file_stamp(jobs_file) + current = _peek_jobs_unlocked() + if current is None: + if base is not None: + raise RuntimeError( + "Cron database changed to an unreadable generation" + ) + reconciled = copy.deepcopy(desired) + elif base is not None: + reconciled = _merge_jobs_three_way( + base, desired, current, removed_ids=removed_ids + ) + else: + reconciled = _merge_jobs_without_base( + desired, current, removed_ids + ) + + # Snapshot the current owner before atomic replace so a privileged + # writer can hand ownership back to the gateway user afterwards. + try: + stat_before = os.stat(jobs_file) + except OSError: + try: + stat_before = os.stat(jobs_file.parent) + except OSError: + stat_before = None + fd, tmp_path = tempfile.mkstemp( dir=str(jobs_file.parent), suffix=".tmp", prefix=".jobs_" ) + published = False try: - with os.fdopen(fd, "w", encoding="utf-8") as f: + with os.fdopen(fd, "w", encoding="utf-8") as handle: json.dump( - {"jobs": jobs, "updated_at": _hermes_now().isoformat()}, - f, + {"jobs": reconciled, "updated_at": _hermes_now().isoformat()}, + handle, indent=2, ) - f.flush() - os.fsync(f.fileno()) - except BaseException: - try: - os.unlink(tmp_path) - except OSError: - pass - tmp_path = None - raise - - if not replace: - # Verify-after-stage: a sibling landing while we serialized - # the payload must trigger another merge round. Same stamp - # fast path as the merge — an unchanged stamp proves nothing - # was written, so the full parse is skipped. - _stamp = getattr(_jobs_lock_state, "load_stamp", None) - _unchanged = ( - _stamp is not None and _jobs_file_stamp(jobs_file) == _stamp - ) - disk_jobs = None if _unchanged else _peek_jobs_unlocked() - if disk_jobs is not None: - payload_ids = { - str(j["id"]) - for j in jobs - if isinstance(j, dict) and j.get("id") - } - intended = {str(i) for i in (removed_ids or ()) if i} - if any( - isinstance(dj, dict) - and dj.get("id") - and str(dj["id"]) not in payload_ids - and str(dj["id"]) not in intended - for dj in disk_jobs - ): - try: - os.unlink(tmp_path) - except OSError: - pass - tmp_path = None - continue + handle.flush() + os.fsync(handle.fileno()) - atomic_replace(tmp_path, jobs_file) - tmp_path = None - _secure_file(jobs_file) - _preserve_file_ownership(jobs_file, _stat_before) - # Invalidate (never refresh) the stamp after writing: the stamp - # certifies "this section's loaded payload still matches disk", - # which stops being provable the moment anyone writes. A refresh - # here would let a nested save (e.g. create_job inside a broader - # section) certify disk against an OUTER caller's stale payload - # and deterministically clobber the nested create; it also races - # a degraded sibling landing between replace and stat. Later - # saves in this section simply take the full merge (fail-safe). - _record_load_stamp(None) - return + if not replace and _jobs_file_stamp(jobs_file) != observed_stamp: + continue - # Exhausted retries — last merge + write without another re-peek. - if not replace: - jobs = _merge_unexpected_disk_jobs(jobs, removed_ids=removed_ids) - fd, tmp_path = tempfile.mkstemp( - dir=str(jobs_file.parent), suffix=".tmp", prefix=".jobs_" - ) - with os.fdopen(fd, "w", encoding="utf-8") as f: - json.dump( - {"jobs": jobs, "updated_at": _hermes_now().isoformat()}, - f, - indent=2, - ) - f.flush() - os.fsync(f.fileno()) - atomic_replace(tmp_path, jobs_file) - tmp_path = None - _secure_file(jobs_file) - _preserve_file_ownership(jobs_file, _stat_before) - except BaseException: - if tmp_path is not None: - try: - os.unlink(tmp_path) - except OSError: - pass - raise + published_path = Path(atomic_replace(tmp_path, jobs_file)) + published = True + _secure_file(published_path) + _preserve_file_ownership(published_path, stat_before) + if loaded is not None: + loaded.base = copy.deepcopy(desired) + loaded.stamp = None + loaded.stamp_trusted = False + return + finally: + if not published: + try: + os.unlink(tmp_path) + except OSError: + pass + + raise RuntimeError( + "jobs.json kept changing while this process prepared a save; " + "refusing to overwrite concurrent updates" + ) def save_jobs( @@ -1387,7 +1575,7 @@ def save_jobs( """Save all jobs to storage. See ``_save_jobs_unlocked`` for ``removed_ids`` / ``replace`` semantics - (shrink-merge guard against concurrent-create clobber, #80624). + (reconciliation guard against concurrent-update clobber). """ with _jobs_lock(): _save_jobs_unlocked(jobs, removed_ids=removed_ids, replace=replace) @@ -2001,7 +2189,7 @@ def remove_job(job_id: str) -> bool: with _jobs_lock(): jobs = load_jobs() original_len = len(jobs) - jobs = [j for j in jobs if j["id"] != canonical_id] + jobs[:] = [j for j in jobs if j["id"] != canonical_id] if len(jobs) < original_len: # Resolve the output dir BEFORE saving so a legacy unsafe ID (e.g. # left over from before the create-time guard) fails closed without diff --git a/hermes_cli/backup.py b/hermes_cli/backup.py index 43748fd3e207c..27e49ee920848 100644 --- a/hermes_cli/backup.py +++ b/hermes_cli/backup.py @@ -87,6 +87,8 @@ # File names to skip (runtime state that's meaningless on another machine) _EXCLUDED_NAMES = { ".backup.lock", + ".jobs.commit.lock", + ".jobs.lock", "gateway.pid", "cron.pid", } @@ -117,6 +119,8 @@ # Older backups predate the backup-side exclusions, so we filter on import too # rather than trusting the archive's contents. _IMPORT_SKIP_NAMES = { + ".jobs.commit.lock", + ".jobs.lock", "gateway_state.json", "gateway.pid", "cron.pid", diff --git a/tests/cron/test_jobs_conflict_safe_save.py b/tests/cron/test_jobs_conflict_safe_save.py new file mode 100644 index 0000000000000..84a8b3bf654b3 --- /dev/null +++ b/tests/cron/test_jobs_conflict_safe_save.py @@ -0,0 +1,353 @@ +from __future__ import annotations + +import copy +import errno +import importlib +import json +import multiprocessing +import os +from types import SimpleNamespace + +import pytest + + +@pytest.fixture +def jobs_env(tmp_path, monkeypatch): + home = tmp_path / ".hermes" + (home / "cron").mkdir(parents=True) + monkeypatch.setenv("HERMES_HOME", str(home)) + + import hermes_constants + import cron.jobs + + importlib.reload(hermes_constants) + importlib.reload(cron.jobs) + return cron.jobs + + +def _job(_jobs, job_id: str): + return { + "id": job_id, + "name": job_id, + "enabled": True, + "prompt": "test", + } + + +def _write_jobs(jobs, rows) -> None: + jobs._current_cron_store().jobs_file.write_text( + json.dumps({"jobs": rows}), encoding="utf-8" + ) + + +def test_three_way_merge_preserves_disjoint_changes(jobs_env): + jobs = jobs_env + base = [_job(jobs, job_id) for job_id in ("a", "b", "e")] + desired = [copy.deepcopy(base[1]), copy.deepcopy(base[2]), _job(jobs, "c")] + current = copy.deepcopy(base) + [_job(jobs, "d")] + desired[0]["enabled"] = False + current[2]["prompt"] = "current" + merged = jobs._merge_jobs_three_way(base, desired, current) + + by_id = {row["id"]: row for row in merged} + assert list(by_id) == ["b", "e", "c", "d"] + assert by_id["b"]["enabled"] is False + assert by_id["e"]["prompt"] == "current" + + +@pytest.mark.parametrize( + ("desired", "current"), + [ + ({"enabled": False}, {"prompt": "other"}), + (None, {"enabled": False}), + ({"enabled": False}, None), + ], +) +def test_three_way_merge_rejects_same_id_conflicts(jobs_env, desired, current): + jobs = jobs_env + base = [_job(jobs, "a")] + + def changed(delta): + if delta is None: + return [] + row = copy.deepcopy(base[0]) + row.update(delta) + return [row] + + with pytest.raises(RuntimeError, match="conflicting concurrent cron job"): + jobs._merge_jobs_three_way(base, changed(desired), changed(current)) + + +@pytest.mark.parametrize("ambiguous", ["missing", "duplicate"]) +def test_ambiguous_rows_fail_closed_during_concurrent_merge(jobs_env, ambiguous): + jobs = jobs_env + first = _job(jobs, "a") if ambiguous == "duplicate" else "legacy" + base = [first, _job(jobs, "a"), _job(jobs, "b")] + desired = copy.deepcopy(base) + current = copy.deepcopy(base) + desired[1]["enabled"] = False + current[2]["prompt"] = "current" + + with pytest.raises(RuntimeError, match="missing or duplicate ids"): + jobs._merge_jobs_three_way(base, desired, current) + + +@pytest.mark.parametrize("sibling_change", ["pause", "delete"]) +def test_loaded_save_preserves_sibling_pause_or_delete(jobs_env, sibling_change): + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + + with jobs._jobs_lock(): + stale = jobs.load_jobs() + current = copy.deepcopy(stale) + if sibling_change == "pause": + current[0]["enabled"] = False + else: + current = [] + _write_jobs(jobs, current) + jobs._save_jobs_unlocked(stale) + + assert jobs.load_jobs() == current + + +def test_unknown_stamps_never_certify_a_stale_loaded_snapshot(jobs_env, monkeypatch): + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + + with jobs._jobs_lock(): + monkeypatch.setattr(jobs, "_jobs_file_stamp", lambda _path: None) + stale = jobs.load_jobs() + current = copy.deepcopy(stale) + current[0]["enabled"] = False + _write_jobs(jobs, current) + jobs._save_jobs_unlocked(stale) + + assert jobs.load_jobs() == current + + +@pytest.mark.skipif(os.name != "posix", reason="POSIX ownership semantics") +def test_commit_lock_preserves_cron_directory_owner(jobs_env, monkeypatch): + jobs = jobs_env + calls = [] + monkeypatch.setattr(jobs.os, "geteuid", lambda: 0) + monkeypatch.setattr( + jobs.os, + "fchown", + lambda fd, uid, gid: calls.append((fd, uid, gid)), + ) + + with jobs._jobs_commit_lock(): + pass + + directory = os.stat(jobs._current_cron_store().cron_dir) + assert len(calls) == 1 + assert calls[0][1:] == (directory.st_uid, directory.st_gid) + assert jobs._jobs_commit_lock_file().stat().st_mode & 0o777 == 0o600 + + +@pytest.mark.skipif( + os.name != "posix" or not hasattr(os, "O_NOFOLLOW"), + reason="POSIX no-follow semantics", +) +def test_commit_lock_refuses_symlink_before_chmod(jobs_env): + jobs = jobs_env + victim = jobs._current_cron_store().cron_dir / "victim" + victim.write_text("unchanged", encoding="utf-8") + victim.chmod(0o644) + jobs._jobs_commit_lock_file().symlink_to(victim) + + with pytest.raises(RuntimeError, match="Unable to open"): + with jobs._jobs_commit_lock(): + pass + + assert victim.stat().st_mode & 0o777 == 0o644 + + +def test_windows_commit_lock_seeds_and_locks_byte_zero(jobs_env, monkeypatch): + jobs = jobs_env + calls = [] + fake = SimpleNamespace(LK_NBLCK=1, LK_UNLCK=2) + + def locking(fd, mode, size): + calls.append((os.lseek(fd, 0, os.SEEK_CUR), mode, size)) + + fake.locking = locking + monkeypatch.setattr(jobs, "fcntl", None) + monkeypatch.setattr(jobs, "msvcrt", fake) + + with jobs._jobs_commit_lock(): + assert jobs._jobs_commit_lock_file().read_bytes() == b" " + + assert calls == [(0, fake.LK_NBLCK, 1), (0, fake.LK_UNLCK, 1)] + + +def test_commit_lock_without_backend_uses_historical_fallback(jobs_env, monkeypatch): + jobs = jobs_env + monkeypatch.setattr(jobs, "fcntl", None) + monkeypatch.setattr(jobs, "msvcrt", None) + + with jobs._jobs_commit_lock(): + pass + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +def test_unsupported_flock_uses_generation_fallback(jobs_env, monkeypatch): + jobs = jobs_env + + def unsupported(_fd, _operation): + raise OSError(errno.ENOTSUP, "unsupported") + + monkeypatch.setattr(jobs.fcntl, "flock", unsupported) + with jobs._jobs_commit_lock(): + pass + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +def test_unlock_error_does_not_fail_published_save(jobs_env, monkeypatch): + jobs = jobs_env + real_flock = jobs.fcntl.flock + + def fail_unlock(fd, operation): + if operation == jobs.fcntl.LOCK_UN: + raise OSError("synthetic unlock failure") + return real_flock(fd, operation) + + with jobs._jobs_lock(): + monkeypatch.setattr(jobs.fcntl, "flock", fail_unlock) + jobs._save_jobs_unlocked([_job(jobs, "saved")], replace=True) + + payload = json.loads(jobs._current_cron_store().jobs_file.read_text()) + assert [row["id"] for row in payload["jobs"]] == ["saved"] + + +def test_loaded_save_honors_removed_ids_for_current_only_job(jobs_env): + jobs = jobs_env + jobs.save_jobs([_job(jobs, "a")], replace=True) + + with jobs._jobs_lock(): + desired = jobs.load_jobs() + _write_jobs(jobs, desired + [_job(jobs, "external")]) + jobs._save_jobs_unlocked(desired, removed_ids={"external"}) + + assert [row["id"] for row in jobs.load_jobs()] == ["a"] + + +def test_generation_churn_fails_closed_without_final_unchecked_write( + jobs_env, monkeypatch +): + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + + with jobs._jobs_lock(): + desired = jobs.load_jobs() + desired[0]["prompt"] = "desired" + real_stamp = jobs._jobs_file_stamp + calls = 0 + + def churning_stamp(path): + nonlocal calls + calls += 1 + if calls in {2, 5, 8}: + _write_jobs(jobs, [seed, _job(jobs, f"external-{calls}")]) + return real_stamp(path) + + monkeypatch.setattr(jobs, "_jobs_file_stamp", churning_stamp) + with pytest.raises(RuntimeError, match="kept changing"): + jobs._save_jobs_unlocked(desired) + + stored = jobs.load_jobs() + assert stored[0]["prompt"] == "test" + assert stored[1]["id"].startswith("external-") + assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +@pytest.mark.parametrize( + "lock_error", + [BlockingIOError(), OSError(errno.ENOLCK, "lock records unavailable")], + ids=["contended", "enolck"], +) +def test_commit_lock_contention_fails_closed(jobs_env, monkeypatch, lock_error): + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + before = jobs._current_cron_store().jobs_file.read_bytes() + real_flock = jobs.fcntl.flock + + def contend(fd, operation): + if operation & jobs.fcntl.LOCK_NB: + raise lock_error + return real_flock(fd, operation) + + with jobs._jobs_lock(): + monkeypatch.setattr(jobs.fcntl, "flock", contend) + monkeypatch.setattr(jobs, "_JOBS_COMMIT_LOCK_TIMEOUT_SECONDS", 0.0) + with pytest.raises(RuntimeError, match="Timed out"): + jobs._save_jobs_unlocked([_job(jobs, "b")]) + + assert jobs._current_cron_store().jobs_file.read_bytes() == before + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX flock subprocess test") +def test_real_degraded_writers_serialize_publication(jobs_env): + jobs = jobs_env + jobs.save_jobs([_job(jobs, "a"), _job(jobs, "b")], replace=True) + context = multiprocessing.get_context("fork") + a_loaded = context.Event() + b_loaded = context.Event() + start_saves = context.Event() + a_committing = context.Event() + b_committing = context.Event() + a_done = context.Event() + b_done = context.Event() + + def writer(loaded_event, committing_event, done_event, first): + jobs._JOBS_LOCK_TIMEOUT_SECONDS = 0.05 + with jobs._jobs_lock(): + loaded = jobs.load_jobs() + if first: + loaded[0]["prompt"] = "writer-a" + loaded.append(_job(jobs, "c")) + else: + loaded[1]["enabled"] = False + loaded.append(_job(jobs, "d")) + loaded_event.set() + if not start_saves.wait(10): + raise RuntimeError("writer start timed out") + committing_event.set() + jobs.save_jobs(loaded) + done_event.set() + + processes = [ + context.Process( + target=writer, args=(a_loaded, a_committing, a_done, True) + ), + context.Process( + target=writer, args=(b_loaded, b_committing, b_done, False) + ), + ] + processes[0].start() + assert a_loaded.wait(5) + processes[1].start() + assert b_loaded.wait(5) + with jobs._jobs_commit_lock(): + start_saves.set() + assert a_committing.wait(5) + assert b_committing.wait(5) + assert not a_done.is_set() + assert not b_done.is_set() + + for process in processes: + process.join(15) + if process.is_alive(): + process.terminate() + process.join(5) + assert process.exitcode == 0 + + stored = {row["id"]: row for row in jobs.load_jobs()} + assert stored["a"]["prompt"] == "writer-a" + assert stored["b"]["enabled"] is False + assert {"c", "d"} <= stored.keys() diff --git a/tests/cron/test_jobs_shrink_merge_80624.py b/tests/cron/test_jobs_shrink_merge_80624.py index 423c3fe93d337..5fe2b1ab59a10 100644 --- a/tests/cron/test_jobs_shrink_merge_80624.py +++ b/tests/cron/test_jobs_shrink_merge_80624.py @@ -139,7 +139,7 @@ def test_jobs_json_on_disk_matches_merge(hermes_env): def test_stamp_fast_path_skips_merge_when_file_unchanged(hermes_env, monkeypatch): """Inside a critical section whose load stamp still matches, the save - must not re-read jobs.json at all (#80703's single-stat fast path).""" + must not re-read jobs.json at all (#80703's no-parse fast path).""" import cron.jobs as jobs from cron.jobs import create_job diff --git a/tests/hermes_cli/test_backup.py b/tests/hermes_cli/test_backup.py index 0fb69caa3ec64..c13cb4722f748 100644 --- a/tests/hermes_cli/test_backup.py +++ b/tests/hermes_cli/test_backup.py @@ -133,6 +133,12 @@ def test_excludes_sqlite_sidecars(self): # The .db itself is still included (and safe-copied separately) assert not _should_exclude(Path("state.db")) + def test_excludes_cron_runtime_locks(self): + from hermes_cli.backup import _should_exclude + + assert _should_exclude(Path("cron/.jobs.lock")) + assert _should_exclude(Path("cron/.jobs.commit.lock")) + # --------------------------------------------------------------------------- # Backup tests @@ -303,7 +309,9 @@ def test_preserves_per_profile_gateway_state(self, tmp_path, monkeypatch): def test_preserves_runtime_pid_and_process_files(self, tmp_path, monkeypatch): """gateway.pid / cron.pid / gateway.lock / processes.json from a backup reference the source machine's process namespace and must never be - written over the target's.""" + written over the target's. Cron lock files are likewise runtime state; + replacing them would split the target's live lock domain. + """ hermes_home = tmp_path / ".hermes" hermes_home.mkdir() monkeypatch.setenv("HERMES_HOME", str(hermes_home)) @@ -320,6 +328,8 @@ def test_preserves_runtime_pid_and_process_files(self, tmp_path, monkeypatch): "cron.pid": "8888", "gateway.lock": "7777", "processes.json": '{"stale": true}', + "cron/.jobs.lock": "stale", + "cron/.jobs.commit.lock": "stale", }) args = Namespace(zipfile=str(zip_path), force=True) @@ -333,6 +343,8 @@ def test_preserves_runtime_pid_and_process_files(self, tmp_path, monkeypatch): # cron.pid / gateway.lock had no live copy and were not seeded. assert not (hermes_home / "cron.pid").exists() assert not (hermes_home / "gateway.lock").exists() + assert not (hermes_home / "cron" / ".jobs.lock").exists() + assert not (hermes_home / "cron" / ".jobs.commit.lock").exists() @@ -1245,4 +1257,3 @@ def test_import_restores_external_to_home_relative_location(self, tmp_path, monk - From 0bb09dd4602c63f56584234830716f6812da44f2 Mon Sep 17 00:00:00 2001 From: poisdahl <4091911+poisdahl@users.noreply.github.com> Date: Wed, 12 Aug 2026 22:07:39 +0200 Subject: [PATCH 2/5] fix(cron): fail closed without commit locking --- cron/jobs.py | 25 ++++--- tests/cron/test_jobs_conflict_safe_save.py | 85 ++++++++++++++++++++-- 2 files changed, 94 insertions(+), 16 deletions(-) diff --git a/cron/jobs.py b/cron/jobs.py index 9627558eba0a9..4a62f9973145d 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -314,10 +314,18 @@ def _prepare_commit_lock_file(handle, owner_source: Optional[os.stat_result]) -> @contextlib.contextmanager def _jobs_commit_lock(): - """Serialize the final read, reconcile, and publication of jobs.json.""" + """Serialize the final read, reconcile, and publication of jobs.json. + + Unlike the broad scheduler lock, this short publication lock must never + degrade to process-only locking. The generation check cannot close the + check-to-replace window by itself, so publishing without a cross-process + backend would reintroduce the lost-update race this lock exists to close. + """ if fcntl is None and msvcrt is None: - yield - return + raise RuntimeError( + "Cron jobs commit lock is unavailable; refusing to publish " + "without cross-process serialization" + ) ensure_dirs() lock_path = _jobs_commit_lock_file() @@ -379,13 +387,10 @@ def _jobs_commit_lock(): time.sleep(0.01) if backend_error is not None: - logger.warning( - "Cron jobs commit lock is unavailable (%s); proceeding without " - "cross-process publication locking", - backend_error, - ) - yield - return + raise RuntimeError( + "Cron jobs commit lock is unsupported; refusing to publish " + "without cross-process serialization" + ) from backend_error yield finally: try: diff --git a/tests/cron/test_jobs_conflict_safe_save.py b/tests/cron/test_jobs_conflict_safe_save.py index 84a8b3bf654b3..678e3ee44514c 100644 --- a/tests/cron/test_jobs_conflict_safe_save.py +++ b/tests/cron/test_jobs_conflict_safe_save.py @@ -183,25 +183,98 @@ def locking(fd, mode, size): assert calls == [(0, fake.LK_NBLCK, 1), (0, fake.LK_UNLCK, 1)] -def test_commit_lock_without_backend_uses_historical_fallback(jobs_env, monkeypatch): +def test_commit_lock_without_backend_fails_closed(jobs_env, monkeypatch): jobs = jobs_env monkeypatch.setattr(jobs, "fcntl", None) monkeypatch.setattr(jobs, "msvcrt", None) - with jobs._jobs_commit_lock(): - pass + with pytest.raises(RuntimeError, match="unavailable.*refusing to publish"): + with jobs._jobs_commit_lock(): + pass + + +@pytest.mark.parametrize("replace", [False, True]) +def test_missing_commit_lock_backend_never_publishes( + jobs_env, monkeypatch, replace +): + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + before = jobs._current_cron_store().jobs_file.read_bytes() + + with jobs._jobs_lock(): + desired = jobs.load_jobs() + desired[0]["prompt"] = "must-not-publish" + replace_calls = [] + monkeypatch.setattr(jobs, "fcntl", None) + monkeypatch.setattr(jobs, "msvcrt", None) + monkeypatch.setattr( + jobs, + "atomic_replace", + lambda *args: replace_calls.append(args), + ) + with pytest.raises(RuntimeError, match="unavailable.*refusing to publish"): + jobs._save_jobs_unlocked(desired, replace=replace) + + assert replace_calls == [] + assert jobs._current_cron_store().jobs_file.read_bytes() == before + assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) @pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") -def test_unsupported_flock_uses_generation_fallback(jobs_env, monkeypatch): +def test_unsupported_flock_fails_closed(jobs_env, monkeypatch): jobs = jobs_env def unsupported(_fd, _operation): raise OSError(errno.ENOTSUP, "unsupported") monkeypatch.setattr(jobs.fcntl, "flock", unsupported) - with jobs._jobs_commit_lock(): - pass + with pytest.raises(RuntimeError, match="unsupported.*refusing to publish"): + with jobs._jobs_commit_lock(): + pass + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +@pytest.mark.parametrize("replace", [False, True]) +@pytest.mark.parametrize( + "lock_errno_name", + [ + name + for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP") + if getattr(errno, name, None) is not None + ], +) +def test_unsupported_commit_lock_never_publishes( + jobs_env, monkeypatch, replace, lock_errno_name +): + """A generation check alone cannot protect check-to-replace publication.""" + jobs = jobs_env + seed = _job(jobs, "a") + jobs.save_jobs([seed], replace=True) + before = jobs._current_cron_store().jobs_file.read_bytes() + + def unsupported(_fd, _operation): + raise OSError(getattr(errno, lock_errno_name), "unsupported") + + with jobs._jobs_lock(): + desired = jobs.load_jobs() + desired[0]["prompt"] = "must-not-publish" + replace_calls = [] + # Acquire the broad mutation lock first; only the short publication + # lock is under test here. Patching flock before _jobs_lock() would + # exercise its bounded retry path and obscure the commit invariant. + monkeypatch.setattr(jobs.fcntl, "flock", unsupported) + monkeypatch.setattr( + jobs, + "atomic_replace", + lambda *args: replace_calls.append(args), + ) + with pytest.raises(RuntimeError, match="unsupported.*refusing to publish"): + jobs._save_jobs_unlocked(desired, replace=replace) + + assert replace_calls == [] + assert jobs._current_cron_store().jobs_file.read_bytes() == before + assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) @pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") From 06ebd6c8132a71b20e29ea60e6f0934f4ad19c84 Mon Sep 17 00:00:00 2001 From: poisdahl <4091911+poisdahl@users.noreply.github.com> Date: Thu, 13 Aug 2026 02:37:00 +0200 Subject: [PATCH 3/5] fix(cron): preserve degraded save availability --- cron/jobs.py | 29 ++++--- tests/cron/test_jobs_conflict_safe_save.py | 98 +++++++++------------- 2 files changed, 58 insertions(+), 69 deletions(-) diff --git a/cron/jobs.py b/cron/jobs.py index 4a62f9973145d..6c8ac24dd47cc 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -316,16 +316,18 @@ def _prepare_commit_lock_file(handle, owner_source: Optional[os.stat_result]) -> def _jobs_commit_lock(): """Serialize the final read, reconcile, and publication of jobs.json. - Unlike the broad scheduler lock, this short publication lock must never - degrade to process-only locking. The generation check cannot close the - check-to-replace window by itself, so publishing without a cross-process - backend would reintroduce the lost-update race this lock exists to close. + This short lock closes the generation-check-to-replace window when the + platform supports it. Match the broad scheduler lock's availability + contract when no cross-process backend exists: keep the in-process lock and + bounded generation reconciliation rather than disabling cron writes. """ if fcntl is None and msvcrt is None: - raise RuntimeError( - "Cron jobs commit lock is unavailable; refusing to publish " - "without cross-process serialization" + logger.warning( + "Cron jobs commit lock is unavailable; publishing with " + "process-local locking and generation checks only" ) + yield False + return ensure_dirs() lock_path = _jobs_commit_lock_file() @@ -387,11 +389,14 @@ def _jobs_commit_lock(): time.sleep(0.01) if backend_error is not None: - raise RuntimeError( - "Cron jobs commit lock is unsupported; refusing to publish " - "without cross-process serialization" - ) from backend_error - yield + logger.warning( + "Cron jobs commit lock is unsupported; publishing with " + "process-local locking and generation checks only: %s", + backend_error, + ) + yield False + return + yield True finally: try: if acquired: diff --git a/tests/cron/test_jobs_conflict_safe_save.py b/tests/cron/test_jobs_conflict_safe_save.py index 678e3ee44514c..d53d36a30c9fb 100644 --- a/tests/cron/test_jobs_conflict_safe_save.py +++ b/tests/cron/test_jobs_conflict_safe_save.py @@ -6,8 +6,6 @@ import json import multiprocessing import os -from types import SimpleNamespace - import pytest @@ -127,7 +125,7 @@ def test_unknown_stamps_never_certify_a_stale_loaded_snapshot(jobs_env, monkeypa assert jobs.load_jobs() == current -@pytest.mark.skipif(os.name != "posix", reason="POSIX ownership semantics") +@pytest.mark.linux_only def test_commit_lock_preserves_cron_directory_owner(jobs_env, monkeypatch): jobs = jobs_env calls = [] @@ -147,10 +145,7 @@ def test_commit_lock_preserves_cron_directory_owner(jobs_env, monkeypatch): assert jobs._jobs_commit_lock_file().stat().st_mode & 0o777 == 0o600 -@pytest.mark.skipif( - os.name != "posix" or not hasattr(os, "O_NOFOLLOW"), - reason="POSIX no-follow semantics", -) +@pytest.mark.linux_only def test_commit_lock_refuses_symlink_before_chmod(jobs_env): jobs = jobs_env victim = jobs._current_cron_store().cron_dir / "victim" @@ -165,76 +160,74 @@ def test_commit_lock_refuses_symlink_before_chmod(jobs_env): assert victim.stat().st_mode & 0o777 == 0o644 +@pytest.mark.windows_only def test_windows_commit_lock_seeds_and_locks_byte_zero(jobs_env, monkeypatch): jobs = jobs_env calls = [] - fake = SimpleNamespace(LK_NBLCK=1, LK_UNLCK=2) + real_locking = jobs.msvcrt.locking def locking(fd, mode, size): calls.append((os.lseek(fd, 0, os.SEEK_CUR), mode, size)) + return real_locking(fd, mode, size) - fake.locking = locking - monkeypatch.setattr(jobs, "fcntl", None) - monkeypatch.setattr(jobs, "msvcrt", fake) + monkeypatch.setattr(jobs.msvcrt, "locking", locking) with jobs._jobs_commit_lock(): assert jobs._jobs_commit_lock_file().read_bytes() == b" " - assert calls == [(0, fake.LK_NBLCK, 1), (0, fake.LK_UNLCK, 1)] + assert calls == [ + (0, jobs.msvcrt.LK_NBLCK, 1), + (0, jobs.msvcrt.LK_UNLCK, 1), + ] -def test_commit_lock_without_backend_fails_closed(jobs_env, monkeypatch): +def test_commit_lock_without_backend_degrades_with_warning( + jobs_env, monkeypatch, caplog +): jobs = jobs_env monkeypatch.setattr(jobs, "fcntl", None) monkeypatch.setattr(jobs, "msvcrt", None) - with pytest.raises(RuntimeError, match="unavailable.*refusing to publish"): - with jobs._jobs_commit_lock(): - pass + with jobs._jobs_commit_lock() as acquired: + assert acquired is False + + assert "process-local locking and generation checks only" in caplog.text @pytest.mark.parametrize("replace", [False, True]) -def test_missing_commit_lock_backend_never_publishes( - jobs_env, monkeypatch, replace +def test_missing_commit_lock_backend_still_publishes( + jobs_env, monkeypatch, caplog, replace ): jobs = jobs_env seed = _job(jobs, "a") jobs.save_jobs([seed], replace=True) - before = jobs._current_cron_store().jobs_file.read_bytes() - with jobs._jobs_lock(): desired = jobs.load_jobs() - desired[0]["prompt"] = "must-not-publish" - replace_calls = [] + desired[0]["prompt"] = "degraded-save" monkeypatch.setattr(jobs, "fcntl", None) monkeypatch.setattr(jobs, "msvcrt", None) - monkeypatch.setattr( - jobs, - "atomic_replace", - lambda *args: replace_calls.append(args), - ) - with pytest.raises(RuntimeError, match="unavailable.*refusing to publish"): - jobs._save_jobs_unlocked(desired, replace=replace) - - assert replace_calls == [] - assert jobs._current_cron_store().jobs_file.read_bytes() == before + jobs._save_jobs_unlocked(desired, replace=replace) + + assert jobs.load_jobs()[0]["prompt"] == "degraded-save" + assert "process-local locking and generation checks only" in caplog.text assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) -@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") -def test_unsupported_flock_fails_closed(jobs_env, monkeypatch): +@pytest.mark.linux_only +def test_unsupported_flock_degrades_with_warning(jobs_env, monkeypatch, caplog): jobs = jobs_env def unsupported(_fd, _operation): raise OSError(errno.ENOTSUP, "unsupported") monkeypatch.setattr(jobs.fcntl, "flock", unsupported) - with pytest.raises(RuntimeError, match="unsupported.*refusing to publish"): - with jobs._jobs_commit_lock(): - pass + with jobs._jobs_commit_lock() as acquired: + assert acquired is False + assert "process-local locking and generation checks only" in caplog.text -@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") + +@pytest.mark.linux_only @pytest.mark.parametrize("replace", [False, True]) @pytest.mark.parametrize( "lock_errno_name", @@ -244,40 +237,31 @@ def unsupported(_fd, _operation): if getattr(errno, name, None) is not None ], ) -def test_unsupported_commit_lock_never_publishes( - jobs_env, monkeypatch, replace, lock_errno_name +def test_unsupported_commit_lock_still_publishes( + jobs_env, monkeypatch, caplog, replace, lock_errno_name ): - """A generation check alone cannot protect check-to-replace publication.""" jobs = jobs_env seed = _job(jobs, "a") jobs.save_jobs([seed], replace=True) - before = jobs._current_cron_store().jobs_file.read_bytes() def unsupported(_fd, _operation): raise OSError(getattr(errno, lock_errno_name), "unsupported") with jobs._jobs_lock(): desired = jobs.load_jobs() - desired[0]["prompt"] = "must-not-publish" - replace_calls = [] + desired[0]["prompt"] = "degraded-save" # Acquire the broad mutation lock first; only the short publication # lock is under test here. Patching flock before _jobs_lock() would # exercise its bounded retry path and obscure the commit invariant. monkeypatch.setattr(jobs.fcntl, "flock", unsupported) - monkeypatch.setattr( - jobs, - "atomic_replace", - lambda *args: replace_calls.append(args), - ) - with pytest.raises(RuntimeError, match="unsupported.*refusing to publish"): - jobs._save_jobs_unlocked(desired, replace=replace) - - assert replace_calls == [] - assert jobs._current_cron_store().jobs_file.read_bytes() == before + jobs._save_jobs_unlocked(desired, replace=replace) + + assert jobs.load_jobs()[0]["prompt"] == "degraded-save" + assert "process-local locking and generation checks only" in caplog.text assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) -@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +@pytest.mark.linux_only def test_unlock_error_does_not_fail_published_save(jobs_env, monkeypatch): jobs = jobs_env real_flock = jobs.fcntl.flock @@ -337,7 +321,7 @@ def churning_stamp(path): assert not list(jobs._current_cron_store().cron_dir.glob(".jobs_*.tmp")) -@pytest.mark.skipif(os.name == "nt", reason="POSIX flock semantics") +@pytest.mark.linux_only @pytest.mark.parametrize( "lock_error", [BlockingIOError(), OSError(errno.ENOLCK, "lock records unavailable")], @@ -364,7 +348,7 @@ def contend(fd, operation): assert jobs._current_cron_store().jobs_file.read_bytes() == before -@pytest.mark.skipif(os.name == "nt", reason="POSIX flock subprocess test") +@pytest.mark.linux_only def test_real_degraded_writers_serialize_publication(jobs_env): jobs = jobs_env jobs.save_jobs([_job(jobs, "a"), _job(jobs, "b")], replace=True) From e023ce7a8b21d2a8a106003e3af2d8d37c10f225 Mon Sep 17 00:00:00 2001 From: poisdahl <4091911+poisdahl@users.noreply.github.com> Date: Thu, 13 Aug 2026 02:40:00 +0200 Subject: [PATCH 4/5] fix(cron): degrade when commit locking is unavailable --- cron/jobs.py | 31 +++++++++++++++------ tests/cron/test_jobs_conflict_safe_save.py | 32 ++++++++++------------ 2 files changed, 37 insertions(+), 26 deletions(-) diff --git a/cron/jobs.py b/cron/jobs.py index 6c8ac24dd47cc..9f042de4e3c08 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -116,9 +116,9 @@ def _ensure_croniter() -> bool: _JOBS_LOCK_TIMEOUT_SECONDS = 30.0 _JOBS_COMMIT_LOCK_TIMEOUT_SECONDS = 5.0 _JOBS_GENERATION_MAX_ATTEMPTS = 3 -_UNSUPPORTED_LOCK_ERRNOS = frozenset( +_UNAVAILABLE_LOCK_ERRNOS = frozenset( value - for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP") + for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP", "ENOLCK") if (value := getattr(errno, name, None)) is not None ) OUTPUT_DIR = CRON_DIR / "output" @@ -346,12 +346,27 @@ def _jobs_commit_lock(): handle = os.fdopen(fd, "r+b") fd = None except OSError as exc: - raise RuntimeError("Unable to open the cron jobs commit lock") from exc + logger.warning( + "Cron jobs commit lock is unavailable; publishing with " + "process-local locking and generation checks only: %s", + exc, + ) + yield False + return finally: if fd is not None: os.close(fd) try: _prepare_commit_lock_file(handle, owner_source) + except OSError as exc: + handle.close() + logger.warning( + "Cron jobs commit lock is unavailable; publishing with " + "process-local locking and generation checks only: %s", + exc, + ) + yield False + return except BaseException: handle.close() raise @@ -378,19 +393,19 @@ def _jobs_commit_lock(): acquired = True break except (BlockingIOError, OSError) as exc: - if getattr(exc, "errno", None) in _UNSUPPORTED_LOCK_ERRNOS: + if getattr(exc, "errno", None) in _UNAVAILABLE_LOCK_ERRNOS: backend_error = exc break if time.monotonic() >= deadline: - raise RuntimeError( - "Timed out waiting for the cron jobs commit lock; " - "refusing to publish a potentially stale snapshot" + backend_error = RuntimeError( + "timed out waiting for the cron jobs commit lock" ) + break time.sleep(0.01) if backend_error is not None: logger.warning( - "Cron jobs commit lock is unsupported; publishing with " + "Cron jobs commit lock is unavailable; publishing with " "process-local locking and generation checks only: %s", backend_error, ) diff --git a/tests/cron/test_jobs_conflict_safe_save.py b/tests/cron/test_jobs_conflict_safe_save.py index d53d36a30c9fb..5be6a3034a6b1 100644 --- a/tests/cron/test_jobs_conflict_safe_save.py +++ b/tests/cron/test_jobs_conflict_safe_save.py @@ -146,18 +146,18 @@ def test_commit_lock_preserves_cron_directory_owner(jobs_env, monkeypatch): @pytest.mark.linux_only -def test_commit_lock_refuses_symlink_before_chmod(jobs_env): +def test_commit_lock_symlink_degrades_without_chmod(jobs_env, caplog): jobs = jobs_env victim = jobs._current_cron_store().cron_dir / "victim" victim.write_text("unchanged", encoding="utf-8") victim.chmod(0o644) jobs._jobs_commit_lock_file().symlink_to(victim) - with pytest.raises(RuntimeError, match="Unable to open"): - with jobs._jobs_commit_lock(): - pass + with jobs._jobs_commit_lock() as acquired: + assert acquired is False assert victim.stat().st_mode & 0o777 == 0o644 + assert "process-local locking and generation checks only" in caplog.text @pytest.mark.windows_only @@ -214,7 +214,7 @@ def test_missing_commit_lock_backend_still_publishes( @pytest.mark.linux_only -def test_unsupported_flock_degrades_with_warning(jobs_env, monkeypatch, caplog): +def test_unavailable_flock_degrades_with_warning(jobs_env, monkeypatch, caplog): jobs = jobs_env def unsupported(_fd, _operation): @@ -233,11 +233,11 @@ def unsupported(_fd, _operation): "lock_errno_name", [ name - for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP") + for name in ("ENOSYS", "ENOTSUP", "EOPNOTSUPP", "ENOLCK") if getattr(errno, name, None) is not None ], ) -def test_unsupported_commit_lock_still_publishes( +def test_unavailable_commit_lock_still_publishes( jobs_env, monkeypatch, caplog, replace, lock_errno_name ): jobs = jobs_env @@ -322,30 +322,26 @@ def churning_stamp(path): @pytest.mark.linux_only -@pytest.mark.parametrize( - "lock_error", - [BlockingIOError(), OSError(errno.ENOLCK, "lock records unavailable")], - ids=["contended", "enolck"], -) -def test_commit_lock_contention_fails_closed(jobs_env, monkeypatch, lock_error): +def test_commit_lock_timeout_degrades_with_warning(jobs_env, monkeypatch, caplog): jobs = jobs_env seed = _job(jobs, "a") jobs.save_jobs([seed], replace=True) - before = jobs._current_cron_store().jobs_file.read_bytes() real_flock = jobs.fcntl.flock def contend(fd, operation): if operation & jobs.fcntl.LOCK_NB: - raise lock_error + raise BlockingIOError() return real_flock(fd, operation) with jobs._jobs_lock(): + desired = jobs.load_jobs() + desired[0]["prompt"] = "degraded-save" monkeypatch.setattr(jobs.fcntl, "flock", contend) monkeypatch.setattr(jobs, "_JOBS_COMMIT_LOCK_TIMEOUT_SECONDS", 0.0) - with pytest.raises(RuntimeError, match="Timed out"): - jobs._save_jobs_unlocked([_job(jobs, "b")]) + jobs._save_jobs_unlocked(desired) - assert jobs._current_cron_store().jobs_file.read_bytes() == before + assert jobs.load_jobs()[0]["prompt"] == "degraded-save" + assert "process-local locking and generation checks only" in caplog.text @pytest.mark.linux_only From d108a49d9c22ae7968f9dfd3b0887c4a9264616b Mon Sep 17 00:00:00 2001 From: poisdahl <4091911+poisdahl@users.noreply.github.com> Date: Thu, 13 Aug 2026 02:54:50 +0200 Subject: [PATCH 5/5] test(cron): inspect Windows lock file after unlock --- tests/cron/test_jobs_conflict_safe_save.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/cron/test_jobs_conflict_safe_save.py b/tests/cron/test_jobs_conflict_safe_save.py index 5be6a3034a6b1..fff8e0b9db7a2 100644 --- a/tests/cron/test_jobs_conflict_safe_save.py +++ b/tests/cron/test_jobs_conflict_safe_save.py @@ -173,8 +173,9 @@ def locking(fd, mode, size): monkeypatch.setattr(jobs.msvcrt, "locking", locking) with jobs._jobs_commit_lock(): - assert jobs._jobs_commit_lock_file().read_bytes() == b" " + pass + assert jobs._jobs_commit_lock_file().read_bytes() == b" " assert calls == [ (0, jobs.msvcrt.LK_NBLCK, 1), (0, jobs.msvcrt.LK_UNLCK, 1),