feat(flow): fix(flow): DPLAN-0141 — reach 100% seedgo (vestigial monitoring + lock ops + CLAUDE.md)

Co-Authored-By: @flow <flow@aipass>
This commit is contained in:
AIOSAI
2026-04-22 12:04:21 -07:00
co-authored by @flow
parent aa1fd6f1db
commit 53b5037af4
7 changed files with 173 additions and 615 deletions
+15 -12
View File
@@ -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": {
+42
View File
@@ -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
```
@@ -0,0 +1 @@
"""Flow Handlers - Runner infrastructure (lock files, process management)."""
@@ -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)
@@ -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)
@@ -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()
+14 -446
View File
@@ -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