From 53b5037af41a79f289a2ed4e19cb7cf9549649f3 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Wed, 22 Apr 2026 12:04:21 -0700 Subject: [PATCH] =?UTF-8?q?feat(flow):=20fix(flow):=20DPLAN-0141=20?= =?UTF-8?q?=E2=80=94=20reach=20100%=20seedgo=20(vestigial=20monitoring=20+?= =?UTF-8?q?=20lock=20ops=20+=20CLAUDE.md)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: @flow --- src/aipass/flow/.seedgo/bypass.json | 27 +- src/aipass/flow/CLAUDE.md | 42 ++ .../flow/apps/handlers/runner/__init__.py | 1 + .../flow/apps/handlers/runner/lock_ops.py | 84 ++++ .../flow/apps/modules/post_close_runner.py | 69 +-- .../flow/apps/modules/registry_monitor.py | 105 +--- .../flow/tests/test_monitor_registry.py | 460 +----------------- 7 files changed, 173 insertions(+), 615 deletions(-) create mode 100644 src/aipass/flow/CLAUDE.md create mode 100644 src/aipass/flow/apps/handlers/runner/__init__.py create mode 100644 src/aipass/flow/apps/handlers/runner/lock_ops.py diff --git a/src/aipass/flow/.seedgo/bypass.json b/src/aipass/flow/.seedgo/bypass.json index 988f37ee..57bca401 100644 --- a/src/aipass/flow/.seedgo/bypass.json +++ b/src/aipass/flow/.seedgo/bypass.json @@ -408,25 +408,28 @@ "reason": "Deferred per devpulse dispatch \u2014 test infrastructure not yet in place for flow branch" }, { - "file": "apps/handlers/registry/monitor_ops.py", - "standard": "unused_function", + "file": "tests/test_monitor_registry.py", + "standard": "architecture", "category": "false_positive", - "lines": [82], - "reason": "on_created is a watchdog FileSystemEventHandler callback invoked by the filesystem monitoring framework, not called directly by our code" + "reason": "Test files live in tests/ by convention, not in the apps/ 3-layer structure" }, { - "file": "apps/handlers/registry/monitor_ops.py", - "standard": "unused_function", + "file": "tests/test_monitor_registry.py", + "standard": "encapsulation", "category": "false_positive", - "lines": [92], - "reason": "on_deleted is a watchdog FileSystemEventHandler callback invoked by the filesystem monitoring framework, not called directly by our code" + "reason": "Unit tests import handlers directly to test implementation functions in isolation" }, { - "file": "apps/handlers/registry/monitor_ops.py", - "standard": "unused_function", + "file": "tests/*", + "standard": "architecture", "category": "false_positive", - "lines": [102], - "reason": "on_moved is a watchdog FileSystemEventHandler callback invoked by the filesystem monitoring framework, not called directly by our code" + "reason": "Test files live in tests/ by convention, not in the apps/ 3-layer structure" + }, + { + "file": "tests/*", + "standard": "encapsulation", + "category": "false_positive", + "reason": "Unit tests import handlers directly to test implementation functions in isolation" } ], "notes": { diff --git a/src/aipass/flow/CLAUDE.md b/src/aipass/flow/CLAUDE.md new file mode 100644 index 00000000..77d84fc7 --- /dev/null +++ b/src/aipass/flow/CLAUDE.md @@ -0,0 +1,42 @@ +# flow + +## Startup + +On any greeting, silently read these files and run the commands — no narration, no announcing steps. Just do it and respond with the status. + +**Read:** `.trinity/passport.json`, `.trinity/local.json`, `.trinity/observations.json`, `README.md`, `STATUS.local.md` +**Check:** If `.ai_mail.local/inbox.json` exists, read it. Process any mail. +**Run:** `git status` + +## Identity + +You are **flow** — an AIPass citizen. + +- **Module:** `aipass.flow` +- **Role:** workflow_planner +- **Purpose:** Workflow planning and tracking + +## Memories + +Update `.trinity/` at natural breakpoints, after milestones, and on `/memo`. + +- `local.json` — Session history, key learnings, active tasks +- `observations.json` — Collaboration patterns, insights +- `passport.json` — Identity (rarely changes) + +## AIPass Context + +This branch is part of the AIPass multi-agent framework. Key concepts: + +- **Branch** — your directory (`src/aipass/flow/`). Your home. +- **Citizen** — the identity that lives in a branch. Has a passport, memories, mailbox. +- **Agent** — a disposable worker spawned for a task. No passport, no memory. + +## Commands + +``` +drone systems # List available infrastructure +drone @ai_mail inbox # Check mailbox +drone @ai_mail send @branch "Subject" "Body" # Send mail +drone @seedgo audit @flow # Run standards audit +``` diff --git a/src/aipass/flow/apps/handlers/runner/__init__.py b/src/aipass/flow/apps/handlers/runner/__init__.py new file mode 100644 index 00000000..e0c733ce --- /dev/null +++ b/src/aipass/flow/apps/handlers/runner/__init__.py @@ -0,0 +1 @@ +"""Flow Handlers - Runner infrastructure (lock files, process management).""" diff --git a/src/aipass/flow/apps/handlers/runner/lock_ops.py b/src/aipass/flow/apps/handlers/runner/lock_ops.py new file mode 100644 index 00000000..6d21388c --- /dev/null +++ b/src/aipass/flow/apps/handlers/runner/lock_ops.py @@ -0,0 +1,84 @@ +# =================== AIPass ==================== +# Name: lock_ops.py +# Description: Process lock file operations handler +# Version: 1.0.0 +# Created: 2026-04-22 +# Modified: 2026-04-22 +# ============================================= + +""" +Lock File Operations Handler + +Provides atomic lock file management for background runners. +Uses O_CREAT | O_EXCL to avoid TOCTOU races between existence check and write. + +Usage: + from aipass.flow.apps.handlers.runner.lock_ops import ( + acquire_lock, release_lock + ) +""" + +import os +from pathlib import Path + +from aipass.prax import logger + +from aipass.flow.apps.handlers.json import json_handler + + +def try_create_lock(lock_file: Path) -> bool: + """Atomically create lock file with current PID. Returns True on success.""" + try: + fd = os.open(str(lock_file), os.O_CREAT | os.O_EXCL | os.O_WRONLY) + os.write(fd, str(os.getpid()).encode()) + os.close(fd) + return True + except FileExistsError: + logger.info("Lock file already exists, cannot acquire: %s", lock_file) + return False + + +def is_lock_stale(lock_file: Path) -> bool: + """Check if existing lock file belongs to a dead process.""" + try: + pid = int(lock_file.read_text(encoding="utf-8").strip()) + os.kill(pid, 0) + logger.info("Another instance running (PID %d), lock valid: %s", pid, lock_file) + return False + except (ValueError, ProcessLookupError, PermissionError): + logger.info("Stale lock found, taking over: %s", lock_file) + return True + + +def acquire_lock(lock_file: Path) -> bool: + """Try to acquire lock file. Returns True if acquired. + + Uses atomic O_CREAT | O_EXCL to avoid TOCTOU race. + """ + if try_create_lock(lock_file): + json_handler.log_operation("lock_acquired", {"lock_file": str(lock_file)}) + return True + + if not is_lock_stale(lock_file): + return False + + try: + lock_file.unlink() + except OSError as exc: + logger.warning("Failed to remove stale lock %s: %s", lock_file, exc) + return False + + if not try_create_lock(lock_file): + logger.info("Another process grabbed lock during retry: %s", lock_file) + return False + + json_handler.log_operation("lock_acquired", {"lock_file": str(lock_file), "stale_recovery": True}) + return True + + +def release_lock(lock_file: Path) -> None: + """Release the lock file.""" + try: + lock_file.unlink(missing_ok=True) + except OSError as exc: + logger.warning("Failed to release lock file %s: %s", lock_file, exc) diff --git a/src/aipass/flow/apps/modules/post_close_runner.py b/src/aipass/flow/apps/modules/post_close_runner.py index 93a0b9d2..7f53b56a 100644 --- a/src/aipass/flow/apps/modules/post_close_runner.py +++ b/src/aipass/flow/apps/modules/post_close_runner.py @@ -19,7 +19,6 @@ all unprocessed plans since it scans ALL of them). This script lives inside the flow branch so handler import guards allow it. """ -import os import sys from pathlib import Path @@ -41,6 +40,7 @@ LOCK_FILE = FLOW_ROOT / ".post_close_runner.lock" # AI summarization removed — plans vectorized directly from flow/processed_plans/ # from aipass.flow.apps.handlers.summary.generate import generate_summaries from aipass.flow.apps.handlers.mbank.process import process_closed_plans +from aipass.flow.apps.handlers.runner.lock_ops import acquire_lock, release_lock def handle_command(command: str, args: list) -> bool: @@ -73,7 +73,7 @@ def handle_command(command: str, args: list) -> bool: json_handler.log_operation("post_close_processed", {"command": command, "args": args}) # Run the post-close processing directly (foreground) - if not _acquire_lock(): + if not acquire_lock(LOCK_FILE): warning("Another instance is already running") return True @@ -84,69 +84,11 @@ def handle_command(command: str, args: list) -> bool: logger.error(f"[{MODULE_NAME}] Background processing failed: {e}") error(f"Processing failed: {e}") finally: - _release_lock() + release_lock(LOCK_FILE) return True -def _try_create_lock() -> bool: - """Atomically create lock file with current PID. Returns True on success.""" - try: - fd = os.open(str(LOCK_FILE), os.O_CREAT | os.O_EXCL | os.O_WRONLY) - os.write(fd, str(os.getpid()).encode()) - os.close(fd) - return True - except FileExistsError: - logger.info("[%s] Lock file already exists, cannot acquire", MODULE_NAME) - return False - - -def _is_lock_stale() -> bool: - """Check if existing lock file belongs to a dead process.""" - try: - pid = int(LOCK_FILE.read_text(encoding="utf-8").strip()) - os.kill(pid, 0) - logger.info("[%s] Another instance running (PID %d), exiting", MODULE_NAME, pid) - return False # Process alive — lock is valid - except (ValueError, ProcessLookupError, PermissionError): - logger.info("[%s] Stale lock found, taking over", MODULE_NAME) - return True - - -def _acquire_lock() -> bool: - """Try to acquire lock file. Returns True if acquired, False if another instance is running. - - Uses atomic O_CREAT | O_EXCL to avoid TOCTOU race between existence check and write. - """ - if _try_create_lock(): - return True - - # Lock exists — check if stale - if not _is_lock_stale(): - return False - - # Remove stale lock and retry - try: - LOCK_FILE.unlink() - except OSError as exc: - logger.warning("[%s] Failed to remove stale lock: %s", MODULE_NAME, exc) - return False - - if not _try_create_lock(): - logger.info("[%s] Another process grabbed lock during retry", MODULE_NAME) - return False - - return True - - -def _release_lock(): - """Release the lock file.""" - try: - LOCK_FILE.unlink(missing_ok=True) - except OSError as e: - logger.warning(f"[{MODULE_NAME}] Failed to release lock file: {e}") - - def print_introspection(): """Display module introspection info.""" console.print() @@ -182,13 +124,12 @@ if __name__ == "__main__": print_help() sys.exit(0) - if not _acquire_lock(): + if not acquire_lock(LOCK_FILE): sys.exit(0) try: - # generate_summaries() — removed, AI summarization no longer needed process_closed_plans() except Exception as e: logger.error(f"[{MODULE_NAME}] Background processing failed: {e}") finally: - _release_lock() + release_lock(LOCK_FILE) diff --git a/src/aipass/flow/apps/modules/registry_monitor.py b/src/aipass/flow/apps/modules/registry_monitor.py index aa9ca0e9..18400df6 100644 --- a/src/aipass/flow/apps/modules/registry_monitor.py +++ b/src/aipass/flow/apps/modules/registry_monitor.py @@ -1,9 +1,9 @@ # =================== AIPass ==================== # Name: registry_monitor.py -# Description: Registry auto-healing and file watching module -# Version: 2.1.0 +# Description: Registry auto-healing module +# Version: 3.0.0 # Created: 2025-11-21 -# Modified: 2025-11-21 +# Modified: 2026-04-22 # ============================================= """ @@ -93,42 +93,8 @@ def scan_plan_files() -> Dict[str, Any]: ) -def start_monitoring(): - """Start PLAN file monitoring with watchdog (thin orchestrator) - - Returns: - True if started successfully, False otherwise - """ - result = start_monitoring_impl(ecosystem_root=ECOSYSTEM_ROOT) - # Module handles display - status = result.get("status", "") - if status == "already_running": - warning("Monitor is already running") - elif status == "started": - console.print(f"[green]OK[/green] {result['message']}") - elif status == "error": - error(result["message"]) - return result.get("success", False) - - -def stop_monitoring(): - """Stop PLAN file monitoring (thin orchestrator) - - Returns: - True if stopped successfully, False otherwise - """ - result = stop_monitoring_impl() - # Module handles display - status = result.get("status", "") - if status == "stopped": - console.print("[green]OK[/green] Monitor stopped") - elif status == "not_running": - warning("Monitor is not running") - return result.get("success", False) - - def get_status() -> Dict[str, Any]: - """Get monitoring status (thin orchestrator)""" + """Get registry status (thin orchestrator)""" return get_status_impl( ecosystem_root=ECOSYSTEM_ROOT, load_registry=load_registry, @@ -147,9 +113,7 @@ def handle_command(command: str, args: List[str]) -> bool: Commands: scan - One-time scan and heal heal - Alias for scan - start - Start watchdog monitoring - stop - Stop watchdog monitoring - status - Show monitoring status + status - Show registry status Args: command: Command name @@ -200,54 +164,19 @@ def handle_command(command: str, args: List[str]) -> bool: console.print() return True - elif subcommand == "start": - console.print("[bold]Starting registry monitor...[/bold]") - console.print() - - # Run initial scan before starting monitor - console.print("[dim]Running initial scan...[/dim]") - scan_result = scan_plan_files() - console.print(f"[dim]Found {scan_result['total_plans']} PLAN files[/dim]") - console.print() - - success = start_monitoring() - if success: - console.print() - warning("Monitor is running. Press Ctrl+C to stop.") - console.print() - - # Keep script alive - try: - while True: - time.sleep(1) - except KeyboardInterrupt: - logger.info("[registry_monitor] Keyboard interrupt received, stopping monitor") - console.print() - console.print("[bold]Stopping monitor...[/bold]") - stop_monitoring() - console.print() - - return success - - elif subcommand == "stop": - return stop_monitoring() - elif subcommand == "status": status = get_status() console.print() - console.print("[bold cyan]Registry Monitor Status[/bold cyan]") + console.print("[bold cyan]Registry Status[/bold cyan]") console.print() console.print(f" • Version: {status['version']}") - console.print( - f" • Monitoring: {'[green]Active[/green]' if status['monitoring_active'] else '[yellow]Inactive[/yellow]'}" - ) console.print(f" • Watch location: {status['watch_location']}") console.print(f" • Total plans: {status['total_plans']}") console.print(f" • Open plans: {status['open_plans']}") console.print(f" • Ignored folders: {status['ignore_folders']}") console.print() - console.print("[dim]Commands: scan | start | stop | status[/dim]") + console.print("[dim]Commands: scan | status[/dim]") console.print() return True @@ -258,9 +187,7 @@ def handle_command(command: str, args: List[str]) -> bool: console.print("Available commands:") console.print(" • scan - One-time scan and heal registry") console.print(" • heal - Alias for scan") - console.print(" • start - Start watchdog monitoring") - console.print(" • stop - Stop watchdog monitoring") - console.print(" • status - Show monitoring status") + console.print(" • status - Show registry status") console.print() return False @@ -281,19 +208,15 @@ def print_introspection(): console.print() console.print("[yellow]Features:[/yellow]") - console.print(" • Real-time file watching (watchdog)") - console.print(" • Auto-detect create/move/delete events") console.print(" • Scan and heal registry") console.print(" • Duplicate detection with auto-renumbering") - console.print(" • Metadata preservation on moves") + console.print(" • Registry status reporting") console.print() console.print("[yellow]Commands:[/yellow]") console.print(" • scan - One-time scan and heal registry") console.print(" • heal - Alias for scan") - console.print(" • start - Start watchdog monitoring (runs until Ctrl+C)") - console.print(" • stop - Stop watchdog monitoring") - console.print(" • status - Show monitoring status") + console.print(" • status - Show registry status") console.print() console.print("[yellow]Connected Handlers:[/yellow]") @@ -317,15 +240,11 @@ def print_help(): console.print("[yellow]SUBCOMMANDS:[/yellow]") console.print(" scan One-time scan and heal") console.print(" heal Alias for scan") - console.print(" start Start persistent watchdog monitoring") - console.print(" stop Stop watchdog monitoring") - console.print(" status Check monitoring status") + console.print(" status Show registry status") console.print() console.print("[yellow]EXAMPLES:[/yellow]") console.print(" [dim]drone @flow registry scan[/dim] # One-time scan and heal") - console.print(" [dim]drone @flow registry start[/dim] # Start persistent monitoring") - console.print(" [dim]drone @flow registry status[/dim] # Check monitoring status") - console.print(" [dim]drone @flow registry stop[/dim] # Stop monitoring") + console.print(" [dim]drone @flow registry status[/dim] # Check registry status") console.print() diff --git a/src/aipass/flow/tests/test_monitor_registry.py b/src/aipass/flow/tests/test_monitor_registry.py index 2b02e473..670348b1 100644 --- a/src/aipass/flow/tests/test_monitor_registry.py +++ b/src/aipass/flow/tests/test_monitor_registry.py @@ -1,8 +1,15 @@ +# =================== AIPass ==================== +# Name: test_monitor_registry.py +# Description: Tests for monitor_ops handler and registry_monitor module +# Version: 2.0.0 +# Created: 2026-03-08 +# Modified: 2026-04-22 +# ============================================= + """Tests for monitor_ops handler and registry_monitor module.""" import builtins import os -import time import types from collections.abc import Mapping, Sequence from pathlib import Path @@ -34,16 +41,6 @@ def _make_plan_file(directory: Path, number: str) -> Path: return plan_file -def _make_event(src_path: str, dest_path: str | None = None, is_directory: bool = False): - """Create a mock watchdog event object.""" - event = MagicMock() - event.src_path = src_path - event.is_directory = is_directory - if dest_path is not None: - event.dest_path = dest_path - return event - - # ═══════════════════════════════════════════════════════════ # 1. handle_walk_error # ═══════════════════════════════════════════════════════════ @@ -56,21 +53,16 @@ class TestHandleWalkError: """PermissionError should not trigger a warning log.""" mod = _import_monitor_ops() with patch.object(mod, "_fire_event", return_value=False): - # The handle_walk_error function is defined inside scan_plan_files_impl. - # We exercise it by creating a directory we cannot read. restricted = tmp_path / "restricted" restricted.mkdir() - # Create a plan file in a readable subdirectory so scan itself works _make_plan_file(tmp_path, "0001") - # Make the restricted dir unreadable os.chmod(str(restricted), 0o000) try: result = mod.scan_plan_files_impl( ecosystem_root=tmp_path, load_registry=lambda: {"plans": {}}, ) - # Scan should complete without crashing assert isinstance(result, dict) assert "total_plans" in result finally: @@ -79,16 +71,12 @@ class TestHandleWalkError: def test_generic_os_error_logs_warning(self, tmp_path, mock_logger): """Non-PermissionError OSError should be logged as warning.""" mod = _import_monitor_ops() - # We cannot easily trigger a generic OSError from os.walk, but we - # can directly call the handle_walk_error closure pattern by - # simulating a scan on a non-existent directory. missing = tmp_path / "nonexistent_root" with patch.object(mod, "_fire_event", return_value=False): result = mod.scan_plan_files_impl( ecosystem_root=missing, load_registry=lambda: {"plans": {}}, ) - # Should not crash, just return empty results assert result["total_plans"] == 0 assert result["added"] == [] @@ -133,10 +121,8 @@ class TestScanPlanFilesImpl: def test_ignores_non_plan_files(self, tmp_path): """Non-plan files should be ignored even if they look similar.""" mod = _import_monitor_ops() - # Valid plan files (any prefix: FPLAN, DPLAN, etc.) _make_plan_file(tmp_path, "0001") (tmp_path / "DPLAN-0002.md").write_text("also a plan", encoding="utf-8") - # Invalid files that should not match (tmp_path / "FPLAN-ABC.md").write_text("bad number", encoding="utf-8") (tmp_path / "README.md").write_text("readme", encoding="utf-8") (tmp_path / "NOTES-0003.md").write_text("not a plan prefix", encoding="utf-8") @@ -151,7 +137,6 @@ class TestScanPlanFilesImpl: def test_skips_ignored_folders(self, tmp_path): """Directories in IGNORE_FOLDERS should be skipped.""" mod = _import_monitor_ops() - # Plan file in ignored directory git_dir = tmp_path / ".git" git_dir.mkdir() _make_plan_file(git_dir, "0001") @@ -160,7 +145,6 @@ class TestScanPlanFilesImpl: pycache_dir.mkdir() _make_plan_file(pycache_dir, "0002") - # Plan file in non-ignored directory good_dir = tmp_path / "active" good_dir.mkdir() _make_plan_file(good_dir, "0003") @@ -177,7 +161,6 @@ class TestScanPlanFilesImpl: def test_detects_orphaned_registry_entries(self, tmp_path): """Registry entries with no matching file should fire deleted events.""" mod = _import_monitor_ops() - # No plan files on disk, but registry has entries registry = { "plans": { "0001": {"file_path": str(tmp_path / "FPLAN-0001.md"), "status": "open"}, @@ -239,7 +222,6 @@ class TestScanPlanFilesImpl: def test_duplicate_plan_files_renumbered(self, tmp_path): """Duplicate plan numbers should be auto-renumbered.""" mod = _import_monitor_ops() - # Create two directories with same plan number dir_a = tmp_path / "project_a" dir_a.mkdir() dir_b = tmp_path / "project_b" @@ -287,265 +269,16 @@ class TestScanPlanFilesImpl: # ═══════════════════════════════════════════════════════════ -# 3. PlanFileWatcher event handlers +# 3. get_status_impl # ═══════════════════════════════════════════════════════════ -class TestPlanFileWatcherOnCreated: - """Tests for PlanFileWatcher.on_created.""" - - def setup_method(self): - """Clear deduplication state before each test.""" - mod = _import_monitor_ops() - mod._recent_events.clear() - - def test_created_event_for_plan_file(self, tmp_path): - """on_created should log and schedule fire for a valid FPLAN file.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0042.md")) - - with patch.object(watcher, "_schedule_fire_created") as mock_fire: - watcher.on_created(event) - mock_fire.assert_called_once() - call_path = mock_fire.call_args[0][0] - assert call_path.name == "FPLAN-0042.md" - - def test_created_event_ignores_directory(self, tmp_path): - """on_created should ignore directory events.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0042.md"), is_directory=True) - - with patch.object(watcher, "_schedule_fire_created") as mock_fire: - watcher.on_created(event) - mock_fire.assert_not_called() - - def test_created_event_ignores_non_plan_file(self, tmp_path): - """on_created should ignore non-FPLAN files.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "README.md")) - - with patch.object(watcher, "_schedule_fire_created") as mock_fire: - watcher.on_created(event) - mock_fire.assert_not_called() - - def test_created_event_deduplication(self, tmp_path): - """Duplicate create events for the same plan within DEDUPE_WINDOW should be ignored.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0042.md")) - - with patch.object(watcher, "_schedule_fire_created") as mock_fire: - watcher.on_created(event) - watcher.on_created(event) # duplicate - assert mock_fire.call_count == 1 - - -class TestPlanFileWatcherOnDeleted: - """Tests for PlanFileWatcher.on_deleted.""" - - def setup_method(self): - mod = _import_monitor_ops() - mod._recent_events.clear() - - def test_deleted_event_for_plan_file(self, tmp_path): - """on_deleted should log and schedule fire for a valid FPLAN file.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0007.md")) - - with patch.object(watcher, "_schedule_fire_deleted") as mock_fire: - watcher.on_deleted(event) - mock_fire.assert_called_once() - call_path = mock_fire.call_args[0][0] - assert call_path.name == "FPLAN-0007.md" - - def test_deleted_event_ignores_directory(self, tmp_path): - """on_deleted should ignore directory events.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0007.md"), is_directory=True) - - with patch.object(watcher, "_schedule_fire_deleted") as mock_fire: - watcher.on_deleted(event) - mock_fire.assert_not_called() - - def test_deleted_event_ignores_non_plan_file(self, tmp_path): - """on_deleted should ignore non-FPLAN files.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "notes.txt")) - - with patch.object(watcher, "_schedule_fire_deleted") as mock_fire: - watcher.on_deleted(event) - mock_fire.assert_not_called() - - def test_deleted_event_deduplication(self, tmp_path): - """Duplicate delete events within DEDUPE_WINDOW should be ignored.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event(str(tmp_path / "FPLAN-0007.md")) - - with patch.object(watcher, "_schedule_fire_deleted") as mock_fire: - watcher.on_deleted(event) - watcher.on_deleted(event) # duplicate - assert mock_fire.call_count == 1 - - -class TestPlanFileWatcherOnMoved: - """Tests for PlanFileWatcher.on_moved.""" - - def setup_method(self): - mod = _import_monitor_ops() - mod._recent_events.clear() - - def test_moved_event_for_plan_file(self, tmp_path): - """on_moved should log and schedule fire when dest is a valid FPLAN file.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - src = str(tmp_path / "old" / "FPLAN-0003.md") - dest = str(tmp_path / "new" / "FPLAN-0003.md") - event = _make_event(src, dest_path=dest) - - with patch.object(watcher, "_schedule_fire_moved") as mock_fire: - watcher.on_moved(event) - mock_fire.assert_called_once() - call_src = mock_fire.call_args[0][0] - call_dest = mock_fire.call_args[0][1] - assert call_src == Path(src) - assert call_dest == Path(dest) - - def test_moved_event_ignores_directory(self, tmp_path): - """on_moved should ignore directory events.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event( - str(tmp_path / "FPLAN-0003.md"), - dest_path=str(tmp_path / "new" / "FPLAN-0003.md"), - is_directory=True, - ) - - with patch.object(watcher, "_schedule_fire_moved") as mock_fire: - watcher.on_moved(event) - mock_fire.assert_not_called() - - def test_moved_event_ignores_non_plan_dest(self, tmp_path): - """on_moved should ignore moves where dest is not a FPLAN file.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event( - str(tmp_path / "FPLAN-0003.md"), - dest_path=str(tmp_path / "renamed.txt"), - ) - - with patch.object(watcher, "_schedule_fire_moved") as mock_fire: - watcher.on_moved(event) - mock_fire.assert_not_called() - - def test_moved_event_deduplication(self, tmp_path): - """Duplicate move events within DEDUPE_WINDOW should be ignored.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - event = _make_event( - str(tmp_path / "old" / "FPLAN-0003.md"), - dest_path=str(tmp_path / "new" / "FPLAN-0003.md"), - ) - - with patch.object(watcher, "_schedule_fire_moved") as mock_fire: - watcher.on_moved(event) - watcher.on_moved(event) # duplicate - assert mock_fire.call_count == 1 - - -# ═══════════════════════════════════════════════════════════ -# 4. start_monitoring_impl / stop_monitoring_impl / get_status_impl -# ═══════════════════════════════════════════════════════════ - - -class TestStartMonitoringImpl: - """Tests for start_monitoring_impl in monitor_ops.""" - - def teardown_method(self): - """Ensure observer is stopped after each test.""" - mod = _import_monitor_ops() - if mod._observer and mod._observer.is_alive(): - mod._observer.stop() - mod._observer.join() - mod._observer = None - - def test_start_returns_success(self, tmp_path): - """Starting monitor on a valid directory should succeed.""" - mod = _import_monitor_ops() - mod._observer = None - result = mod.start_monitoring_impl(tmp_path) - assert result["success"] is True - assert result["status"] == "started" - assert str(tmp_path) in result["message"] - - def test_start_when_already_running(self, tmp_path): - """Starting monitor when already running should return already_running.""" - mod = _import_monitor_ops() - mod._observer = None - # Start once - mod.start_monitoring_impl(tmp_path) - # Start again - result = mod.start_monitoring_impl(tmp_path) - assert result["success"] is False - assert result["status"] == "already_running" - - def test_start_handles_observer_exception(self, tmp_path): - """If Observer raises, start should return error status.""" - mod = _import_monitor_ops() - mod._observer = None - with patch("aipass.flow.apps.handlers.registry.monitor_ops.Observer") as mock_obs: - mock_obs.return_value.start.side_effect = RuntimeError("Cannot start") - result = mod.start_monitoring_impl(tmp_path) - assert result["success"] is False - assert result["status"] == "error" - assert "Cannot start" in result["message"] - - -class TestStopMonitoringImpl: - """Tests for stop_monitoring_impl in monitor_ops.""" - - def teardown_method(self): - mod = _import_monitor_ops() - mod._observer = None - - def test_stop_running_observer(self, tmp_path): - """Stopping a running observer should succeed.""" - mod = _import_monitor_ops() - mod._observer = None - mod.start_monitoring_impl(tmp_path) - result = mod.stop_monitoring_impl() - assert result["success"] is True - assert result["status"] == "stopped" - - def test_stop_when_not_running(self): - """Stopping when no observer is running should return not_running.""" - mod = _import_monitor_ops() - mod._observer = None - result = mod.stop_monitoring_impl() - assert result["success"] is False - assert result["status"] == "not_running" - - class TestGetStatusImpl: """Tests for get_status_impl in monitor_ops.""" - def teardown_method(self): + def test_status_returns_correct_fields(self, tmp_path): + """Status should return all expected fields.""" mod = _import_monitor_ops() - if mod._observer and mod._observer.is_alive(): - mod._observer.stop() - mod._observer.join() - mod._observer = None - - def test_status_when_not_monitoring(self, tmp_path): - """Status should report inactive when no observer is running.""" - mod = _import_monitor_ops() - mod._observer = None registry = { "plans": { "0001": {"status": "open"}, @@ -554,7 +287,7 @@ class TestGetStatusImpl: } } result = mod.get_status_impl(tmp_path, load_registry=lambda: registry) - assert not result["monitoring_active"] + assert result["monitoring_active"] is False assert result["total_plans"] == 3 assert result["open_plans"] == 2 assert result["watch_location"] == str(tmp_path) @@ -562,116 +295,19 @@ class TestGetStatusImpl: assert result["version"] == "2.0.0" assert result["ignore_folders"] == len(mod.IGNORE_FOLDERS) - def test_status_when_monitoring_active(self, tmp_path): - """Status should report active when observer is running.""" - mod = _import_monitor_ops() - mod._observer = None - mod.start_monitoring_impl(tmp_path) - result = mod.get_status_impl(tmp_path, load_registry=lambda: {"plans": {}}) - assert result["monitoring_active"] is True - def test_status_with_empty_registry(self, tmp_path): """Status should handle empty registry.""" mod = _import_monitor_ops() - mod._observer = None result = mod.get_status_impl(tmp_path, load_registry=lambda: {"plans": {}}) assert result["total_plans"] == 0 assert result["open_plans"] == 0 # ═══════════════════════════════════════════════════════════ -# 5. registry_monitor module wrappers +# 4. registry_monitor module wrappers # ═══════════════════════════════════════════════════════════ -class TestRegistryMonitorStartMonitoring: - """Tests for registry_monitor.start_monitoring wrapper.""" - - def test_start_monitoring_success(self): - """Successful start should print success message and return True.""" - mod = _import_registry_monitor() - mock_con = MagicMock() - with ( - patch.object( - mod, - "start_monitoring_impl", - return_value={"success": True, "status": "started", "message": "Monitor started"}, - ), - patch.object(mod, "console", mock_con), - ): - result = mod.start_monitoring() - assert result is True - mock_con.print.assert_called() - - def test_start_monitoring_already_running(self): - """Already-running should trigger warning and return False.""" - mod = _import_registry_monitor() - mock_warn = MagicMock() - with ( - patch.object( - mod, - "start_monitoring_impl", - return_value={"success": False, "status": "already_running", "message": "Already running"}, - ), - patch.object(mod, "warning", mock_warn), - ): - result = mod.start_monitoring() - assert result is False - mock_warn.assert_called_once_with("Monitor is already running") - - def test_start_monitoring_error(self): - """Error status should trigger error display and return False.""" - mod = _import_registry_monitor() - mock_err = MagicMock() - with ( - patch.object( - mod, - "start_monitoring_impl", - return_value={"success": False, "status": "error", "message": "Observer failed"}, - ), - patch.object(mod, "error", mock_err), - ): - result = mod.start_monitoring() - assert result is False - mock_err.assert_called_once_with("Observer failed") - - -class TestRegistryMonitorStopMonitoring: - """Tests for registry_monitor.stop_monitoring wrapper.""" - - def test_stop_monitoring_success(self): - """Successful stop should print message and return True.""" - mod = _import_registry_monitor() - mock_con = MagicMock() - with ( - patch.object( - mod, - "stop_monitoring_impl", - return_value={"success": True, "status": "stopped", "message": "Monitor stopped"}, - ), - patch.object(mod, "console", mock_con), - ): - result = mod.stop_monitoring() - assert result is True - mock_con.print.assert_called() - - def test_stop_monitoring_not_running(self): - """Stopping when not running should trigger warning and return False.""" - mod = _import_registry_monitor() - mock_warn = MagicMock() - with ( - patch.object( - mod, - "stop_monitoring_impl", - return_value={"success": False, "status": "not_running", "message": "Not running"}, - ), - patch.object(mod, "warning", mock_warn), - ): - result = mod.stop_monitoring() - assert result is False - mock_warn.assert_called_once_with("Monitor is not running") - - class TestRegistryMonitorGetStatus: """Tests for registry_monitor.get_status wrapper.""" @@ -697,7 +333,7 @@ class TestRegistryMonitorGetStatus: # ═══════════════════════════════════════════════════════════ -# 6. _fire_event helper +# 5. _fire_event helper # ═══════════════════════════════════════════════════════════ @@ -736,71 +372,3 @@ class TestFireEvent: with patch.object(builtins, "__import__", side_effect=_failing_import): result = mod._fire_event("test_event") assert result is False - - -# ═══════════════════════════════════════════════════════════ -# 7. PlanFileWatcher internal methods -# ═══════════════════════════════════════════════════════════ - - -class TestPlanFileWatcherInternals: - """Tests for PlanFileWatcher helper methods.""" - - def test_is_plan_file_valid(self): - """Valid plan filenames (any prefix) should return True.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - assert watcher._is_plan_file("/some/path/FPLAN-0001.md") is True - assert watcher._is_plan_file("/some/path/FPLAN-9999.md") is True - assert watcher._is_plan_file("/some/path/DPLAN-0001.md") is True - assert watcher._is_plan_file("/some/path/APLAN-0042.md") is True - assert watcher._is_plan_file("/some/path/RPLAN-0100.md") is True - assert watcher._is_plan_file("/some/path/TDPLAN-0002.md") is True - - def test_is_plan_file_invalid(self): - """Invalid filenames should return False.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - assert watcher._is_plan_file("/some/path/FPLAN-ABC.md") is False - assert watcher._is_plan_file("/some/path/FPLAN-00001.md") is False - assert watcher._is_plan_file("/some/path/README.md") is False - assert watcher._is_plan_file("/some/path/FPLAN-0001.txt") is False - assert watcher._is_plan_file("/some/path/plan-0001.md") is False - - def test_get_plan_number(self): - """Should extract the 4-digit number from plan filename.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - assert watcher._get_plan_number(Path("FPLAN-0042.md")) == "0042" - assert watcher._get_plan_number(Path("FPLAN-0001.md")) == "0001" - assert watcher._get_plan_number(Path("/deep/path/FPLAN-1234.md")) == "1234" - assert watcher._get_plan_number(Path("DPLAN-0005.md")) == "0005" - assert watcher._get_plan_number(Path("TDPLAN-0002.md")) == "0002" - - def test_get_plan_number_invalid(self): - """Invalid filenames should return None.""" - mod = _import_monitor_ops() - watcher = mod.PlanFileWatcher() - assert watcher._get_plan_number(Path("README.md")) is None - assert watcher._get_plan_number(Path("FPLAN-ABC.md")) is None - - def test_deduplication_window_expires(self): - """Events outside DEDUPE_WINDOW should not be considered duplicates.""" - mod = _import_monitor_ops() - mod._recent_events.clear() - watcher = mod.PlanFileWatcher() - - # Add an event with an old timestamp - mod._recent_events.append(("created", "0001", time.time() - 10.0)) - - # Should not be duplicate since old event is beyond DEDUPE_WINDOW - assert watcher._is_duplicate_event("created", "0001") is False - - def test_different_event_types_not_deduplicated(self): - """Different event types for same plan should not be deduplicated.""" - mod = _import_monitor_ops() - mod._recent_events.clear() - watcher = mod.PlanFileWatcher() - - assert watcher._is_duplicate_event("created", "0001") is False - assert watcher._is_duplicate_event("deleted", "0001") is False # different type