diff --git a/CHANGELOG.md b/CHANGELOG.md index 50de4fb..482b799 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,7 @@ and versions are tracked in the repo-root `VERSION` file. ### Changed +- Skip contended retention passes instead of blocking CLI invocations on housekeeping locks (#386). - Reuse secure log lock descriptors and cache source paths per invocation; logging I/O failures stay inside logging (#381). - Align the Typer support floor with the tested matrix and cover representative minimum/maximum Typer and Click version pairings. diff --git a/docs/performance.md b/docs/performance.md index 5bb71f8..81a1e9e 100644 --- a/docs/performance.md +++ b/docs/performance.md @@ -137,3 +137,17 @@ with an advisory lock around each append and a fresh descriptor after fork. Human formatters cache up to 256 source paths for the current invocation and project binding. Repeated paths require no filesystem resolution. Sidecar I/O errors are routed through `logging.Handler.handleError` and do not fail commands. + +### Concurrent retention + +Retention takes a nonblocking maintenance lock on POSIX and Windows. A busy lock +skips that pass at debug level; the next successful invocation reconciles policy +debt. Startup and teardown both acquire the lock before scanning. Command work +never waits for a stopped lock holder. Deletion still revalidates metadata and +leases under the lock. + +`RetentionPolicy.safe_defaults()` includes `max_total_bytes=512 MiB`, so the +default policy uses the byte-policy recursive-walk bounds described above. A +consumer that needs only count/age retention can explicitly omit the byte cap. +The concurrent benchmark in #391 measures twelve processes against one warmed +cache and gates the p95-to-serial ratio. diff --git a/lib/python/base_cli/_runtime.py b/lib/python/base_cli/_runtime.py index 3d31fdf..2b5e603 100644 --- a/lib/python/base_cli/_runtime.py +++ b/lib/python/base_cli/_runtime.py @@ -446,6 +446,8 @@ def prune_run_bundles( now=clock, size_scan_cursor=size_scan_cursor, ) + except BlockingIOError: + log.debug("Skipping run bundle retention: another invocation holds the maintenance lock.") except (OSError, RuntimeError) as exc: # Retention is maintenance. An unavailable lock or a transient # filesystem failure must not turn an otherwise valid invocation into @@ -459,23 +461,31 @@ def refresh_run_bundle_index( current_run_root: Path | None = None, logger: logging.Logger | None = None, ) -> None: - """Refresh the diagnostic bundle index after a run becomes terminal.""" + """Refresh the index after a run becomes terminal. + + A concurrent maintenance pass may win the nonblocking lock. In that case + the index is eventually consistent and the next foreground pass retries it. + """ log = logger or logging.getLogger(__name__) runs_root = Path(runs_root) if not runs_root.exists() or runs_root.is_symlink(): return try: - bundles, _size_scan_cursor = _discover_run_bundles( - runs_root, - protected=set(), - max_age_seconds=None, - now=time.time(), - measure_sizes=False, - size_budget=0, - ) with _retention_lock(runs_root): + bundles, _size_scan_cursor = _discover_run_bundles( + runs_root, + protected=set(), + max_age_seconds=None, + now=time.time(), + measure_sizes=False, + size_budget=0, + ) _write_run_index(runs_root, bundles, log, current_run_root=current_run_root) + except BlockingIOError: + log.debug( + "Skipping run bundle index refresh under '%s': another invocation holds the maintenance lock.", runs_root + ) except (OSError, RuntimeError) as exc: log.debug("Could not refresh run bundle index under '%s': %s", runs_root, exc) @@ -922,12 +932,15 @@ def _retention_lock(runs_root: Path) -> Iterator[None]: pass restrict_file(lock_path) stream = lock_path.open("a+b") + locked = False try: _lock_retention_stream(stream) + locked = True yield finally: try: - _unlock_retention_stream(stream) + if locked: + _unlock_retention_stream(stream) finally: stream.close() @@ -935,10 +948,15 @@ def _retention_lock(runs_root: Path) -> Iterator[None]: def _lock_retention_stream(stream: object) -> None: fd = stream.fileno() # type: ignore[attr-defined] if _fcntl is not None: - _fcntl.flock(fd, _fcntl.LOCK_EX) + _fcntl.flock(fd, _fcntl.LOCK_EX | _fcntl.LOCK_NB) elif _msvcrt is not None: # pragma: no cover - Windows stream.seek(0) - _msvcrt.locking(fd, _msvcrt.LK_LOCK, 1) + try: + _msvcrt.locking(fd, _msvcrt.LK_NBLCK, 1) + except OSError as exc: + if exc.errno in {11, 13, 36}: + raise BlockingIOError("retention lock is busy") from exc + raise def _unlock_retention_stream(stream: object) -> None: diff --git a/tests/test_retention_contention.py b/tests/test_retention_contention.py new file mode 100644 index 0000000..fecb0c7 --- /dev/null +++ b/tests/test_retention_contention.py @@ -0,0 +1,93 @@ +from __future__ import annotations + +import logging +import os +import subprocess +import sys +from pathlib import Path +from unittest.mock import patch + +from base_cli import _runtime as runtime + + +def test_contended_retention_returns_without_waiting(tmp_path: Path) -> None: + # A separate process faithfully models advisory lock contention on both OSes. + code = """ +import sys +from pathlib import Path +from base_cli._runtime import _retention_lock +with _retention_lock(Path(sys.argv[1])): + print('locked', flush=True) + sys.stdin.readline() +""" + env = {**os.environ, "PYTHONPATH": str(Path(runtime.__file__).resolve().parents[1])} + child = subprocess.Popen( + [sys.executable, "-c", code, str(tmp_path)], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + env=env, + ) + try: + assert child.stdout.readline().strip() == "locked" + probe = subprocess.run( + [ + sys.executable, + "-c", + """ +import sys +from pathlib import Path +from base_cli._runtime import prune_run_bundles, refresh_run_bundle_index +root = Path(sys.argv[1]) +prune_run_bundles(root, max_bundles=1) +refresh_run_bundle_index(root) +""", + str(tmp_path), + ], + env=env, + capture_output=True, + text=True, + timeout=5, + ) + assert probe.returncode == 0, probe.stderr + finally: + child.communicate("done\n", timeout=5) + + +def test_contended_index_refresh_is_eventually_consistent(tmp_path: Path) -> None: + logger = logging.getLogger("retention-index-refresh") + with patch.object(runtime, "_retention_lock", side_effect=BlockingIOError("busy")): + runtime.refresh_run_bundle_index(tmp_path, logger=logger) + + +def test_concurrent_passes_converge_after_a_serial_pass(tmp_path: Path) -> None: + for index in range(8): + bundle = tmp_path / f"run-{index:02d}" + bundle.mkdir() + (bundle / "run.json").write_text( + f'{{"run_id": "run-{index:02d}", "status": "ok", "started_at": "2020-01-01T00:00:00Z"}}', + encoding="utf-8", + ) + code = """ +import sys +from pathlib import Path +from base_cli._runtime import prune_run_bundles +prune_run_bundles(Path(sys.argv[1]), max_bundles=3) +""" + env = {**os.environ, "PYTHONPATH": str(Path(runtime.__file__).resolve().parents[1])} + children = [ + subprocess.Popen( + [sys.executable, "-c", code, str(tmp_path)], + env=env, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + for _ in range(8) + ] + for child in children: + _stdout, stderr = child.communicate(timeout=10) + assert child.returncode == 0, stderr.decode() + + runtime.prune_run_bundles(tmp_path, max_bundles=3) + assert len([path for path in tmp_path.iterdir() if path.is_dir() and path.name.startswith("run-")]) <= 3