#671 prax: single-instance lock on the monitor — stop duplicate/orphan monitors from double-sending Telegram relay. New instance_lock handler writes a liveness-checked pidfile (prax_json/monitor.pid, outside the tailed system_logs/): acquire() before relay init refuses to start (fail-loud, names the holding PID) when a live monitor holds the lock, reclaims a stale pidfile on a dead PID, release() clears it on shutdown. Liveness probe platform-branched — POSIX os.kill(pid,0), Windows OpenProcess/GetExitCodeProcess (a raw os.kill(pid,0) TERMINATES the target on Windows; reused the canonical devpulse/watchdog/agent.py:137 impl). monitor.py split under the 600-line limit (pid_cache extracted). Built by @prax, verified by devpulse: 120 tests green across the 4 touched files (+25 new incl 3 Windows-path), seedgo 31/31 on all 3 sources. Verify caught the Windows os.kill hazard on the first pass; filed the seedgo windows_compat detector gap as #682.

This commit is contained in:
AIOSAI
2026-07-10 02:53:29 -07:00
parent 10a8f738a0
commit 739dada015
8 changed files with 675 additions and 339 deletions
+12
View File
@@ -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]
@@ -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()})
@@ -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)
+45 -148
View File
@@ -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")
+207
View File
@@ -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
+3 -191
View File
@@ -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()
+188
View File
@@ -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
@@ -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: