From d4aa8330dccbac5acefc698f2ce3812557b69248 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Fri, 10 Apr 2026 01:47:59 -0700 Subject: [PATCH] =?UTF-8?q?feat(system):=20S84:=20CRITICAL=20fixes=20?= =?UTF-8?q?=E2=80=94=20memory=20(path+security),=20ai=5Fmail=20(locking),?= =?UTF-8?q?=20prax=20(thread=20safety),=20global=20prompt=20patrick-file?= =?UTF-8?q?=20rule?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: @devpulse --- .../apps/handlers/email/inbox_cleanup.py | 153 +++++++++--------- .../ai_mail/apps/handlers/email/inbox_ops.py | 19 ++- .../cli/apps/handlers/init/bootstrap.py | 3 +- .../memory/apps/handlers/central_writer.py | 2 +- .../memory/apps/handlers/dashboard_push.py | 2 +- .../memory/apps/handlers/learnings/manager.py | 11 +- .../memory/apps/handlers/schema/normalize.py | 11 +- .../apps/handlers/symbolic/deduplicator.py | 26 +-- .../apps/handlers/symbolic/extractor.py | 27 +--- src/aipass/memory/tests/conftest.py | 8 +- .../prax/apps/handlers/logging/setup.py | 12 +- .../apps/handlers/monitoring/event_queue.py | 37 +++-- src/aipass/prax/apps/modules/logger.py | 8 +- src/aipass/prax/apps/modules/monitor.py | 41 ++--- src/aipass/prax/tests/test_event_queue.py | 2 +- 15 files changed, 194 insertions(+), 168 deletions(-) diff --git a/src/aipass/ai_mail/apps/handlers/email/inbox_cleanup.py b/src/aipass/ai_mail/apps/handlers/email/inbox_cleanup.py index 3b482ac2..baf707e0 100644 --- a/src/aipass/ai_mail/apps/handlers/email/inbox_cleanup.py +++ b/src/aipass/ai_mail/apps/handlers/email/inbox_cleanup.py @@ -227,31 +227,32 @@ def mark_all_read_and_archive(branch_path: Path) -> Tuple[bool, str, int]: _migrate_deleted_json_if_exists(mailbox_path) try: - # Load inbox - with open(inbox_file, 'r', encoding='utf-8') as f: - inbox_data = json.load(f) + with _get_inbox_lock()(inbox_file): + # Load inbox + with open(inbox_file, 'r', encoding='utf-8') as f: + inbox_data = json.load(f) - messages = inbox_data.get("messages", []) - count = len(messages) + messages = inbox_data.get("messages", []) + count = len(messages) - if count == 0: - return True, "Inbox already empty", 0 + if count == 0: + return True, "Inbox already empty", 0 - # Mark all as read and save to deleted/ folder - for msg in messages: - msg["read"] = True - _save_to_deleted_folder(mailbox_path, msg) + # Mark all as read and save to deleted/ folder + for msg in messages: + msg["read"] = True + _save_to_deleted_folder(mailbox_path, msg) - # Clear inbox - inbox_data["messages"] = [] - inbox_data["total_messages"] = 0 - inbox_data["unread_count"] = 0 + # Clear inbox + inbox_data["messages"] = [] + inbox_data["total_messages"] = 0 + inbox_data["unread_count"] = 0 - # Save inbox - with open(inbox_file, 'w', encoding='utf-8') as f: - json.dump(inbox_data, f, indent=2, ensure_ascii=False) + # Save inbox + with open(inbox_file, 'w', encoding='utf-8') as f: + json.dump(inbox_data, f, indent=2, ensure_ascii=False) - # Update dashboard (all zeros - inbox cleared) + # Update dashboard (outside lock - not inbox.json) _update_dashboard(branch_path, 0, 0, 0) # Trigger auto-purge of deleted folder @@ -315,37 +316,38 @@ def mark_as_opened(branch_path: Path, message_id: str) -> Tuple[bool, str, Optio return False, f"Inbox not found: {inbox_file}", None try: - with open(inbox_file, 'r', encoding='utf-8') as f: - inbox_data = json.load(f) + with _get_inbox_lock()(inbox_file): + with open(inbox_file, 'r', encoding='utf-8') as f: + inbox_data = json.load(f) - messages = inbox_data.get("messages", []) - target_msg = None + messages = inbox_data.get("messages", []) + target_msg = None - for msg in messages: - if msg.get("id") == message_id: - target_msg = msg - break + for msg in messages: + if msg.get("id") == message_id: + target_msg = msg + break - if target_msg is None: - return False, f"Message not found: {message_id}", None + if target_msg is None: + return False, f"Message not found: {message_id}", None - # Update status to opened (v2 schema) - target_msg["status"] = "opened" - # Keep backward compat - target_msg["read"] = True + # Update status to opened (v2 schema) + target_msg["status"] = "opened" + # Keep backward compat + target_msg["read"] = True - # Recalculate status counts (v2 schema) - new_count = sum( - 1 for m in messages - if m.get("status") == "new" or (m.get("status") is None and not m.get("read", False)) - ) - opened_count = sum(1 for m in messages if m.get("status") == "opened") - inbox_data["unread_count"] = new_count + # Recalculate status counts (v2 schema) + new_count = sum( + 1 for m in messages + if m.get("status") == "new" or (m.get("status") is None and not m.get("read", False)) + ) + opened_count = sum(1 for m in messages if m.get("status") == "opened") + inbox_data["unread_count"] = new_count - with open(inbox_file, 'w', encoding='utf-8') as f: - json.dump(inbox_data, f, indent=2, ensure_ascii=False) + with open(inbox_file, 'w', encoding='utf-8') as f: + json.dump(inbox_data, f, indent=2, ensure_ascii=False) - # Update dashboard + # Update dashboard (outside lock - not inbox.json) _update_dashboard(branch_path, new_count, opened_count, inbox_data["total_messages"]) return True, f"Message {message_id} marked as opened", target_msg @@ -379,48 +381,49 @@ def mark_as_closed_and_archive(branch_path: Path, message_id: str, skip_post_ops _migrate_deleted_json_if_exists(mailbox_path) try: - with open(inbox_file, 'r', encoding='utf-8') as f: - inbox_data = json.load(f) + with _get_inbox_lock()(inbox_file): + with open(inbox_file, 'r', encoding='utf-8') as f: + inbox_data = json.load(f) - messages = inbox_data.get("messages", []) - message_to_archive = None - message_index = None + messages = inbox_data.get("messages", []) + message_to_archive = None + message_index = None - for i, msg in enumerate(messages): - if msg.get("id") == message_id: - message_to_archive = msg - message_index = i - break + for i, msg in enumerate(messages): + if msg.get("id") == message_id: + message_to_archive = msg + message_index = i + break - if message_to_archive is None: - return False, f"Message not found: {message_id}" + if message_to_archive is None: + return False, f"Message not found: {message_id}" - # Mark as closed (v2 schema) - message_to_archive["status"] = "closed" - message_to_archive["read"] = True # backward compat + # Mark as closed (v2 schema) + message_to_archive["status"] = "closed" + message_to_archive["read"] = True # backward compat - # Remove from inbox - messages.pop(message_index) + # Remove from inbox + messages.pop(message_index) - # Update inbox counts - inbox_data["messages"] = messages - inbox_data["total_messages"] = len(messages) - # v2 status counts - new_count = sum( - 1 for m in messages - if m.get("status") == "new" or (m.get("status") is None and not m.get("read", False)) - ) - opened_count = sum(1 for m in messages if m.get("status") == "opened") - inbox_data["unread_count"] = new_count + # Update inbox counts + inbox_data["messages"] = messages + inbox_data["total_messages"] = len(messages) + # v2 status counts + new_count = sum( + 1 for m in messages + if m.get("status") == "new" or (m.get("status") is None and not m.get("read", False)) + ) + opened_count = sum(1 for m in messages if m.get("status") == "opened") + inbox_data["unread_count"] = new_count - with open(inbox_file, 'w', encoding='utf-8') as f: - json.dump(inbox_data, f, indent=2, ensure_ascii=False) + with open(inbox_file, 'w', encoding='utf-8') as f: + json.dump(inbox_data, f, indent=2, ensure_ascii=False) - # Save to deleted/ folder (new pattern) - _save_to_deleted_folder(mailbox_path, message_to_archive) + # Save to deleted/ folder (inside lock to ensure consistency) + _save_to_deleted_folder(mailbox_path, message_to_archive) if not skip_post_ops: - # Update dashboard + # Update dashboard (outside lock - not inbox.json) _update_dashboard(branch_path, new_count, opened_count, inbox_data["total_messages"]) # Trigger auto-purge of deleted folder diff --git a/src/aipass/ai_mail/apps/handlers/email/inbox_ops.py b/src/aipass/ai_mail/apps/handlers/email/inbox_ops.py index 07259718..e36a4485 100644 --- a/src/aipass/ai_mail/apps/handlers/email/inbox_ops.py +++ b/src/aipass/ai_mail/apps/handlers/email/inbox_ops.py @@ -20,6 +20,18 @@ from typing import Dict from aipass.prax.apps.modules.logger import system_logger as logger from aipass.ai_mail.apps.handlers.json import json_handler +# Lazy import for inbox file lock +_inbox_lock = None + + +def _get_inbox_lock(): + """Lazy import inbox_lock context manager.""" + global _inbox_lock + if _inbox_lock is None: + from aipass.ai_mail.apps.handlers.email.inbox_lock import inbox_lock + _inbox_lock = inbox_lock + return _inbox_lock + def load_inbox(inbox_file: Path) -> Dict: @@ -73,11 +85,12 @@ def load_inbox(inbox_file: Path) -> Dict: ) migrated = True - # Persist migration + # Persist migration under lock to prevent concurrent write races if migrated: try: - with open(inbox_file, 'w', encoding='utf-8') as f: - json.dump(inbox_data, f, indent=2, ensure_ascii=False) + with _get_inbox_lock()(inbox_file): + with open(inbox_file, 'w', encoding='utf-8') as f: + json.dump(inbox_data, f, indent=2, ensure_ascii=False) except Exception as e: logger.warning("[inbox] Migration persist failed for %s: %s", inbox_file, e) diff --git a/src/aipass/cli/apps/handlers/init/bootstrap.py b/src/aipass/cli/apps/handlers/init/bootstrap.py index b092a905..b130813e 100644 --- a/src/aipass/cli/apps/handlers/init/bootstrap.py +++ b/src/aipass/cli/apps/handlers/init/bootstrap.py @@ -185,7 +185,7 @@ def _readme_md(name: str) -> str: "aipass init agent my_agent\n" "\n" "# 2. Start a session\n" - "cd my_agent/\n" + "cd src/my_agent/\n" "claude # or your preferred AI CLI\n" "\n" "# 3. Check project status\n" @@ -272,6 +272,7 @@ def _gitignore() -> str: ".trinity/\n" ".ai_mail.local/\n" "*.local.*\n" + "!STATUS.local.md\n" "\n" "# Plans (local working docs)\n" "DPLAN-*\n" diff --git a/src/aipass/memory/apps/handlers/central_writer.py b/src/aipass/memory/apps/handlers/central_writer.py index 50c61f40..1fd2be91 100644 --- a/src/aipass/memory/apps/handlers/central_writer.py +++ b/src/aipass/memory/apps/handlers/central_writer.py @@ -35,7 +35,7 @@ logger = get_system_logger() # CONSTANTS # ============================================================================= -_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_MEMORY_ROOT = Path(__file__).resolve().parents[2] def _find_repo_root() -> Path: diff --git a/src/aipass/memory/apps/handlers/dashboard_push.py b/src/aipass/memory/apps/handlers/dashboard_push.py index 40aa96b3..22d7588d 100644 --- a/src/aipass/memory/apps/handlers/dashboard_push.py +++ b/src/aipass/memory/apps/handlers/dashboard_push.py @@ -31,7 +31,7 @@ from aipass.memory.apps.handlers.json import json_handler logger = get_system_logger() # Resolve paths relative to handler location -_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_MEMORY_ROOT = Path(__file__).resolve().parents[2] # ============================================================================= diff --git a/src/aipass/memory/apps/handlers/learnings/manager.py b/src/aipass/memory/apps/handlers/learnings/manager.py index 6fd6f7e8..7fc7f85f 100644 --- a/src/aipass/memory/apps/handlers/learnings/manager.py +++ b/src/aipass/memory/apps/handlers/learnings/manager.py @@ -48,6 +48,15 @@ from aipass.memory.apps.handlers.json.memory_files import ( _MEMORY_ROOT = Path(__file__).resolve().parents[3] CHROMA_SUBPROCESS_SCRIPT = _MEMORY_ROOT / "apps" / "handlers" / "storage" / "chroma_subprocess.py" +def _find_repo_root() -> Path: + """Walk up from this file to find repo root (contains AIPASS_REGISTRY.json).""" + current = Path(__file__).resolve().parent + for parent in [current] + list(current.parents): + if (parent / "AIPASS_REGISTRY.json").exists(): + return parent + return Path.cwd() + + # Defaults DEFAULT_MAX_LEARNINGS = 100 DEFAULT_MAX_RECENTLY_COMPLETED = 20 @@ -973,7 +982,7 @@ def process_all_branches() -> Dict[str, Any]: Returns: Dict with processing summary """ - registry_path = Path.home() / "AIPASS_REGISTRY.json" + registry_path = _find_repo_root() / "AIPASS_REGISTRY.json" if not registry_path.exists(): return {'success': False, 'error': 'AIPASS_REGISTRY.json not found'} diff --git a/src/aipass/memory/apps/handlers/schema/normalize.py b/src/aipass/memory/apps/handlers/schema/normalize.py index d4e89667..3f086812 100644 --- a/src/aipass/memory/apps/handlers/schema/normalize.py +++ b/src/aipass/memory/apps/handlers/schema/normalize.py @@ -33,6 +33,15 @@ from aipass.memory.apps.handlers.json import json_handler logger = get_system_logger() +def _find_repo_root() -> Path: + """Walk up from this file to find repo root (contains AIPASS_REGISTRY.json).""" + current = Path(__file__).resolve().parent + for parent in [current] + list(current.parents): + if (parent / "AIPASS_REGISTRY.json").exists(): + return parent + return Path.cwd() + + def normalize_memory_file(file_path: Path, dry_run: bool = False) -> Dict[str, Any]: """ Normalize schema for a single memory file. @@ -148,7 +157,7 @@ def normalize_all_memory_files(dry_run: bool = False) -> Dict[str, Any]: Dict with statistics """ # Read registry - registry_path = Path.home() / "AIPASS_REGISTRY.json" + registry_path = _find_repo_root() / "AIPASS_REGISTRY.json" if not registry_path.exists(): return {'success': False, 'error': "AIPASS_REGISTRY.json not found"} diff --git a/src/aipass/memory/apps/handlers/symbolic/deduplicator.py b/src/aipass/memory/apps/handlers/symbolic/deduplicator.py index 215e77df..32cbb229 100644 --- a/src/aipass/memory/apps/handlers/symbolic/deduplicator.py +++ b/src/aipass/memory/apps/handlers/symbolic/deduplicator.py @@ -105,33 +105,23 @@ def deduplicate_fragment( # Build prompt and call LLM messages = _build_dedup_prompt(new_fragment, existing_fragments) - # Direct OpenRouter API call via urllib (no cross-branch imports needed) + # Direct OpenRouter API call via urllib import urllib.request import urllib.error - from pathlib import Path - api_key = None - _aipass_root = _MEMORY_ROOT.parent # memory/ -> aipass/ - for env_path in [ - Path.home() / ".secrets" / "aipass" / ".env", - _aipass_root / "api" / "apps" / ".env", - _aipass_root / "api" / ".env", - ]: - if env_path.exists(): - with open(env_path, encoding="utf-8") as f: - for line in f: - if line.strip().startswith("OPENROUTER_API_KEY="): - api_key = line.strip().split("=", 1)[1].strip().strip('"').strip("'") - break - if api_key: - break + # Load API key via api branch's key management + try: + from aipass.api.apps.handlers.auth.keys import get_api_key + api_key = get_api_key("openrouter") + except ImportError: + api_key = None if not api_key: return { 'success': True, 'action': 'ADD', 'fragment': new_fragment, - 'reason': 'No OpenRouter API key found, defaulting to ADD' + 'reason': 'No OpenRouter API key found (api branch unavailable or key missing), defaulting to ADD' } payload = json.dumps({ diff --git a/src/aipass/memory/apps/handlers/symbolic/extractor.py b/src/aipass/memory/apps/handlers/symbolic/extractor.py index 33eabd85..77cf63eb 100644 --- a/src/aipass/memory/apps/handlers/symbolic/extractor.py +++ b/src/aipass/memory/apps/handlers/symbolic/extractor.py @@ -348,31 +348,20 @@ def extract_fragments_llm(chat_history: List[Dict[str, Any]]) -> Dict[str, Any]: if not chat_history: return empty - # Direct OpenRouter API call via urllib (no cross-branch imports needed) + # Direct OpenRouter API call via urllib import urllib.request import urllib.error - from pathlib import Path - # Load API key from env file - api_key = None - _aipass_root = _MEMORY_ROOT.parent # memory/ -> aipass/ - for env_path in [ - Path.home() / ".secrets" / "aipass" / ".env", - _aipass_root / "api" / "apps" / ".env", - _aipass_root / "api" / ".env", - ]: - if env_path.exists(): - with open(env_path, encoding="utf-8") as f: - for line in f: - if line.strip().startswith("OPENROUTER_API_KEY="): - api_key = line.strip().split("=", 1)[1].strip().strip('"').strip("'") - break - if api_key: - break + # Load API key via api branch's key management + try: + from aipass.api.apps.handlers.auth.keys import get_api_key + api_key = get_api_key("openrouter") + except ImportError: + api_key = None if not api_key: return {'success': False, 'fragments': [], 'chunk_count': 0, - 'error': "No OpenRouter API key found in env files"} + 'error': "No OpenRouter API key found (api branch unavailable or key missing)"} chunks = _chunk_messages(chat_history) all_fragments = [] diff --git a/src/aipass/memory/tests/conftest.py b/src/aipass/memory/tests/conftest.py index e8b2c7f1..0932914c 100644 --- a/src/aipass/memory/tests/conftest.py +++ b/src/aipass/memory/tests/conftest.py @@ -112,22 +112,22 @@ def sample_memory_data() -> dict: def sample_registry_data() -> dict: """Provides sample AIPASS_REGISTRY.json data.""" return { - "branches": { - "test_branch": { + "branches": [ + { "name": "TEST_BRANCH", "path": "src/aipass/test_branch", "module": "aipass.test_branch", "email": "@test_branch", "status": "active" }, - "memory": { + { "name": "MEMORY", "path": "src/aipass/memory", "module": "aipass.memory", "email": "@memory", "status": "active" } - } + ] } diff --git a/src/aipass/prax/apps/handlers/logging/setup.py b/src/aipass/prax/apps/handlers/logging/setup.py index d717eec3..df4619e2 100755 --- a/src/aipass/prax/apps/handlers/logging/setup.py +++ b/src/aipass/prax/apps/handlers/logging/setup.py @@ -14,6 +14,7 @@ Handles dual logging (system-wide + branch-local) and terminal output. """ import logging +import threading logger = logging.getLogger(__name__) from pathlib import Path from typing import Dict, Optional @@ -40,6 +41,7 @@ from aipass.prax.apps.handlers.json import json_handler logger = logging.getLogger(__name__) _system_logger: Optional[logging.Logger] = None _captured_loggers: Dict[str, logging.Logger] = {} +_captured_loggers_lock = threading.Lock() _terminal_output_enabled = False # Try to import terminal handler support @@ -87,8 +89,9 @@ def setup_individual_logger( Returns: Configured logger instance """ - if module_name in _captured_loggers: - return _captured_loggers[module_name] + with _captured_loggers_lock: + if module_name in _captured_loggers: + return _captured_loggers[module_name] # Log new logger creation if _system_logger: @@ -146,8 +149,9 @@ def setup_individual_logger( terminal_handler = create_terminal_handler() # type: ignore[misc] logger.addHandler(terminal_handler) - # Store for reuse - _captured_loggers[module_name] = logger + # Store for reuse (thread-safe) + with _captured_loggers_lock: + _captured_loggers[module_name] = logger return logger diff --git a/src/aipass/prax/apps/handlers/monitoring/event_queue.py b/src/aipass/prax/apps/handlers/monitoring/event_queue.py index 84a6b471..464a6ef7 100644 --- a/src/aipass/prax/apps/handlers/monitoring/event_queue.py +++ b/src/aipass/prax/apps/handlers/monitoring/event_queue.py @@ -51,28 +51,28 @@ class MonitoringQueue: self.queue = PriorityQueue(maxsize=maxsize) self.recent_events = [] # For deduplication self.lock = threading.Lock() - self.running = True + self._stopped = threading.Event() def enqueue(self, event: MonitoringEvent) -> bool: """Add event to queue (thread-safe)""" - if not self.running: + if self._stopped.is_set(): return False json_handler.log_operation("event_queued", {"event_type": event.event_type, "branch": event.branch}) - # Simple deduplication - if not self._is_duplicate(event): + # Simple deduplication + queue put under single lock + with self.lock: + if self._is_duplicate(event): + return False try: self.queue.put(event, block=False) - with self.lock: - self.recent_events.append(event) - if len(self.recent_events) > 100: - self.recent_events.pop(0) + self.recent_events.append(event) + if len(self.recent_events) > 100: + self.recent_events.pop(0) return True except Exception as e: logger.warning(f"[event_queue] Failed to enqueue event (type={event.event_type}, branch={event.branch}): {e}") return False - return False def dequeue(self, timeout: float = 0.1) -> Optional[MonitoringEvent]: """Get next event from queue (thread-safe)""" @@ -93,19 +93,18 @@ class MonitoringQueue: def stop(self): """Stop accepting new events""" - self.running = False + self._stopped.set() self.flush() def _is_duplicate(self, event: MonitoringEvent) -> bool: - """Check if event duplicates recent event""" - with self.lock: - for recent in self.recent_events[-10:]: - if (recent.event_type == event.event_type and - recent.branch == event.branch and - recent.action == event.action and - recent.message == event.message and - abs((event.timestamp - recent.timestamp).total_seconds()) < 1): - return True + """Check if event duplicates recent event. Caller must hold self.lock.""" + for recent in self.recent_events[-10:]: + if (recent.event_type == event.event_type and + recent.branch == event.branch and + recent.action == event.action and + recent.message == event.message and + abs((event.timestamp - recent.timestamp).total_seconds()) < 1): + return True return False def size(self) -> int: diff --git a/src/aipass/prax/apps/modules/logger.py b/src/aipass/prax/apps/modules/logger.py index 37483036..5187cfb4 100755 --- a/src/aipass/prax/apps/modules/logger.py +++ b/src/aipass/prax/apps/modules/logger.py @@ -40,6 +40,7 @@ __all__ = [ import logging import sys +import threading from typing import Dict, Any # Stdlib logger for except-block compliance (seedgo requires variable named 'logger') @@ -95,10 +96,15 @@ class SystemLogger: """Auto-routing logger that writes to calling module's log file""" _watcher_started = False + _watcher_lock = threading.Lock() def _ensure_watcher(self): """Lazy-start file watchers on first logger use""" - if not SystemLogger._watcher_started: + if SystemLogger._watcher_started: + return + with SystemLogger._watcher_lock: + if SystemLogger._watcher_started: + return # Double-check after acquiring lock # Set flag FIRST to prevent recursion: trigger.fire() uses logger # internally, which would re-enter _ensure_watcher() before we return SystemLogger._watcher_started = True diff --git a/src/aipass/prax/apps/modules/monitor.py b/src/aipass/prax/apps/modules/monitor.py index 02f57e03..ae600b7e 100755 --- a/src/aipass/prax/apps/modules/monitor.py +++ b/src/aipass/prax/apps/modules/monitor.py @@ -114,9 +114,10 @@ def _refresh_pid_cache() -> None: global _pid_cache_last_refresh import time as _time now = _time.time() - if now - _pid_cache_last_refresh < _PID_CACHE_TTL: - return - _pid_cache_last_refresh = now + 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 @@ -149,7 +150,7 @@ def _get_pid_for_branch(branch: str) -> Optional[int]: # ============================================================================= # Global monitoring state -_monitoring_active = False +_stop_event = threading.Event() # Thread-safe shutdown signal _event_queue: Optional[MonitoringQueue] = None _module_tracker: Optional[ModuleTracker] = None _display_thread: Optional[threading.Thread] = None @@ -199,7 +200,7 @@ def handle_command(command: str, args: List[str]) -> bool: def _run_monitor(args: List[str]) -> bool: """Launch Mission Control live monitoring.""" - global _monitoring_active, _event_queue, _module_tracker + global _event_queue, _module_tracker global _display_thread, _file_watcher_thread, _log_watcher_thread json_handler.log_operation("monitor_started", {"args": args}) @@ -208,7 +209,7 @@ def _run_monitor(args: List[str]) -> bool: # Initialize monitoring subsystems _event_queue = MonitoringQueue() _module_tracker = ModuleTracker() - _monitoring_active = True + _stop_event.clear() _is_tty = sys.stdin.isatty() @@ -256,15 +257,17 @@ def _start_threads(): def _stop_threads(): """Stop all monitoring threads""" - global _monitoring_active, _event_queue + global _event_queue - _monitoring_active = False + _stop_event.set() if _event_queue: _event_queue.stop() - # Give threads time to finish - time.sleep(0.5) + # Join all daemon threads with timeout + for t in (_display_thread, _file_watcher_thread, _log_watcher_thread): + if t is not None and t.is_alive(): + t.join(timeout=2.0) logger.info("All monitoring threads stopped") @@ -287,9 +290,9 @@ def _render_event(event) -> None: def _display_worker(): """Display thread - pulls events from queue and displays them. No filtering.""" - global _monitoring_active, _event_queue + global _event_queue - while _monitoring_active: + while not _stop_event.is_set(): if not _event_queue: time.sleep(0.1) continue @@ -408,7 +411,7 @@ def _start_observer_with_fallback(handler, watch_dirs): def _file_watcher_worker(): """File watcher thread - watches filesystem changes and pushes to queue""" - global _monitoring_active, _event_queue + global _event_queue from aipass.prax.apps.handlers.monitoring.filesystem_handler import MonitoringFileHandler @@ -437,7 +440,7 @@ def _file_watcher_worker(): return try: - while _monitoring_active: + while not _stop_event.is_set(): time.sleep(0.1) finally: observer.stop() @@ -470,7 +473,7 @@ def _start_log_watcher_with_fallback(event_queue) -> bool: def _log_watcher_worker(): """Log watcher thread - uses proper log_watcher.py with all improvements""" - global _monitoring_active, _event_queue + global _event_queue from aipass.prax.apps.handlers.monitoring.log_watcher import stop_log_watcher @@ -482,7 +485,7 @@ def _log_watcher_worker(): return try: - while _monitoring_active: + while not _stop_event.is_set(): time.sleep(0.1) finally: stop_log_watcher() @@ -502,13 +505,13 @@ def _handle_interactive_cmd(cmd: str, get_help_text) -> None: def _interactive_loop(): """Interactive command loop - handles user input, or passive loop if no TTY""" - global _monitoring_active + global _event_queue # Non-TTY mode: just keep alive if not sys.stdin.isatty(): logger.info("[monitor] No TTY detected - passive mode (Ctrl+C to stop)") try: - while _monitoring_active: + while not _stop_event.is_set(): time.sleep(0.5) except KeyboardInterrupt: logger.info("[monitor] Stopped by user (passive mode)") @@ -520,7 +523,7 @@ def _interactive_loop(): get_help_text ) - while _monitoring_active: + while not _stop_event.is_set(): try: user_input = input().strip() if not user_input: diff --git a/src/aipass/prax/tests/test_event_queue.py b/src/aipass/prax/tests/test_event_queue.py index 8e076041..ff1044d5 100644 --- a/src/aipass/prax/tests/test_event_queue.py +++ b/src/aipass/prax/tests/test_event_queue.py @@ -454,5 +454,5 @@ class TestMonitoringQueue: q = MonitoringQueue() q.stop() q.stop() # Should not raise - assert q.running is False + assert q._stopped.is_set() assert q.size() == 0