From 739dada0151eb6e787bd18f2fe2c2b6f9b534db2 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Fri, 10 Jul 2026 02:53:29 -0700 Subject: [PATCH] =?UTF-8?q?#671=20prax:=20single-instance=20lock=20on=20th?= =?UTF-8?q?e=20monitor=20=E2=80=94=20stop=20duplicate/orphan=20monitors=20?= =?UTF-8?q?from=20double-sending=20Telegram=20relay.=20New=20instance=5Flo?= =?UTF-8?q?ck=20handler=20writes=20a=20liveness-checked=20pidfile=20(prax?= =?UTF-8?q?=5Fjson/monitor.pid,=20outside=20the=20tailed=20system=5Flogs/)?= =?UTF-8?q?:=20acquire()=20before=20relay=20init=20refuses=20to=20start=20?= =?UTF-8?q?(fail-loud,=20names=20the=20holding=20PID)=20when=20a=20live=20?= =?UTF-8?q?monitor=20holds=20the=20lock,=20reclaims=20a=20stale=20pidfile?= =?UTF-8?q?=20on=20a=20dead=20PID,=20release()=20clears=20it=20on=20shutdo?= =?UTF-8?q?wn.=20Liveness=20probe=20platform-branched=20=E2=80=94=20POSIX?= =?UTF-8?q?=20os.kill(pid,0),=20Windows=20OpenProcess/GetExitCodeProcess?= =?UTF-8?q?=20(a=20raw=20os.kill(pid,0)=20TERMINATES=20the=20target=20on?= =?UTF-8?q?=20Windows;=20reused=20the=20canonical=20devpulse/watchdog/agen?= =?UTF-8?q?t.py:137=20impl).=20monitor.py=20split=20under=20the=20600-line?= =?UTF-8?q?=20limit=20(pid=5Fcache=20extracted).=20Built=20by=20@prax,=20v?= =?UTF-8?q?erified=20by=20devpulse:=20120=20tests=20green=20across=20the?= =?UTF-8?q?=204=20touched=20files=20(+25=20new=20incl=203=20Windows-path),?= =?UTF-8?q?=20seedgo=2031/31=20on=20all=203=20sources.=20Verify=20caught?= =?UTF-8?q?=20the=20Windows=20os.kill=20hazard=20on=20the=20first=20pass;?= =?UTF-8?q?=20filed=20the=20seedgo=20windows=5Fcompat=20detector=20gap=20a?= =?UTF-8?q?s=20#682.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 12 + .../apps/handlers/monitoring/instance_lock.py | 127 +++++++++++ .../apps/handlers/monitoring/pid_cache.py | 91 ++++++++ src/aipass/prax/apps/modules/monitor.py | 193 ++++------------ src/aipass/prax/tests/test_instance_lock.py | 207 ++++++++++++++++++ src/aipass/prax/tests/test_monitor_module.py | 194 +--------------- src/aipass/prax/tests/test_pid_cache.py | 188 ++++++++++++++++ src/aipass/prax/tests/test_telegram_relay.py | 2 + 8 files changed, 675 insertions(+), 339 deletions(-) create mode 100644 src/aipass/prax/apps/handlers/monitoring/instance_lock.py create mode 100644 src/aipass/prax/apps/handlers/monitoring/pid_cache.py create mode 100644 src/aipass/prax/tests/test_instance_lock.py create mode 100644 src/aipass/prax/tests/test_pid_cache.py diff --git a/CHANGELOG.md b/CHANGELOG.md index f1ead3ee..f172e747 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -51,6 +51,18 @@ PyPI version — not the changelog header. Rerouted to a dim console note; genuine argument errors still `error()` → exit 2. (devpulse) +- **The prax monitor now holds a single-instance lock, so a duplicate/orphan + monitor can't double-send Telegram relay messages (issue #671).** A new + `instance_lock` handler writes a pidfile (`prax_json/monitor.pid`, outside the + tailed `system_logs/`) with a liveness check: `acquire()` runs before relay + init and refuses to start (fail-loud, naming the holding PID) if a live monitor + already holds the lock, reclaims a stale pidfile when the recorded PID is dead, + and `release()` clears it on shutdown. The liveness probe is platform-branched + — POSIX `os.kill(pid, 0)`, Windows `OpenProcess`/`GetExitCodeProcess` (a raw + `os.kill(pid, 0)` *terminates* the target on Windows). `monitor.py` was also + split under the 600-line limit (`pid_cache` extracted). +25 tests. + (built by @prax, verified by devpulse) + --- ## [2026-07-09] diff --git a/src/aipass/prax/apps/handlers/monitoring/instance_lock.py b/src/aipass/prax/apps/handlers/monitoring/instance_lock.py new file mode 100644 index 00000000..7c66ddd8 --- /dev/null +++ b/src/aipass/prax/apps/handlers/monitoring/instance_lock.py @@ -0,0 +1,127 @@ +# =================== AIPass ==================== +# Name: instance_lock.py +# Description: Single-instance lock for the prax monitor +# Version: 1.0.0 +# Created: 2026-07-10 +# Modified: 2026-07-10 +# ============================================= + +"""Single-instance lock for the prax monitor. + +Prevents duplicate monitor processes from running concurrently (and +double-sending Telegram relay messages). Uses a pidfile with liveness +check — cross-platform (Linux / macOS / Windows). + +Lock file lives in prax_json/monitor.pid (outside system_logs/ to avoid +the tailed-directory feedback loop). +""" + +import json as _json +import os +import sys +from pathlib import Path +from typing import Optional + +from aipass.prax.apps.modules.logger import get_direct_logger +from aipass.prax.apps.handlers.json import json_handler + +logger = get_direct_logger() + +_lock_path_override: Optional[Path] = None +_held_lock: Optional[Path] = None + + +def _pid_alive_windows(pid: int) -> bool: + """Windows-safe liveness check via OpenProcess + GetExitCodeProcess.""" + import ctypes + from ctypes import wintypes + + PROCESS_QUERY_LIMITED_INFORMATION = 0x1000 + STILL_ACTIVE = 259 + + kernel32 = ctypes.windll.kernel32 # type: ignore[attr-defined] + kernel32.OpenProcess.argtypes = [wintypes.DWORD, wintypes.BOOL, wintypes.DWORD] + kernel32.OpenProcess.restype = wintypes.HANDLE + kernel32.GetExitCodeProcess.argtypes = [wintypes.HANDLE, ctypes.POINTER(wintypes.DWORD)] + kernel32.GetExitCodeProcess.restype = wintypes.BOOL + kernel32.CloseHandle.argtypes = [wintypes.HANDLE] + kernel32.CloseHandle.restype = wintypes.BOOL + + handle = kernel32.OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, False, pid) + if not handle: + return False + try: + exit_code = wintypes.DWORD() + if not kernel32.GetExitCodeProcess(handle, ctypes.byref(exit_code)): + return False + return exit_code.value == STILL_ACTIVE + finally: + kernel32.CloseHandle(handle) + + +def _is_pid_alive(pid: int) -> bool: + """Check if a process with the given PID is alive. Cross-platform.""" + if sys.platform == "win32": + try: + return _pid_alive_windows(pid) + except Exception as exc: + logger.info("[instance_lock] PID %d Windows check failed (assuming alive): %s", pid, exc) + return True + try: + os.kill(pid, 0) + return True + except ProcessLookupError: + logger.info("[instance_lock] PID %d not found", pid) + return False + except PermissionError: + logger.info("[instance_lock] PID %d alive (permission denied on signal)", pid) + return True + except OSError as exc: + logger.info("[instance_lock] os.kill(%d, 0) raised %s", pid, exc) + return False + + +def get_lock_path() -> Path: + """Return the path for the monitor single-instance lock file.""" + if _lock_path_override is not None: + return _lock_path_override + return Path(__file__).resolve().parent.parent.parent / "prax_json" / "monitor.pid" + + +def acquire(error_fn=None) -> None: + """Acquire single-instance lock. Raises SystemExit(1) if another live instance holds it.""" + global _held_lock + lock_path = get_lock_path() + json_handler.log_operation("instance_lock_acquire", {"pid": os.getpid()}) + + if lock_path.exists(): + try: + data = _json.loads(lock_path.read_text(encoding="utf-8")) + existing_pid = data.get("pid", 0) + if existing_pid and _is_pid_alive(existing_pid): + msg = f"Monitor already running (PID {existing_pid}). Kill the existing process or remove {lock_path}" + if error_fn: + error_fn(msg) + logger.error("[instance_lock] %s", msg) + raise SystemExit(1) + logger.info("[instance_lock] Reclaiming stale lock (PID %d is dead)", existing_pid) + except (ValueError, OSError) as exc: + logger.info("[instance_lock] Removing corrupt lock file: %s", exc) + + lock_path.parent.mkdir(parents=True, exist_ok=True) + lock_path.write_text(_json.dumps({"pid": os.getpid()}), encoding="utf-8") + _held_lock = lock_path + logger.info("[instance_lock] Acquired (PID %d)", os.getpid()) + + +def release() -> None: + """Release the single-instance lock file.""" + global _held_lock + if _held_lock and _held_lock.exists(): + try: + _held_lock.unlink() + logger.info("[instance_lock] Released") + except OSError as exc: + logger.warning("[instance_lock] Failed to remove lock file: %s", exc) + _held_lock = None + json_handler.log_operation("instance_lock_release", {"pid": os.getpid()}) diff --git a/src/aipass/prax/apps/handlers/monitoring/pid_cache.py b/src/aipass/prax/apps/handlers/monitoring/pid_cache.py new file mode 100644 index 00000000..89a26dde --- /dev/null +++ b/src/aipass/prax/apps/handlers/monitoring/pid_cache.py @@ -0,0 +1,91 @@ +# =================== AIPass ==================== +# Name: pid_cache.py +# Description: PID cache for branch-to-agent mapping +# Version: 1.0.0 +# Created: 2026-07-10 +# Modified: 2026-07-10 +# ============================================= + +"""PID cache — maps branch names to active agent PIDs from dispatch lock files. + +Scans .dispatch.lock files in each branch's ai_mail.local/ directory, +verifies the PID is alive via /proc, and caches the mapping with a TTL. +Used by the monitor to attribute events to the owning agent process. +""" + +import json as _json +import sys +import threading +import time as _time +from pathlib import Path +from typing import Optional + +from aipass.prax.apps.modules.logger import get_direct_logger +from aipass.prax.apps.handlers.json import json_handler + +logger = get_direct_logger() + +_pid_cache: dict[str, int] = {} +_pid_cache_lock = threading.Lock() +_pid_cache_last_refresh: float = 0.0 +_PID_CACHE_TTL = 30.0 + + +def parse_lock_pid(branch_entry: dict, new_cache: dict[str, int]) -> None: + """Parse a single dispatch lock file and add to cache if PID is live.""" + branch_path = Path(branch_entry.get("path", "")) + lock_path = branch_path / "ai_mail.local" / ".dispatch.lock" + if not lock_path.exists(): + return + try: + lock_data = _json.loads(lock_path.read_text(encoding="utf-8")) + pid = lock_data.get("pid", 0) + if not pid or not (sys.platform == "linux" and Path(f"/proc/{pid}").exists()): + return + name = branch_entry.get("name", "").upper() + if name: + new_cache[name] = pid + except (ValueError, OSError) as e: + logger.info("[pid_cache] Skipping dispatch lock %s: %s", lock_path, e) + + +def refresh(repo_root: Optional[Path] = None) -> None: + """Scan dispatch lock files to build branch-to-PID mapping. + + Args: + repo_root: Repository root path. When None, walks up from this file. + """ + global _pid_cache_last_refresh + + now = _time.time() + with _pid_cache_lock: + if now - _pid_cache_last_refresh < _PID_CACHE_TTL: + return + _pid_cache_last_refresh = now + + try: + if repo_root is None: + repo_root = Path(__file__).resolve().parent.parent.parent.parent + registry_path = repo_root / "AIPASS_REGISTRY.json" + if not registry_path.exists(): + return + data = _json.loads(registry_path.read_text(encoding="utf-8")) + new_cache: dict[str, int] = {} + for branch in data.get("branches", []): + parse_lock_pid(branch, new_cache) + with _pid_cache_lock: + _pid_cache.clear() + _pid_cache.update(new_cache) + json_handler.log_operation("pid_cache_refresh", {"count": len(new_cache)}) + except Exception as e: + logger.info("[pid_cache] Refresh failed: %s", e) + + +def get_pid_for_branch(branch: str) -> Optional[int]: + """Look up PID for a branch from the cache.""" + refresh() + base = branch.upper() + if base.endswith(" AGENT"): + base = base[:-6] + with _pid_cache_lock: + return _pid_cache.get(base) diff --git a/src/aipass/prax/apps/modules/monitor.py b/src/aipass/prax/apps/modules/monitor.py index d261647b..f70c27af 100755 --- a/src/aipass/prax/apps/modules/monitor.py +++ b/src/aipass/prax/apps/modules/monitor.py @@ -6,17 +6,7 @@ # Modified: 2026-03-09 # ============================================= -""" -PRAX Monitor Module - Mission Control for Autonomous Branches - -Thin orchestration layer for real-time monitoring of file changes, log events, -and agent activity across all AIPass branches. Delegates to handlers in -apps/handlers/monitoring/ (unified_stream, branch_detector, event_queue, etc.) - -Usage: - drone @prax monitor # Show introspection - drone @prax monitor run # Monitor all branches -""" +"""PRAX Monitor Module - Mission Control for Autonomous Branches.""" import os import sys @@ -57,75 +47,11 @@ from aipass.prax.apps.handlers.monitoring.telegram_relay import ( stop_relay, is_relay_enabled_by_env, ) - - -# ============================================================================= -# PID CACHE - Maps branch names to active agent PIDs from dispatch lock files -# ============================================================================= +from aipass.prax.apps.handlers.monitoring.pid_cache import get_pid_for_branch as _get_pid_for_branch +from aipass.prax.apps.handlers.monitoring import instance_lock import json as _json -_pid_cache: dict[str, int] = {} -_pid_cache_lock = threading.Lock() -_pid_cache_last_refresh: float = 0.0 -_PID_CACHE_TTL = 30.0 # Refresh every 30 seconds - - -def _parse_lock_pid(branch_entry: dict, new_cache: dict[str, int]) -> None: - """Parse a single dispatch lock file and add to cache if PID is live.""" - branch_path = Path(branch_entry.get("path", "")) - lock_path = branch_path / "ai_mail.local" / ".dispatch.lock" - if not lock_path.exists(): - return - try: - lock_data = _json.loads(lock_path.read_text(encoding="utf-8")) - pid = lock_data.get("pid", 0) - if not pid or not (sys.platform == "linux" and Path(f"/proc/{pid}").exists()): - return - name = branch_entry.get("name", "").upper() - if name: - new_cache[name] = pid - except (ValueError, OSError) as e: - logger.info("[monitor] Skipping dispatch lock %s: %s", lock_path, e) - - -def _refresh_pid_cache() -> None: - """Scan dispatch lock files to build branch→PID mapping.""" - global _pid_cache_last_refresh - import time as _time - - now = _time.time() - with _pid_cache_lock: - if now - _pid_cache_last_refresh < _PID_CACHE_TTL: - return - _pid_cache_last_refresh = now - - try: - from aipass.prax.apps.handlers.config.load import _find_repo_root - - registry_path = _find_repo_root() / "AIPASS_REGISTRY.json" - if not registry_path.exists(): - return - data = _json.loads(registry_path.read_text(encoding="utf-8")) - new_cache: dict[str, int] = {} - for branch in data.get("branches", []): - _parse_lock_pid(branch, new_cache) - with _pid_cache_lock: - _pid_cache.clear() - _pid_cache.update(new_cache) - except Exception as e: - logger.info(f"[monitor] PID cache refresh failed: {e}") - - -def _get_pid_for_branch(branch: str) -> Optional[int]: - """Look up PID for a branch from the cache.""" - _refresh_pid_cache() - base = branch.upper() - if base.endswith(" AGENT"): - base = base[:-6] - with _pid_cache_lock: - return _pid_cache.get(base) - # ============================================================================= # MODULE STATE @@ -143,6 +69,14 @@ _log_watcher_thread: Optional[threading.Thread] = None def print_introspection(): """Display module introspection - shows connected handlers and architecture.""" json_handler.log_operation("print_introspection", {"module": "monitor"}) + _handlers = [ + ("1. unified_stream.py", "print_event() - Terminal output formatting"), + ("2. branch_detector.py", "detect_branch_from_path() - Path-to-branch mapping"), + ("3. interactive_filter.py", "FilterState, parse_command() - Runtime filtering"), + ("4. monitoring_filters.py", "should_monitor(), get_priority() - Event filtering"), + ("5. event_queue.py", "MonitoringEvent, MonitoringQueue - Event buffering"), + ("6. module_tracker.py", "ModuleTracker - Module execution tracking"), + ] console.print() console.print("[bold cyan]PRAX Monitor Module[/bold cyan]") console.print() @@ -151,79 +85,49 @@ def print_introspection(): console.print(" Unified console for file changes, logs, and module activity") console.print() console.print("[yellow]Connected Handlers (apps/handlers/monitoring/):[/yellow]") - console.print() - console.print(" [cyan]1. unified_stream.py[/cyan]") - console.print(" [dim]→ print_event() - Terminal output formatting[/dim]") - console.print() - console.print(" [cyan]2. branch_detector.py[/cyan]") - console.print(" [dim]→ detect_branch_from_path() - Path-to-branch mapping[/dim]") - console.print() - console.print(" [cyan]3. interactive_filter.py[/cyan]") - console.print(" [dim]→ FilterState, parse_command() - Runtime filtering[/dim]") - console.print() - console.print(" [cyan]4. monitoring_filters.py[/cyan]") - console.print(" [dim]→ should_monitor(), get_priority() - Event filtering[/dim]") - console.print() - console.print(" [cyan]5. event_queue.py[/cyan]") - console.print(" [dim]→ MonitoringEvent, MonitoringQueue - Event buffering[/dim]") - console.print() - console.print(" [cyan]6. module_tracker.py[/cyan]") - console.print(" [dim]→ ModuleTracker - Module execution tracking[/dim]") - console.print() - console.print(" [cyan]7. file watcher (threaded)[/cyan]") - console.print(" [dim]→ Real-time file change detection using watchdog[/dim]") + for name, desc in _handlers: + console.print(f"\n [cyan]{name}[/cyan]\n [dim]{desc}[/dim]") + console.print("\n [cyan]7. file watcher (threaded)[/cyan]") + console.print(" [dim]Real-time file change detection using watchdog[/dim]") console.print(" [green]STATUS: Active - monitors ECOSYSTEM_ROOT recursively[/green]") - console.print() - console.print(" [cyan]8. log monitor (threaded)[/cyan]") - console.print(" [dim]→ Log stream processing from SYSTEM_LOGS_DIR[/dim]") + console.print("\n [cyan]8. log monitor (threaded)[/cyan]") + console.print(" [dim]Log stream processing from SYSTEM_LOGS_DIR[/dim]") console.print(" [green]STATUS: Active - watches *.log files for new entries[/green]") - console.print() - console.print("[dim]Run 'drone @prax monitor --help' for usage[/dim]") - console.print() + console.print("\n[dim]Run 'drone @prax monitor --help' for usage[/dim]\n") def print_help(): """Drone-compliant help output - command syntax and examples.""" console.print() console.print("[bold cyan]PRAX Monitor - Unified Branch Monitoring[/bold cyan]") - console.print() - console.print("[yellow]Commands:[/yellow]") - console.print() - console.print(" [cyan]drone @prax monitor[/cyan]") - console.print(" Show module introspection") - console.print() - console.print(" [cyan]drone @prax monitor run[/cyan]") - console.print(" Start monitoring all branches") - console.print() - console.print(" [cyan]drone @prax monitor run all[/cyan]") - console.print(" Explicit all-branches monitoring") - console.print() - console.print(" [cyan]drone @prax monitor run [branches][/cyan]") - console.print(" Monitor specific branches (comma-separated)") - console.print(" Example: drone @prax monitor run seedgo,cli,flow") - console.print() - console.print(" [cyan]drone @prax monitor run --relay[/cyan]") - console.print(" Enable Telegram relay (mirrors feed to prax_monitor bot)") - console.print(" Also enabled by env AIPASS_PRAX_MONITOR_RELAY=1") - console.print() - console.print(" [cyan]drone @prax monitor --help[/cyan]") - console.print(" Show this help") - console.print() - console.print("[yellow]Interactive Mode Commands:[/yellow]") - console.print() + _cmds = [ + ("drone @prax monitor", "Show module introspection"), + ("drone @prax monitor run", "Start monitoring all branches"), + ("drone @prax monitor run all", "Explicit all-branches monitoring"), + ( + "drone @prax monitor run [branches]", + "Monitor specific branches (comma-separated)\n Example: drone @prax monitor run seedgo,cli,flow", + ), + ( + "drone @prax monitor run --relay", + "Enable Telegram relay (mirrors feed to prax_monitor bot)" + "\n Also enabled by env AIPASS_PRAX_MONITOR_RELAY=1", + ), + ("drone @prax monitor --help", "Show this help"), + ] + console.print("\n[yellow]Commands:[/yellow]") + for cmd, desc in _cmds: + console.print(f"\n [cyan]{cmd}[/cyan]\n {desc}") + console.print("\n[yellow]Interactive Mode Commands:[/yellow]") console.print(" [cyan]help[/cyan] Show available commands") console.print(" [cyan]status[/cyan] Display current monitoring state") console.print(" [cyan]filter [branches][/cyan] Adjust branch filter") console.print(" [cyan]quit/exit[/cyan] Stop monitoring") - console.print() - console.print("[yellow]Examples:[/yellow]") - console.print() - console.print(" [dim]# Monitor all branches[/dim]") + console.print("\n[yellow]Examples:[/yellow]") + console.print("\n [dim]# Monitor all branches[/dim]") console.print(" $ drone @prax monitor run") - console.print() - console.print(" [dim]# Monitor specific branches[/dim]") - console.print(" $ drone @prax monitor run seedgo,cli,flow") - console.print() + console.print("\n [dim]# Monitor specific branches[/dim]") + console.print(" $ drone @prax monitor run seedgo,cli,flow\n") # ============================================================================= @@ -232,17 +136,7 @@ def print_help(): def handle_command(command: str, args: List[str]) -> bool: - """ - Handle monitor command - required for auto-discovery by prax.py - - Args: - command: Command name from prax.py dispatcher - args: Command arguments (branch filters, flags, etc.) - - Returns: - True if command was handled (command == "monitor") - False if not our command (pass to next handler) - """ + """Handle monitor command - required for auto-discovery by prax.py.""" if command != "monitor": return False @@ -283,6 +177,8 @@ def _run_monitor(args: List[str]) -> bool: global _event_queue, _module_tracker global _display_thread, _file_watcher_thread, _log_watcher_thread + instance_lock.acquire(error_fn=error) + json_handler.log_operation("monitor_started", {"args": args}) logger.info(f"Starting unified monitoring (args: {args})") @@ -362,6 +258,7 @@ def _stop_threads(): if t is not None and t.is_alive(): t.join(timeout=2.0) + instance_lock.release() logger.info("All monitoring threads stopped") diff --git a/src/aipass/prax/tests/test_instance_lock.py b/src/aipass/prax/tests/test_instance_lock.py new file mode 100644 index 00000000..77ce1de3 --- /dev/null +++ b/src/aipass/prax/tests/test_instance_lock.py @@ -0,0 +1,207 @@ +# =================== AIPass ==================== +# Name: test_instance_lock.py +# Description: Tests for the monitor single-instance lock +# Version: 1.0.0 +# Created: 2026-07-10 +# Modified: 2026-07-10 +# ============================================= + +"""Tests for apps/handlers/monitoring/instance_lock.py + +Covers: +- _is_pid_alive() cross-platform liveness check +- acquire() creates lock, refuses live duplicate, reclaims stale +- release() removes lock file on clean shutdown +""" + +import json +import os +import sys +from unittest.mock import MagicMock, patch + +_HANDLER_MOCKS = { + "aipass.prax.apps.handlers.json": MagicMock(), + "aipass.prax.apps.handlers.json.json_handler": MagicMock(), +} + + +def _import_lock(): + """Import (or reload) instance_lock with handler mocks.""" + fresh = {k: MagicMock() for k in _HANDLER_MOCKS} + with patch.dict(sys.modules, fresh): + import importlib + + if "aipass.prax.apps.handlers.monitoring.instance_lock" in sys.modules: + mod = importlib.reload(sys.modules["aipass.prax.apps.handlers.monitoring.instance_lock"]) + else: + mod = importlib.import_module("aipass.prax.apps.handlers.monitoring.instance_lock") + return mod + + +class TestIsPidAlive: + """Test cross-platform PID liveness check.""" + + def test_live_pid_returns_true_posix(self): + """os.kill(pid, 0) success means alive on POSIX.""" + mod = _import_lock() + with patch("sys.platform", "linux"), patch("os.kill"): + assert mod._is_pid_alive(os.getpid()) is True + + def test_dead_pid_returns_false(self): + """Non-existent PID returns False on POSIX.""" + mod = _import_lock() + with patch("sys.platform", "linux"), patch("os.kill", side_effect=ProcessLookupError): + assert mod._is_pid_alive(99999999) is False + + def test_permission_error_means_alive(self): + """PermissionError means the process exists but is owned by another user.""" + mod = _import_lock() + with patch("sys.platform", "linux"), patch("os.kill", side_effect=PermissionError): + assert mod._is_pid_alive(1) is True + + def test_generic_oserror_returns_false(self): + """Other OSError returns False.""" + mod = _import_lock() + with patch("sys.platform", "linux"), patch("os.kill", side_effect=OSError(99, "Unknown")): + assert mod._is_pid_alive(12345) is False + + def test_windows_delegates_to_pid_alive_windows(self): + """On win32, _is_pid_alive delegates to _pid_alive_windows.""" + mod = _import_lock() + with ( + patch("sys.platform", "win32"), + patch.object(mod, "_pid_alive_windows", return_value=True) as mock_win, + ): + assert mod._is_pid_alive(1234) is True + mock_win.assert_called_once_with(1234) + + def test_windows_dead_pid(self): + """On win32, dead PID returns False via _pid_alive_windows.""" + mod = _import_lock() + with ( + patch("sys.platform", "win32"), + patch.object(mod, "_pid_alive_windows", return_value=False), + ): + assert mod._is_pid_alive(99999999) is False + + def test_windows_ctypes_failure_assumes_alive(self): + """On win32, if ctypes fails, assume the process is alive (safe default).""" + mod = _import_lock() + with ( + patch("sys.platform", "win32"), + patch.object(mod, "_pid_alive_windows", side_effect=OSError("ctypes failed")), + ): + assert mod._is_pid_alive(1234) is True + + +class TestAcquire: + """Test single-instance lock acquisition.""" + + def test_creates_lock_file(self, tmp_path): + """acquire() creates a lock file with the current PID.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + mod.acquire() + + assert lock_path.exists() + data = json.loads(lock_path.read_text(encoding="utf-8")) + assert data["pid"] == os.getpid() + + def test_refuses_when_live_instance_holds_lock(self, tmp_path): + """acquire() exits with SystemExit(1) when another live process holds the lock.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + lock_path.write_text(json.dumps({"pid": os.getpid()}), encoding="utf-8") + + import pytest + + mock_error = MagicMock() + with pytest.raises(SystemExit) as exc_info: + mod.acquire(error_fn=mock_error) + assert exc_info.value.code == 1 + mock_error.assert_called_once() + assert str(os.getpid()) in mock_error.call_args[0][0] + + def test_reclaims_stale_lock(self, tmp_path): + """acquire() reclaims the lock when the recorded PID is dead.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + lock_path.write_text(json.dumps({"pid": 99999999}), encoding="utf-8") + + with patch.object(mod, "_is_pid_alive", return_value=False): + mod.acquire() + + data = json.loads(lock_path.read_text(encoding="utf-8")) + assert data["pid"] == os.getpid() + + def test_reclaims_corrupt_lock_file(self, tmp_path): + """acquire() overwrites a corrupt lock file.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + lock_path.write_text("{corrupt json", encoding="utf-8") + + mod.acquire() + + data = json.loads(lock_path.read_text(encoding="utf-8")) + assert data["pid"] == os.getpid() + + def test_creates_parent_directories(self, tmp_path): + """acquire() creates parent directories if they don't exist.""" + mod = _import_lock() + lock_path = tmp_path / "nested" / "dir" / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + mod.acquire() + + assert lock_path.exists() + + +class TestRelease: + """Test single-instance lock release.""" + + def test_removes_lock_file(self, tmp_path): + """release() removes the lock file.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + mod.acquire() + assert lock_path.exists() + + mod.release() + assert not lock_path.exists() + + def test_clears_held_lock_state(self, tmp_path): + """release() clears the _held_lock global.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + mod.acquire() + mod.release() + assert mod._held_lock is None + + def test_release_without_acquire_is_safe(self): + """release() is a no-op when no lock is held.""" + mod = _import_lock() + setattr(mod, "_held_lock", None) + mod.release() + + def test_release_handles_already_deleted_file(self, tmp_path): + """release() handles the case where the lock file was already deleted.""" + mod = _import_lock() + lock_path = tmp_path / "monitor.pid" + setattr(mod, "_lock_path_override", lock_path) + + mod.acquire() + lock_path.unlink() + mod.release() + assert mod._held_lock is None diff --git a/src/aipass/prax/tests/test_monitor_module.py b/src/aipass/prax/tests/test_monitor_module.py index 31b7332f..7507759f 100644 --- a/src/aipass/prax/tests/test_monitor_module.py +++ b/src/aipass/prax/tests/test_monitor_module.py @@ -41,6 +41,8 @@ _MONITORING_MOCKS = { "aipass.prax.apps.handlers.monitoring.monitoring_filters": MagicMock(), "aipass.prax.apps.handlers.monitoring.file_watcher_integration": MagicMock(), "aipass.prax.apps.handlers.monitoring.telegram_relay": MagicMock(), + "aipass.prax.apps.handlers.monitoring.pid_cache": MagicMock(), + "aipass.prax.apps.handlers.monitoring.instance_lock": MagicMock(), } @@ -209,196 +211,6 @@ class TestGetWatchDirectories: assert trinity_dir in paths -# --------------------------------------------------------------------------- -# _parse_lock_pid tests (lines 61-74) -# --------------------------------------------------------------------------- - - -class TestParseLockPid: - """Test dispatch lock file parsing for PID cache.""" - - def test_no_lock_file_does_nothing(self, tmp_path): - """Branch entry without a lock file adds nothing to cache.""" - mod = _import_monitor() - new_cache: dict[str, int] = {} - entry = {"path": str(tmp_path / "somebranch"), "name": "flow"} - mod._parse_lock_pid(entry, new_cache) - assert new_cache == {} - - def test_lock_file_with_live_pid_on_linux(self, tmp_path): - """Lock file with a PID that has a /proc entry adds to cache.""" - mod = _import_monitor() - branch_dir = tmp_path / "mybranch" - mail_dir = branch_dir / "ai_mail.local" - mail_dir.mkdir(parents=True) - lock_data = {"pid": 12345} - (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") - - new_cache: dict[str, int] = {} - entry = {"path": str(branch_dir), "name": "flow"} - - # Mock /proc/12345 existence check - with ( - patch("sys.platform", "linux"), - patch("pathlib.Path.exists", side_effect=lambda self=None: True), - ): - mod._parse_lock_pid(entry, new_cache) - - assert new_cache.get("FLOW") == 12345 - - def test_lock_file_with_zero_pid(self, tmp_path): - """Lock file with pid=0 skips entry.""" - mod = _import_monitor() - branch_dir = tmp_path / "mybranch" - mail_dir = branch_dir / "ai_mail.local" - mail_dir.mkdir(parents=True) - lock_data = {"pid": 0} - (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") - - new_cache: dict[str, int] = {} - entry = {"path": str(branch_dir), "name": "flow"} - mod._parse_lock_pid(entry, new_cache) - assert new_cache == {} - - def test_lock_file_with_invalid_json(self, tmp_path): - """Lock file with invalid JSON logs warning and continues.""" - mod = _import_monitor() - branch_dir = tmp_path / "mybranch" - mail_dir = branch_dir / "ai_mail.local" - mail_dir.mkdir(parents=True) - (mail_dir / ".dispatch.lock").write_text("{bad json}", encoding="utf-8") - - new_cache: dict[str, int] = {} - entry = {"path": str(branch_dir), "name": "flow"} - mod._parse_lock_pid(entry, new_cache) - assert new_cache == {} - - def test_lock_file_with_empty_name(self, tmp_path): - """Branch entry with empty name skips cache update.""" - mod = _import_monitor() - branch_dir = tmp_path / "mybranch" - mail_dir = branch_dir / "ai_mail.local" - mail_dir.mkdir(parents=True) - lock_data = {"pid": 99999} - (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") - - new_cache: dict[str, int] = {} - entry = {"path": str(branch_dir), "name": ""} - - with ( - patch("sys.platform", "linux"), - patch("pathlib.Path.exists", return_value=True), - ): - mod._parse_lock_pid(entry, new_cache) - assert new_cache == {} - - -# --------------------------------------------------------------------------- -# _refresh_pid_cache tests (lines 80-102) -# --------------------------------------------------------------------------- - - -class TestRefreshPidCache: - """Test PID cache refresh from registry.""" - - def test_skips_when_within_ttl(self): - """Cache refresh is skipped if within TTL window.""" - mod = _import_monitor() - import time - - # Simulate recent refresh - with mod._pid_cache_lock: - setattr(mod, "_pid_cache_last_refresh", time.time()) - - with patch.object(mod, "_parse_lock_pid") as mock_parse: - mod._refresh_pid_cache() - mock_parse.assert_not_called() - - def test_refreshes_when_ttl_expired(self, tmp_path): - """Cache refresh runs when TTL has expired.""" - mod = _import_monitor() - with mod._pid_cache_lock: - setattr(mod, "_pid_cache_last_refresh", 0.0) - - registry_data = {"branches": [{"name": "flow", "path": str(tmp_path / "flow")}]} - registry_file = tmp_path / "AIPASS_REGISTRY.json" - registry_file.write_text(json.dumps(registry_data), encoding="utf-8") - - mock_find_root = MagicMock(return_value=tmp_path) - with patch.dict( - sys.modules, - {"aipass.prax.apps.handlers.config.load": MagicMock(_find_repo_root=mock_find_root)}, - ): - mod._refresh_pid_cache() - - def test_handles_missing_registry(self, tmp_path): - """Missing registry file does not crash.""" - mod = _import_monitor() - with mod._pid_cache_lock: - setattr(mod, "_pid_cache_last_refresh", 0.0) - - mock_find_root = MagicMock(return_value=tmp_path) - with patch.dict( - sys.modules, - {"aipass.prax.apps.handlers.config.load": MagicMock(_find_repo_root=mock_find_root)}, - ): - mod._refresh_pid_cache() - - def test_handles_exception_in_refresh(self): - """Exception during refresh is caught and logged.""" - mod = _import_monitor() - with mod._pid_cache_lock: - setattr(mod, "_pid_cache_last_refresh", 0.0) - - mock_load = MagicMock() - mock_load._find_repo_root.side_effect = RuntimeError("boom") - with patch.dict( - sys.modules, - {"aipass.prax.apps.handlers.config.load": mock_load}, - ): - # Should not raise - mod._refresh_pid_cache() - - -# --------------------------------------------------------------------------- -# _get_pid_for_branch tests (lines 107-112) -# --------------------------------------------------------------------------- - - -class TestGetPidForBranch: - """Test PID lookup for branch names.""" - - def test_returns_pid_from_cache(self): - """Returns PID when branch is in cache.""" - mod = _import_monitor() - with mod._pid_cache_lock: - mod._pid_cache["FLOW"] = 42 - - with patch.object(mod, "_refresh_pid_cache"): - result = mod._get_pid_for_branch("flow") - assert result == 42 - - def test_strips_agent_suffix(self): - """Branch name ending in ' AGENT' is stripped before lookup.""" - mod = _import_monitor() - with mod._pid_cache_lock: - mod._pid_cache["FLOW"] = 42 - - with patch.object(mod, "_refresh_pid_cache"): - result = mod._get_pid_for_branch("flow agent") - assert result == 42 - - def test_returns_none_when_not_cached(self): - """Returns None when branch is not in cache.""" - mod = _import_monitor() - with mod._pid_cache_lock: - mod._pid_cache.clear() - - with patch.object(mod, "_refresh_pid_cache"): - result = mod._get_pid_for_branch("nonexistent") - assert result is None - - # --------------------------------------------------------------------------- # print_introspection tests (lines 130-167) # --------------------------------------------------------------------------- @@ -844,7 +656,7 @@ class TestFileWatcherWorker: "aipass.prax.apps.handlers.config.load": mock_config, }, ), - patch.object(mod, "_get_watch_directories", return_value=[("/tmp", True)]), + patch.object(mod, "_get_watch_directories", return_value=[("fakedir", True)]), patch.object(mod, "_start_observer_with_fallback", return_value=None), ): mod._file_watcher_worker() diff --git a/src/aipass/prax/tests/test_pid_cache.py b/src/aipass/prax/tests/test_pid_cache.py new file mode 100644 index 00000000..eba85544 --- /dev/null +++ b/src/aipass/prax/tests/test_pid_cache.py @@ -0,0 +1,188 @@ +# =================== AIPass ==================== +# Name: test_pid_cache.py +# Description: Tests for the PID cache handler +# Version: 1.0.0 +# Created: 2026-07-10 +# Modified: 2026-07-10 +# ============================================= + +"""Tests for apps/handlers/monitoring/pid_cache.py""" + +import json +import sys +import time +from unittest.mock import MagicMock, patch + +_HANDLER_MOCKS = { + "aipass.prax.apps.handlers.json": MagicMock(), + "aipass.prax.apps.handlers.json.json_handler": MagicMock(), +} + + +def _import_pid_cache(): + """Import (or reload) the pid_cache module with handler mocks.""" + fresh = {k: MagicMock() for k in _HANDLER_MOCKS} + with patch.dict(sys.modules, fresh): + import importlib + + if "aipass.prax.apps.handlers.monitoring.pid_cache" in sys.modules: + mod = importlib.reload(sys.modules["aipass.prax.apps.handlers.monitoring.pid_cache"]) + else: + mod = importlib.import_module("aipass.prax.apps.handlers.monitoring.pid_cache") + return mod + + +class TestParseLockPid: + """Test dispatch lock file parsing for PID cache.""" + + def test_no_lock_file_does_nothing(self, tmp_path): + """Branch entry without a lock file adds nothing to cache.""" + mod = _import_pid_cache() + new_cache: dict[str, int] = {} + entry = {"path": str(tmp_path / "somebranch"), "name": "flow"} + mod.parse_lock_pid(entry, new_cache) + assert new_cache == {} + + def test_lock_file_with_live_pid_on_linux(self, tmp_path): + """Lock file with a PID that has a /proc entry adds to cache.""" + mod = _import_pid_cache() + branch_dir = tmp_path / "mybranch" + mail_dir = branch_dir / "ai_mail.local" + mail_dir.mkdir(parents=True) + lock_data = {"pid": 12345} + (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") + + new_cache: dict[str, int] = {} + entry = {"path": str(branch_dir), "name": "flow"} + + with ( + patch("sys.platform", "linux"), + patch("pathlib.Path.exists", side_effect=lambda self=None: True), + ): + mod.parse_lock_pid(entry, new_cache) + + assert new_cache.get("FLOW") == 12345 + + def test_lock_file_with_zero_pid(self, tmp_path): + """Lock file with pid=0 skips entry.""" + mod = _import_pid_cache() + branch_dir = tmp_path / "mybranch" + mail_dir = branch_dir / "ai_mail.local" + mail_dir.mkdir(parents=True) + lock_data = {"pid": 0} + (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") + + new_cache: dict[str, int] = {} + entry = {"path": str(branch_dir), "name": "flow"} + mod.parse_lock_pid(entry, new_cache) + assert new_cache == {} + + def test_lock_file_with_invalid_json(self, tmp_path): + """Lock file with invalid JSON logs warning and continues.""" + mod = _import_pid_cache() + branch_dir = tmp_path / "mybranch" + mail_dir = branch_dir / "ai_mail.local" + mail_dir.mkdir(parents=True) + (mail_dir / ".dispatch.lock").write_text("{bad json}", encoding="utf-8") + + new_cache: dict[str, int] = {} + entry = {"path": str(branch_dir), "name": "flow"} + mod.parse_lock_pid(entry, new_cache) + assert new_cache == {} + + def test_lock_file_with_empty_name(self, tmp_path): + """Branch entry with empty name skips cache update.""" + mod = _import_pid_cache() + branch_dir = tmp_path / "mybranch" + mail_dir = branch_dir / "ai_mail.local" + mail_dir.mkdir(parents=True) + lock_data = {"pid": 99999} + (mail_dir / ".dispatch.lock").write_text(json.dumps(lock_data), encoding="utf-8") + + new_cache: dict[str, int] = {} + entry = {"path": str(branch_dir), "name": ""} + + with ( + patch("sys.platform", "linux"), + patch("pathlib.Path.exists", return_value=True), + ): + mod.parse_lock_pid(entry, new_cache) + assert new_cache == {} + + +class TestRefresh: + """Test PID cache refresh from registry.""" + + def test_skips_when_within_ttl(self): + """Cache refresh is skipped if within TTL window.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + setattr(mod, "_pid_cache_last_refresh", time.time()) + + with patch.object(mod, "parse_lock_pid") as mock_parse: + mod.refresh() + mock_parse.assert_not_called() + + def test_refreshes_when_ttl_expired(self, tmp_path): + """Cache refresh runs when TTL has expired.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + setattr(mod, "_pid_cache_last_refresh", 0.0) + + registry_data = {"branches": [{"name": "flow", "path": str(tmp_path / "flow")}]} + registry_file = tmp_path / "AIPASS_REGISTRY.json" + registry_file.write_text(json.dumps(registry_data), encoding="utf-8") + + mod.refresh(repo_root=tmp_path) + + def test_handles_missing_registry(self, tmp_path): + """Missing registry file does not crash.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + setattr(mod, "_pid_cache_last_refresh", 0.0) + + mod.refresh(repo_root=tmp_path) + + def test_handles_exception_in_refresh(self, tmp_path): + """Exception during refresh is caught and logged.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + setattr(mod, "_pid_cache_last_refresh", 0.0) + + registry_file = tmp_path / "AIPASS_REGISTRY.json" + registry_file.write_text("{corrupt", encoding="utf-8") + mod.refresh(repo_root=tmp_path) + + +class TestGetPidForBranch: + """Test PID lookup for branch names.""" + + def test_returns_pid_from_cache(self): + """Returns PID when branch is in cache.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + mod._pid_cache["FLOW"] = 42 + + with patch.object(mod, "refresh"): + result = mod.get_pid_for_branch("flow") + assert result == 42 + + def test_strips_agent_suffix(self): + """Branch name ending in ' AGENT' is stripped before lookup.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + mod._pid_cache["FLOW"] = 42 + + with patch.object(mod, "refresh"): + result = mod.get_pid_for_branch("flow agent") + assert result == 42 + + def test_returns_none_when_not_cached(self): + """Returns None when branch is not in cache.""" + mod = _import_pid_cache() + with mod._pid_cache_lock: + mod._pid_cache.clear() + + with patch.object(mod, "refresh"): + result = mod.get_pid_for_branch("nonexistent") + assert result is None diff --git a/src/aipass/prax/tests/test_telegram_relay.py b/src/aipass/prax/tests/test_telegram_relay.py index fdba8e93..71627774 100644 --- a/src/aipass/prax/tests/test_telegram_relay.py +++ b/src/aipass/prax/tests/test_telegram_relay.py @@ -352,6 +352,8 @@ class TestRenderEventCallsRelay: "aipass.prax.apps.handlers.monitoring.monitoring_filters": MagicMock(), "aipass.prax.apps.handlers.monitoring.file_watcher_integration": MagicMock(), "aipass.prax.apps.handlers.monitoring.telegram_relay": MagicMock(), + "aipass.prax.apps.handlers.monitoring.pid_cache": MagicMock(), + "aipass.prax.apps.handlers.monitoring.instance_lock": MagicMock(), } with patch.dict(sys.modules, fresh_mocks): if "aipass.prax.apps.modules.monitor" in sys.modules: