feat(system): S84: CRITICAL fixes — memory (path+security), ai_mail (locking), prax (thread safety), global prompt patrick-file rule
Co-Authored-By: @devpulse <devpulse@aipass>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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]
|
||||
|
||||
|
||||
# =============================================================================
|
||||
|
||||
@@ -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'}
|
||||
|
||||
@@ -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"}
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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 = []
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user