From 717acf078a983489126e41a7d2e25512a9d04aeb Mon Sep 17 00:00:00 2001 From: AIPass Date: Thu, 12 Mar 2026 19:29:06 -0700 Subject: [PATCH] feat(system): hook containment + memory bank operational (FPLAN-0026) (#35) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Hook containment: AIPass hooks moved from global to project-level .claude/settings.json. Global now only has trinity inject. Non-AIPass projects no longer receive AIPass context. Memory bank (FPLAN-0026): @memory completed all 6 phases autonomously — venv with ML deps, registry fix, rollover E2E, search, plan archival, help cleanup. ChromaDB operational with local + global .chroma stores. Co-authored-by: Claude Opus 4.6 --- .claude/settings.json | 60 +++- src/aipass/memory/.gitignore | 1 + .../memory/apps/handlers/dashboard_push.py | 29 +- .../memory/apps/handlers/intake/__init__.py | 0 .../apps/handlers/intake/plans_processor.py | 324 ++++++++++++++++++ .../memory/apps/handlers/monitor/detector.py | 47 ++- .../apps/handlers/monitor/memory_watcher.py | 37 +- .../apps/handlers/rollover/extractor.py | 38 +- .../apps/handlers/rollover/orchestrator.py | 65 +++- .../apps/handlers/search/query_executor.py | 120 ++++--- .../handlers/storage/chroma_subprocess.py | 161 ++++++++- .../apps/handlers/vector/embed_subprocess.py | 80 +++++ src/aipass/memory/apps/memory.py | 19 +- src/aipass/memory/apps/modules/rollover.py | 2 +- .../memory/config/.plans_processed.json | 3 + .../memory/config/memory_bank.config.json | 15 + 16 files changed, 889 insertions(+), 112 deletions(-) create mode 100644 src/aipass/memory/apps/handlers/intake/__init__.py create mode 100644 src/aipass/memory/apps/handlers/intake/plans_processor.py create mode 100644 src/aipass/memory/apps/handlers/vector/embed_subprocess.py create mode 100644 src/aipass/memory/config/.plans_processed.json create mode 100644 src/aipass/memory/config/memory_bank.config.json diff --git a/.claude/settings.json b/.claude/settings.json index bc2214f4..44853105 100644 --- a/.claude/settings.json +++ b/.claude/settings.json @@ -10,5 +10,63 @@ ], "defaultMode": "acceptEdits" }, - "skipDangerousModePermissionPrompt": true + "skipDangerousModePermissionPrompt": true, + "hooks": { + "UserPromptSubmit": [ + { + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh global_prompt && cat /home/patrick/Projects/AIPass/.aipass/aipass_global_prompt.md 2>/dev/null || true" + } + ] + }, + { + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh branch_prompt && python3 /home/patrick/Projects/AIPass/.claude/hooks/branch_prompt_loader.py" + } + ] + }, + { + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh identity && python3 /home/patrick/Projects/AIPass/.claude/hooks/identity_injector.py" + } + ] + }, + { + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh email_check && python3 /home/patrick/Projects/AIPass/.claude/hooks/email_notification.py" + } + ] + } + ], + "PreCompact": [ + { + "matcher": "manual", + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh pre_compact && python3 /home/patrick/Projects/AIPass/.claude/hooks/pre_compact.py", + "timeout": 60 + } + ] + }, + { + "matcher": "auto", + "hooks": [ + { + "type": "command", + "command": "/home/patrick/.claude/hook_logger.sh pre_compact && python3 /home/patrick/Projects/AIPass/.claude/hooks/pre_compact.py", + "timeout": 60 + } + ] + } + ] + } } diff --git a/src/aipass/memory/.gitignore b/src/aipass/memory/.gitignore index 9cf1dfc4..e0a8e539 100644 --- a/src/aipass/memory/.gitignore +++ b/src/aipass/memory/.gitignore @@ -2,6 +2,7 @@ __pycache__/ *.pyc *.pyo .env +.venv/ *.egg-info/ .coverage htmlcov/ diff --git a/src/aipass/memory/apps/handlers/dashboard_push.py b/src/aipass/memory/apps/handlers/dashboard_push.py index 96b4962d..173d0d04 100644 --- a/src/aipass/memory/apps/handlers/dashboard_push.py +++ b/src/aipass/memory/apps/handlers/dashboard_push.py @@ -38,10 +38,21 @@ _MEMORY_ROOT = Path(__file__).resolve().parents[3] # ============================================================================= CENTRAL_FILE = _MEMORY_ROOT / "central" / "memory_bank.central.json" -AIPASS_REGISTRY = Path.home() / "AIPASS_REGISTRY.json" CONFIG_PATH = _MEMORY_ROOT / "config" / "memory_bank.config.json" TEMPLATE_VERSION_FILE = _MEMORY_ROOT / "templates" / ".template_version.json" + +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() + + +AIPASS_REGISTRY = _find_repo_root() / "AIPASS_REGISTRY.json" + # Near-rollover threshold: branches with fewer than this many lines remaining NEAR_ROLLOVER_THRESHOLD = 100 @@ -154,18 +165,22 @@ def _find_branches_near_rollover() -> List[Dict[str, Any]]: branches = registry.get("branches", []) rollover_config = _get_rollover_config() + repo_root = _find_repo_root() for branch in branches: branch_name = branch.get("name", "") - branch_path = Path(branch.get("path", "")) + raw_path = branch.get("path", "") + branch_path = Path(raw_path) + if not branch_path.is_absolute(): + branch_path = repo_root / raw_path if not branch_path.exists(): continue max_lines = _get_max_lines_for_branch(branch_name, rollover_config) - # Check both .local.json and .observations.json + # Memory files live in .trinity/ subdirectory for suffix in ["local", "observations"]: - memory_file = branch_path / f"{branch_name}.{suffix}.json" + memory_file = branch_path / ".trinity" / f"{suffix}.json" if not memory_file.exists(): continue @@ -283,9 +298,13 @@ def _get_all_branch_paths() -> List[Path]: return [] registry = json_loads(AIPASS_REGISTRY.read_text(encoding="utf-8")) + repo_root = _find_repo_root() paths = [] for branch in registry.get("branches", []): - branch_path = Path(branch.get("path", "")) + raw_path = branch.get("path", "") + branch_path = Path(raw_path) + if not branch_path.is_absolute(): + branch_path = repo_root / raw_path if branch_path.exists(): paths.append(branch_path) return paths diff --git a/src/aipass/memory/apps/handlers/intake/__init__.py b/src/aipass/memory/apps/handlers/intake/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/src/aipass/memory/apps/handlers/intake/plans_processor.py b/src/aipass/memory/apps/handlers/intake/plans_processor.py new file mode 100644 index 00000000..0952f9e3 --- /dev/null +++ b/src/aipass/memory/apps/handlers/intake/plans_processor.py @@ -0,0 +1,324 @@ +# =================== AIPass ==================== +# Name: plans_processor.py +# Description: Plan Archival Vectorization Handler +# Version: 1.0.0 +# Created: 2026-03-12 +# Modified: 2026-03-12 +# ============================================= + +""" +Plan Archival Vectorization Handler + +Reads closed plan files from flow/processed_plans/, chunks them, +generates embeddings via subprocess, and stores vectors in ChromaDB. + +Called by memory_watcher._check_plans() during watch mode +and can be invoked directly for manual processing. + +Uses subprocess pattern for ML operations (memory venv isolation). +""" + +import json +import os +import subprocess +import sys +from pathlib import Path +from datetime import datetime +from typing import Dict, Any, List + +from aipass.prax import logger + +# Subprocess scripts +_HANDLERS_DIR = Path(__file__).resolve().parent.parent +EMBED_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "vector" / "embed_subprocess.py" +CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py" + +# Memory venv python +_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_MEMORY_VENV_PYTHON = _MEMORY_ROOT / ".venv" / "bin" / "python" + + +def _find_repo_root() -> Path: + """Walk up from this file to find repo root.""" + current = Path(__file__).resolve().parent + for parent in [current] + list(current.parents): + if (parent / "AIPASS_REGISTRY.json").exists(): + return parent + return Path.cwd() + + +def _get_memory_python() -> str: + env_override = os.environ.get("AIPASS_MEMORY_PYTHON") + if env_override: + return env_override + if _MEMORY_VENV_PYTHON.exists(): + return str(_MEMORY_VENV_PYTHON) + return sys.executable + + +MEMORY_PYTHON = _get_memory_python() + +# Track which files have been processed +_PROCESSED_MANIFEST = _MEMORY_ROOT / "config" / ".plans_processed.json" + +# Chunk settings +MAX_CHUNK_CHARS = 1500 # ~375 tokens, fits well with all-MiniLM-L6-v2 + + +# ============================================================================= +# CHUNKING +# ============================================================================= + +def _chunk_plan_text(text: str, filename: str) -> List[Dict[str, str]]: + """ + Chunk plan text into sections for vectorization. + + Splits on markdown headers (## / ###), with fallback to paragraph splitting. + Each chunk gets metadata about its source. + + Args: + text: Full plan text + filename: Source filename for metadata + + Returns: + List of dicts with 'text' and 'section' keys + """ + chunks = [] + + # Split by markdown headers + lines = text.split('\n') + current_section = filename + current_lines = [] + + for line in lines: + if line.startswith('## ') or line.startswith('### '): + # Flush previous section + if current_lines: + section_text = '\n'.join(current_lines).strip() + if section_text and len(section_text) > 30: + chunks.append({'text': section_text, 'section': current_section}) + current_section = line.lstrip('#').strip() + current_lines = [line] + else: + current_lines.append(line) + + # Flush last section + if current_lines: + section_text = '\n'.join(current_lines).strip() + if section_text and len(section_text) > 30: + chunks.append({'text': section_text, 'section': current_section}) + + # If no headers found, chunk by size + if not chunks: + full_text = text.strip() + if len(full_text) > MAX_CHUNK_CHARS: + for i in range(0, len(full_text), MAX_CHUNK_CHARS): + chunk_text = full_text[i:i + MAX_CHUNK_CHARS].strip() + if chunk_text and len(chunk_text) > 30: + chunks.append({'text': chunk_text, 'section': f'{filename}_part{i // MAX_CHUNK_CHARS}'}) + elif len(full_text) > 30: + chunks.append({'text': full_text, 'section': filename}) + + # Split oversized chunks + final_chunks = [] + for chunk in chunks: + if len(chunk['text']) > MAX_CHUNK_CHARS * 2: + text_content = chunk['text'] + for i in range(0, len(text_content), MAX_CHUNK_CHARS): + part = text_content[i:i + MAX_CHUNK_CHARS].strip() + if part and len(part) > 30: + final_chunks.append({ + 'text': part, + 'section': f"{chunk['section']}_part{i // MAX_CHUNK_CHARS}" + }) + else: + final_chunks.append(chunk) + + return final_chunks + + +# ============================================================================= +# PROCESSED MANIFEST +# ============================================================================= + +def _load_manifest() -> Dict[str, str]: + """Load processed files manifest.""" + if _PROCESSED_MANIFEST.exists(): + try: + return json.loads(_PROCESSED_MANIFEST.read_text(encoding='utf-8')) + except Exception: + return {} + return {} + + +def _save_manifest(manifest: Dict[str, str]) -> None: + """Save processed files manifest.""" + _PROCESSED_MANIFEST.parent.mkdir(parents=True, exist_ok=True) + _PROCESSED_MANIFEST.write_text(json.dumps(manifest, indent=2), encoding='utf-8') + + +# ============================================================================= +# SUBPROCESS WRAPPERS +# ============================================================================= + +def _embed_texts(texts: List[str]) -> dict: + """Encode texts via subprocess.""" + input_data = json.dumps({'texts': texts}) + try: + result = subprocess.run( + [str(MEMORY_PYTHON), str(EMBED_SUBPROCESS_SCRIPT)], + input=input_data, + capture_output=True, text=True, timeout=120 + ) + if result.returncode != 0: + return {'success': False, 'error': result.stderr or 'Embedding failed'} + return json.loads(result.stdout) + except Exception as e: + return {'success': False, 'error': str(e)} + + +def _store_vectors(embeddings, documents, metadatas, collection_name="flow_plans") -> dict: + """Store vectors via subprocess.""" + input_data = { + 'operation': 'store_vectors', + 'branch': 'FLOW', + 'memory_type': collection_name, + 'embeddings': embeddings, + 'documents': documents, + 'metadatas': metadatas, + 'db_path': None # global + } + try: + result = subprocess.run( + [str(MEMORY_PYTHON), str(CHROMA_SUBPROCESS_SCRIPT)], + input=json.dumps(input_data), + capture_output=True, text=True, timeout=60 + ) + if result.returncode != 0: + return {'success': False, 'error': result.stderr or 'Storage failed'} + return json.loads(result.stdout) + except Exception as e: + return {'success': False, 'error': str(e)} + + +# ============================================================================= +# PUBLIC API +# ============================================================================= + +def process_plans() -> Dict[str, Any]: + """ + Process plan files from flow/processed_plans/ into vector storage. + + Workflow: + 1. Load config to find plans directory + 2. Scan for unprocessed .md files + 3. Chunk each file into sections + 4. Embed all chunks via subprocess + 5. Store vectors in ChromaDB + 6. Update processed manifest + + Returns: + Dict with success, files_processed, total_chunks + """ + # Load config + config_path = _MEMORY_ROOT / "config" / "memory_bank.config.json" + try: + config = json.loads(config_path.read_text(encoding='utf-8')) + plans_config = config.get('plans', {}) + except Exception as e: + return {'success': False, 'error': f'Config load failed: {e}'} + + if not plans_config.get('enabled', False): + return {'success': True, 'skipped': True, 'reason': 'plans disabled'} + + # Resolve plans directory (relative to repo root) + plans_dir = plans_config.get('path', 'src/aipass/flow/processed_plans') + repo_root = _find_repo_root() + plans_path = Path(plans_dir) if Path(plans_dir).is_absolute() else repo_root / plans_dir + extensions = plans_config.get('supported_extensions', ['.md']) + collection_name = plans_config.get('collection_name', 'flow_plans') + + if not plans_path.exists(): + return {'success': True, 'files_processed': 0, 'total_chunks': 0, 'reason': 'plans dir not found'} + + # Get plan files + files = [] + for ext in extensions: + files.extend(plans_path.glob(f'*{ext}')) + + if not files: + return {'success': True, 'files_processed': 0, 'total_chunks': 0} + + # Load manifest to skip already-processed files + manifest = _load_manifest() + unprocessed = [f for f in files if f.name not in manifest] + + if not unprocessed: + return {'success': True, 'files_processed': 0, 'total_chunks': 0, 'reason': 'all files already processed'} + + logger.info(f"[plans] Found {len(unprocessed)} unprocessed plan files") + + total_chunks = 0 + files_processed = 0 + errors = [] + + for plan_file in unprocessed: + try: + text = plan_file.read_text(encoding='utf-8') + except Exception as e: + errors.append(f'{plan_file.name}: read error: {e}') + continue + + # Chunk the plan + chunks = _chunk_plan_text(text, plan_file.name) + if not chunks: + manifest[plan_file.name] = datetime.now().isoformat() + continue + + # Extract texts and build metadata + texts = [c['text'] for c in chunks] + metadatas = [ + { + 'source_file': plan_file.name, + 'section': c['section'], + 'processed_at': datetime.now().isoformat(), + 'type': 'plan' + } + for c in chunks + ] + + # Embed + embed_result = _embed_texts(texts) + if not embed_result.get('success'): + errors.append(f"{plan_file.name}: embed error: {embed_result.get('error')}") + continue + + embeddings = embed_result.get('embeddings', []) + if not embeddings: + errors.append(f'{plan_file.name}: no embeddings returned') + continue + + # Store + store_result = _store_vectors(embeddings, texts, metadatas, collection_name) + if not store_result.get('success'): + errors.append(f"{plan_file.name}: store error: {store_result.get('error')}") + continue + + # Mark as processed + manifest[plan_file.name] = datetime.now().isoformat() + files_processed += 1 + total_chunks += len(chunks) + logger.info(f"[plans] Processed {plan_file.name}: {len(chunks)} chunks vectorized") + + # Save manifest + _save_manifest(manifest) + + result = { + 'success': files_processed > 0 or not errors, + 'files_processed': files_processed, + 'total_chunks': total_chunks, + } + if errors: + result['errors'] = errors + + return result diff --git a/src/aipass/memory/apps/handlers/monitor/detector.py b/src/aipass/memory/apps/handlers/monitor/detector.py index 651643c5..ab9b66c1 100644 --- a/src/aipass/memory/apps/handlers/monitor/detector.py +++ b/src/aipass/memory/apps/handlers/monitor/detector.py @@ -34,6 +34,18 @@ logger = get_system_logger() # No module imports (handler independence) +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() + + +_REPO_ROOT = _find_repo_root() + + # ============================================================================= # DATA STRUCTURES # ============================================================================= @@ -61,12 +73,14 @@ class RolloverTrigger: def _read_registry() -> List[Dict[str, Any]]: """ - Read AIPASS_REGISTRY.json + Read AIPASS_REGISTRY.json from repo root. + + Registry paths are relative — resolved against repo root. Returns: - List of branch dictionaries + List of branch dictionaries with absolute paths """ - registry_path = Path.home() / "AIPASS_REGISTRY.json" + registry_path = _REPO_ROOT / "AIPASS_REGISTRY.json" if not registry_path.exists(): return [] @@ -74,18 +88,30 @@ def _read_registry() -> List[Dict[str, Any]]: try: with open(registry_path, 'r', encoding='utf-8') as f: data = json.load(f) - return data.get('branches', []) - except Exception as e: - # No logging in handlers - let caller handle errors + branches = data.get('branches', []) + + # Resolve relative paths against repo root + for branch in branches: + raw_path = branch.get('path', '') + resolved = Path(raw_path) + if not resolved.is_absolute(): + resolved = _REPO_ROOT / raw_path + branch['path'] = str(resolved) + + return branches + except Exception: return [] def _get_memory_file_path(branch: Dict, memory_type: str) -> Path | None: """ - Get path to memory file for branch + Get path to memory file for branch. + + Memory files live in .trinity/ subdirectory as {memory_type}.json + (e.g., .trinity/local.json, .trinity/observations.json). Args: - branch: Branch dict from registry + branch: Branch dict from registry (path already resolved to absolute) memory_type: 'observations' or 'local' Returns: @@ -95,9 +121,8 @@ def _get_memory_file_path(branch: Dict, memory_type: str) -> Path | None: if not branch_path.exists(): return None - branch_name = branch.get('name', '').upper() - file_name = f"{branch_name}.{memory_type}.json" - file_path = branch_path / file_name + # Memory files are in .trinity/ subdirectory + file_path = branch_path / '.trinity' / f'{memory_type}.json' return file_path if file_path.exists() else None diff --git a/src/aipass/memory/apps/handlers/monitor/memory_watcher.py b/src/aipass/memory/apps/handlers/monitor/memory_watcher.py index d79a7dbd..5c0aa44b 100644 --- a/src/aipass/memory/apps/handlers/monitor/memory_watcher.py +++ b/src/aipass/memory/apps/handlers/monitor/memory_watcher.py @@ -160,11 +160,12 @@ def check_and_rollover() -> Dict[str, Any]: # Extract branch name from path (last component, uppercase) branch_name = branch.name.upper() - # Find memory files in this branch (skip DASHBOARD files - no rollover metadata) - for pattern in ['*.local.json', '*.observations.json']: - for memory_file in branch.glob(pattern): - if memory_file.name.startswith('DASHBOARD'): - continue + # Find memory files in .trinity/ subdirectory + trinity_dir = branch / '.trinity' + if not trinity_dir.exists(): + continue + for pattern in ['local.json', 'observations.json']: + for memory_file in trinity_dir.glob(pattern): results['files_checked'] += 1 @@ -360,16 +361,28 @@ def _check_code_archive() -> Dict[str, Any]: # UTILITY FUNCTIONS # ============================================================================= +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 _get_branch_paths() -> list[Path]: """ - Get all branch paths from AIPASS_REGISTRY.json (silent - no logging) + Get all branch paths from AIPASS_REGISTRY.json (silent - no logging). + + Registry paths are relative — resolved against repo root. Returns: List of Path objects for each branch """ import json - registry_path = Path.home() / "AIPASS_REGISTRY.json" + repo_root = _find_repo_root() + registry_path = repo_root / "AIPASS_REGISTRY.json" if not registry_path.exists(): return [] @@ -381,7 +394,10 @@ def _get_branch_paths() -> list[Path]: paths = [] for branch in branches: - branch_path = Path(branch.get('path', '')) + raw_path = branch.get('path', '') + branch_path = Path(raw_path) + if not branch_path.is_absolute(): + branch_path = repo_root / raw_path if branch_path.exists(): paths.append(branch_path) @@ -392,7 +408,7 @@ def _get_branch_paths() -> list[Path]: def _is_memory_file(file_path: Path) -> bool: """ - Check if file is a memory file (*.local.json or *.observations.json) + Check if file is a memory file in .trinity/ (local.json or observations.json) Args: file_path: Path to check @@ -401,7 +417,8 @@ def _is_memory_file(file_path: Path) -> bool: True if memory file, False otherwise """ name = file_path.name - return (name.endswith('.local.json') or name.endswith('.observations.json')) + parent = file_path.parent.name + return parent == '.trinity' and name in ('local.json', 'observations.json') # ============================================================================= diff --git a/src/aipass/memory/apps/handlers/rollover/extractor.py b/src/aipass/memory/apps/handlers/rollover/extractor.py index 37c2055e..63726f64 100644 --- a/src/aipass/memory/apps/handlers/rollover/extractor.py +++ b/src/aipass/memory/apps/handlers/rollover/extractor.py @@ -141,6 +141,31 @@ def _count_file_lines(file_path: Path) -> int: return len(f.readlines()) +# ============================================================================= +# PATH HELPERS +# ============================================================================= + +def _derive_branch_and_type(file_path: Path) -> tuple[str, str]: + """ + Derive branch name and memory type from file path. + + Handles both naming conventions: + - .trinity/local.json → branch from parent.parent.name, type from stem + - BRANCH.type.json (legacy) → parsed from stem + + Returns: + Tuple of (branch_name, memory_type) e.g. ("DEVPULSE", "local") + """ + if file_path.parent.name == '.trinity': + branch_name = file_path.parent.parent.name.upper() + memory_type = file_path.stem # "local" or "observations" + else: + parts = file_path.stem.split('.') + branch_name = parts[0] if len(parts) > 0 else "UNKNOWN" + memory_type = parts[1] if len(parts) > 1 else "unknown" + return branch_name, memory_type + + # ============================================================================= # STRUCTURE DETECTION # ============================================================================= @@ -307,10 +332,9 @@ def _extract_items_v2(file_path: Path, data: Dict[str, Any]) -> Dict[str, Any]: 'error': f"Failed to write file: {e}" } - # Parse branch and type from filename - parts = file_path.stem.split('.') - branch_name = parts[0] if len(parts) > 0 else "UNKNOWN" - memory_type = parts[1] if len(parts) > 1 else "unknown" + # Derive branch and type from path + # .trinity/local.json → branch = parent.parent.name, type = stem + branch_name, memory_type = _derive_branch_and_type(file_path) return { 'success': True, @@ -420,10 +444,8 @@ def extract_items( 'error': f"Failed to write file: {e}" } - # Parse branch and type from filename - parts = file_path.stem.split('.') - branch_name = parts[0] if len(parts) > 0 else "UNKNOWN" - memory_type = parts[1] if len(parts) > 1 else "unknown" + # Derive branch and type from path + branch_name, memory_type = _derive_branch_and_type(file_path) return { 'success': True, diff --git a/src/aipass/memory/apps/handlers/rollover/orchestrator.py b/src/aipass/memory/apps/handlers/rollover/orchestrator.py index 04d97e18..cbce635b 100644 --- a/src/aipass/memory/apps/handlers/rollover/orchestrator.py +++ b/src/aipass/memory/apps/handlers/rollover/orchestrator.py @@ -32,15 +32,27 @@ from aipass.prax import logger # Handler imports (relative within the memory package) from aipass.memory.apps.handlers.monitor import detector from aipass.memory.apps.handlers.rollover import extractor -from aipass.memory.apps.handlers.vector import embedder from aipass.memory.apps.handlers.tracking import line_counter -# ChromaDB storage via subprocess +# Subprocess scripts for ML operations (run in memory venv) _HANDLERS_DIR = Path(__file__).resolve().parent.parent CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py" +EMBED_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "vector" / "embed_subprocess.py" -# Use system python by default; can be overridden via environment variable -MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable) +# Memory venv python — auto-detect from memory/.venv/ or use env var override +_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_MEMORY_VENV_PYTHON = _MEMORY_ROOT / ".venv" / "bin" / "python" + +def _get_memory_python() -> str: + """Get the Python executable for memory ML operations.""" + env_override = os.environ.get("AIPASS_MEMORY_PYTHON") + if env_override: + return env_override + if _MEMORY_VENV_PYTHON.exists(): + return str(_MEMORY_VENV_PYTHON) + return sys.executable + +MEMORY_PYTHON = _get_memory_python() # ============================================================================= @@ -118,6 +130,43 @@ def store_vectors_subprocess(branch: str, memory_type: str, embeddings: list, return {'success': False, 'error': str(e)} +# ============================================================================= +# EMBEDDING VIA SUBPROCESS +# ============================================================================= + +def encode_batch_subprocess(texts: list) -> dict: + """ + Encode texts via subprocess using memory venv's sentence-transformers. + + Args: + texts: List of text strings to encode + + Returns: + Dict with success, embeddings, count, dimension + """ + input_data = json.dumps({'texts': texts}) + + try: + result = subprocess.run( + [str(MEMORY_PYTHON), str(EMBED_SUBPROCESS_SCRIPT)], + input=input_data, + capture_output=True, + text=True, + timeout=120 + ) + + if result.returncode != 0: + return {'success': False, 'error': result.stderr or 'Embedding subprocess failed'} + + return json.loads(result.stdout) + except subprocess.TimeoutExpired: + return {'success': False, 'error': 'Embedding timed out'} + except json.JSONDecodeError as e: + return {'success': False, 'error': f'Invalid JSON from embedder: {e}'} + except Exception as e: + return {'success': False, 'error': str(e)} + + # ============================================================================= # PATH HELPERS # ============================================================================= @@ -280,8 +329,8 @@ def execute_rollover() -> Dict[str, Any]: continue memories = extract_result.get('entries', []) - branch = extract_result.get('branch', '') - memory_type = extract_result.get('type', 'unknown') + branch = extract_result.get('branch', '') or trigger.branch + memory_type = extract_result.get('type', 'unknown') or trigger.memory_type old_lines = extract_result.get('old_lines', 0) new_lines = extract_result.get('new_lines', 0) @@ -295,8 +344,8 @@ def execute_rollover() -> Dict[str, Any]: # Convert memory items to text for vectorization texts = extract_text_from_memories(memories) - # Step 3: Generate embeddings - embed_result = embedder.encode_batch(texts) + # Step 3: Generate embeddings (via subprocess in memory venv) + embed_result = encode_batch_subprocess(texts) if not embed_result['success']: error_msg = embed_result.get('error', 'Unknown error') diff --git a/src/aipass/memory/apps/handlers/search/query_executor.py b/src/aipass/memory/apps/handlers/search/query_executor.py index 409b87c3..0fcb46e6 100644 --- a/src/aipass/memory/apps/handlers/search/query_executor.py +++ b/src/aipass/memory/apps/handlers/search/query_executor.py @@ -1,16 +1,16 @@ # =================== AIPass ==================== # Name: query_executor.py # Description: Search Query Execution Handler -# Version: 1.0.0 +# Version: 1.1.0 # Created: 2026-03-08 -# Modified: 2026-03-08 +# Modified: 2026-03-12 # ============================================= """ Search Query Execution Handler Contains the core search execution logic: subprocess-based vector search, -query encoding via embedder, similarity calculation, and result filtering. +query encoding via subprocess embedder, similarity calculation, and result filtering. Called by the search module which handles display/CLI concerns. Purpose: @@ -27,20 +27,81 @@ from typing import Dict, Any from aipass.prax import logger -# Handler imports -from aipass.memory.apps.handlers.vector import embedder - -# ChromaDB search via subprocess +# Subprocess scripts for ML operations (run in memory venv) _HANDLERS_DIR = Path(__file__).resolve().parent.parent CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py" +EMBED_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "vector" / "embed_subprocess.py" -# Use system python by default; can be overridden via environment variable -MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable) +# Memory venv python — auto-detect from memory/.venv/ or use env var override +_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_MEMORY_VENV_PYTHON = _MEMORY_ROOT / ".venv" / "bin" / "python" + + +def _get_memory_python() -> str: + """Get the Python executable for memory ML operations.""" + env_override = os.environ.get("AIPASS_MEMORY_PYTHON") + if env_override: + return env_override + if _MEMORY_VENV_PYTHON.exists(): + return str(_MEMORY_VENV_PYTHON) + return sys.executable + + +MEMORY_PYTHON = _get_memory_python() # Minimum similarity threshold - filter out irrelevant results MIN_SIMILARITY_THRESHOLD = 0.40 # 40% minimum relevance +# ============================================================================= +# SUBPROCESS EMBEDDING +# ============================================================================= + +def encode_query_subprocess(query: str) -> dict: + """ + Encode query text via subprocess using memory venv's sentence-transformers. + + Args: + query: Search query text + + Returns: + Dict with success, embedding (list of floats), dimension + """ + input_data = json.dumps({'texts': [query]}) + + try: + result = subprocess.run( + [str(MEMORY_PYTHON), str(EMBED_SUBPROCESS_SCRIPT)], + input=input_data, + capture_output=True, + text=True, + timeout=120 + ) + + if result.returncode != 0: + return {'success': False, 'error': result.stderr or 'Embedding subprocess failed'} + + data = json.loads(result.stdout) + if not data.get('success'): + return data + + embeddings = data.get('embeddings', []) + if not embeddings: + return {'success': False, 'error': 'No embedding generated'} + + return { + 'success': True, + 'embedding': embeddings[0], + 'dimension': data.get('dimension', 384) + } + except subprocess.TimeoutExpired: + return {'success': False, 'error': 'Embedding timed out'} + except json.JSONDecodeError as e: + return {'success': False, 'error': f'Invalid JSON from embedder: {e}'} + except Exception as e: + return {'success': False, 'error': str(e)} + + # ============================================================================= # SUBPROCESS VECTOR SEARCH # ============================================================================= @@ -55,8 +116,6 @@ def search_vectors_subprocess( """ Search vectors via subprocess. - This ensures ChromaDB compatibility regardless of calling Python version. - Args: query_embedding: Query embedding vector (list of floats) branch: Optional branch filter @@ -103,12 +162,12 @@ def search_vectors_subprocess( def _calculate_similarity(distance: float) -> float: """ - Calculate similarity from ChromaDB L2 distance. + Calculate similarity from ChromaDB cosine distance. - ChromaDB L2 distance: 0=identical, ~2=very different. + ChromaDB cosine distance: 0=identical, 2=opposite. Args: - distance: L2 distance from ChromaDB + distance: Cosine distance from ChromaDB Returns: Similarity score between 0 and 1 @@ -120,8 +179,6 @@ def _filter_results(results: list, n_results: int) -> list: """ Filter search results by similarity threshold and quality. - Removes empty documents and results below the minimum similarity threshold. - Args: results: Raw search results from subprocess n_results: Maximum number of results to return @@ -136,7 +193,6 @@ def _filter_results(results: list, n_results: int) -> list: similarity = _calculate_similarity(distance) - # Skip empty documents and low-relevance results if not document or not document.strip(): continue if similarity < MIN_SIMILARITY_THRESHOLD: @@ -162,7 +218,7 @@ def execute_search( Execute semantic search: encode query, search vectors, filter results. Workflow: - 1. Encode query to embedding vector via embedder handler + 1. Encode query to embedding vector via subprocess (memory venv) 2. Search ChromaDB via subprocess 3. Filter and score results by similarity @@ -173,18 +229,10 @@ def execute_search( n_results: Number of results to return Returns: - Dict with: - - success: bool - - query: original query text - - branch: branch filter (if any) - - memory_type: memory type filter (if any) - - results: list of filtered result dicts with similarity scores - - collections_searched: number of collections searched - - total_results: total raw results before filtering - - error: error message (on failure) + Dict with success, results, collections_searched, total_results """ - # Step 1: Encode query - embed_result = embedder.encode_batch([query]) + # Step 1: Encode query via subprocess + embed_result = encode_query_subprocess(query) if not embed_result['success']: error_msg = embed_result.get('error', 'Unknown error') @@ -195,19 +243,7 @@ def execute_search( 'query': query, } - embeddings = embed_result.get('embeddings', []) - if not embeddings: - return { - 'success': False, - 'error': 'No embedding generated', - 'query': query, - } - - query_embedding = embeddings[0] - # Convert numpy array to list for JSON serialization - if hasattr(query_embedding, 'tolist'): - query_embedding = query_embedding.tolist() - + query_embedding = embed_result['embedding'] logger.info(f"[search] Encoded query to {len(query_embedding)}-dim vector") # Step 2: Search via subprocess diff --git a/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py b/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py index 811d43b2..45eecc4e 100755 --- a/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py +++ b/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py @@ -1,16 +1,18 @@ # =================== AIPass ==================== # Name: chroma_subprocess.py # Description: ChromaDB Subprocess Handler -# Version: 1.1.0 +# Version: 1.2.0 # Created: 2025-11-27 -# Modified: 2026-03-06 +# Modified: 2026-03-12 # ============================================= """ ChromaDB Subprocess Handler -Called via subprocess from rollover module to ensure ChromaDB operations -run in an isolated process. +Called via subprocess from rollover orchestrator to ensure ChromaDB operations +run in the memory-specific venv (AIPASS_MEMORY_PYTHON). + +Self-contained — does NOT import from aipass package (not available in memory venv). Input: JSON on stdin with operation and parameters Output: JSON on stdout with result @@ -18,20 +20,153 @@ Output: JSON on stdout with result import sys import json +from pathlib import Path +from datetime import datetime -from aipass.memory.apps.handlers.storage.chroma import store_vectors, list_all_collections, search_vectors +# ============================================================================= +# CHROMADB OPERATIONS (inline — no aipass imports) +# ============================================================================= + +# Default global chroma path: memory/.chroma +_MEMORY_ROOT = Path(__file__).resolve().parents[3] +_DEFAULT_DB_PATH = _MEMORY_ROOT / ".chroma" + +# Singleton clients per path +_clients = {} + + +def _get_client(db_path=None): + """Get or create ChromaDB PersistentClient.""" + import chromadb + + if db_path is None: + db_path = _DEFAULT_DB_PATH + + db_path = Path(db_path) + path_str = str(db_path) + + if path_str not in _clients: + db_path.mkdir(parents=True, exist_ok=True) + _clients[path_str] = chromadb.PersistentClient(path=path_str) + + return _clients[path_str] + + +def _store_vectors(branch, memory_type, embeddings, documents, metadatas, db_path=None): + """Store vectors in branch-specific collection.""" + client = _get_client(db_path) + + collection_name = f"{branch.lower()}_{memory_type.lower()}" + collection = client.get_or_create_collection( + name=collection_name, + metadata={"hnsw:space": "cosine", "branch": branch, "type": memory_type}, + embedding_function=None + ) + + existing_count = collection.count() + timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") + ids = [ + f"{branch}_{memory_type}_{existing_count + i}_{timestamp}" + for i in range(len(embeddings)) + ] + + # Chroma expects lists, not numpy arrays + embeddings_list = [ + emb.tolist() if hasattr(emb, 'tolist') else emb + for emb in embeddings + ] + + collection.add( + embeddings=embeddings_list, + documents=documents, + metadatas=metadatas, + ids=ids + ) + + new_count = collection.count() + + return { + 'success': True, + 'collection': collection_name, + 'count': len(embeddings), + 'total_vectors': new_count, + 'ids': ids + } + + +def _list_collections(db_path=None): + """List all collections.""" + client = _get_client(db_path) + collections = client.list_collections() + names = [col.name for col in collections] + return { + 'success': True, + 'collections': names, + 'count': len(names) + } + + +def _search_vectors(query_embedding, branch=None, memory_type=None, n_results=5, db_path=None): + """Search for similar vectors.""" + client = _get_client(db_path) + + # Determine which collections to search + if branch and memory_type: + collection_names = [f"{branch.lower()}_{memory_type.lower()}"] + else: + all_collections = client.list_collections() + collection_names = [col.name for col in all_collections] + if branch: + collection_names = [c for c in collection_names if c.startswith(branch.lower())] + if memory_type: + collection_names = [c for c in collection_names if c.endswith(memory_type.lower())] + + if not collection_names: + return {'success': True, 'results': [], 'message': 'No matching collections'} + + all_results = [] + for cname in collection_names: + try: + collection = client.get_collection(cname, embedding_function=None) + results = collection.query( + query_embeddings=[query_embedding], + n_results=n_results + ) + if results['documents'] and results['documents'][0]: + for i, doc in enumerate(results['documents'][0]): + all_results.append({ + 'collection': cname, + 'document': doc, + 'metadata': results['metadatas'][0][i] if results['metadatas'] else {}, + 'distance': results['distances'][0][i] if results['distances'] else None, + 'id': results['ids'][0][i] if results['ids'] else None + }) + except Exception: + continue + + all_results.sort(key=lambda x: x['distance'] if x['distance'] is not None else float('inf')) + + return { + 'success': True, + 'results': all_results, + 'collections_searched': len(collection_names), + 'total_results': len(all_results) + } + + +# ============================================================================= +# MAIN +# ============================================================================= def main(): - """Process ChromaDB operation from stdin JSON""" + """Process ChromaDB operation from stdin JSON.""" try: - # Read JSON input from stdin input_data = json.load(sys.stdin) - operation = input_data.get('operation') if operation == 'store_vectors': - result = store_vectors( + result = _store_vectors( branch=input_data.get('branch'), memory_type=input_data.get('memory_type'), embeddings=input_data.get('embeddings'), @@ -40,9 +175,11 @@ def main(): db_path=input_data.get('db_path') ) elif operation == 'list_collections': - result = list_all_collections() + result = _list_collections( + db_path=input_data.get('db_path') + ) elif operation == 'search_vectors': - result = search_vectors( + result = _search_vectors( query_embedding=input_data.get('query_embedding'), branch=input_data.get('branch'), memory_type=input_data.get('memory_type'), @@ -52,11 +189,9 @@ def main(): else: result = {'success': False, 'error': f'Unknown operation: {operation}'} - # Output result as JSON print(json.dumps(result)) except Exception as e: - # Output error as JSON print(json.dumps({'success': False, 'error': str(e)})) sys.exit(1) diff --git a/src/aipass/memory/apps/handlers/vector/embed_subprocess.py b/src/aipass/memory/apps/handlers/vector/embed_subprocess.py new file mode 100644 index 00000000..eb376225 --- /dev/null +++ b/src/aipass/memory/apps/handlers/vector/embed_subprocess.py @@ -0,0 +1,80 @@ +# =================== AIPass ==================== +# Name: embed_subprocess.py +# Description: Embedding Subprocess Handler +# Version: 1.0.0 +# Created: 2026-03-12 +# Modified: 2026-03-12 +# ============================================= + +""" +Embedding Subprocess Handler + +Called via subprocess from rollover orchestrator to ensure sentence-transformers +and torch run in the memory-specific venv (AIPASS_MEMORY_PYTHON). + +Input: JSON on stdin with texts to encode +Output: JSON on stdout with embeddings +""" + +import sys +import json + + +def main(): + """Process embedding request from stdin JSON""" + try: + input_data = json.load(sys.stdin) + texts = input_data.get('texts', []) + + if not texts: + print(json.dumps({'success': True, 'embeddings': [], 'count': 0, 'dimension': 384})) + return + + # Import here — runs in memory venv where these are installed + from sentence_transformers import SentenceTransformer + import torch + + model = SentenceTransformer('all-MiniLM-L6-v2') + + use_gpu = torch.cuda.is_available() + if use_gpu: + model = model.to('cuda') + batch_size = 64 + else: + batch_size = 16 + + # Pre-sort by length (reduces padding waste) + sorted_pairs = sorted(enumerate(texts), key=lambda x: len(x[1])) + sorted_indices, sorted_texts = zip(*sorted_pairs) + + # Encode + embeddings = model.encode( + list(sorted_texts), + batch_size=batch_size, + convert_to_tensor=False, + normalize_embeddings=True, + show_progress_bar=False + ) + + # Restore original order + ordered = [None] * len(texts) + for orig_idx, sorted_idx in enumerate(sorted_indices): + ordered[sorted_idx] = embeddings[orig_idx].tolist() + + if use_gpu: + torch.cuda.empty_cache() + + print(json.dumps({ + 'success': True, + 'embeddings': ordered, + 'count': len(ordered), + 'dimension': 384 + })) + + except Exception as e: + print(json.dumps({'success': False, 'error': str(e)})) + sys.exit(1) + + +if __name__ == '__main__': + main() diff --git a/src/aipass/memory/apps/memory.py b/src/aipass/memory/apps/memory.py index 9a475ad7..685a003a 100755 --- a/src/aipass/memory/apps/memory.py +++ b/src/aipass/memory/apps/memory.py @@ -89,8 +89,7 @@ def print_help(): what_content = ( "Memory is the [bold]central memory archive[/bold] that:\n\n" " [green]>[/green] Provides semantic search across all branch memories\n" - " [green]>[/green] Archives memories when branches hit rollover limits\n" - " [green]>[/green] Extracts symbolic dimensions from conversations" + " [green]>[/green] Archives memories when branches hit rollover limits" ) console.print(Panel( what_content, @@ -108,19 +107,13 @@ def print_help(): table.add_column("Command", style="green") table.add_column("Description", style="dim") - # Core commands + # Core commands (only implemented ones) table.add_row("search ", "Semantic search across all branch memories") - table.add_row("rollover", "Execute memory rollover for files over 600 lines") + table.add_row("rollover", "Execute memory rollover for files exceeding limits") table.add_row("status", "Show rollover statistics for all branches") table.add_row("check", "Check which files need rollover (dry run)") table.add_row("watch", "Start memory watcher (auto-rollover on changes)") table.add_row("sync-lines", "Update line count metadata for all branches") - table.add_row("push-templates", "Push template updates to all branches") - table.add_row("push-templates --dry-run", "Preview template changes without writing") - table.add_row("diff-templates", "Show template differences per branch") - table.add_row("template-status", "Show template version and push status") - table.add_row("symbolic demo", "Run fragmented memory demonstration") - table.add_row("symbolic fragments ", "Search symbolic fragments") console.print(table) @@ -133,7 +126,7 @@ def print_help(): console.print(" [yellow]Via Drone (recommended):[/yellow]") console.print(" [dim]drone @memory search \"error handling\"[/dim]") console.print(" [dim]drone @memory status[/dim]") - console.print(" [dim]drone @memory symbolic demo[/dim]") + console.print(" [dim]drone @memory rollover[/dim]") console.print() console.print(" [yellow]Direct execution:[/yellow]") console.print(" [dim]python3 -m aipass.memory.apps.memory search \"query\"[/dim]") @@ -165,7 +158,7 @@ def print_help(): console.print("-" * 70) console.print() - console.print("Commands: search, rollover, status, check, watch, sync-lines, push-templates, diff-templates, template-status, symbolic") + console.print("Commands: search, rollover, status, check, watch, sync-lines") console.print() @@ -286,7 +279,7 @@ def start_watch() -> None: return console.print(f"[green]>[/green] Watching {result.get('count', 0)} branch directories") - console.print("[dim]Auto-rollover enabled when files exceed 600 lines[/dim]") + console.print("[dim]Auto-rollover enabled when files exceed limits[/dim]") console.print("[dim]Press Ctrl+C to stop[/dim]") console.print() diff --git a/src/aipass/memory/apps/modules/rollover.py b/src/aipass/memory/apps/modules/rollover.py index 967c8846..01ca653c 100755 --- a/src/aipass/memory/apps/modules/rollover.py +++ b/src/aipass/memory/apps/modules/rollover.py @@ -265,7 +265,7 @@ def show_status() -> None: status_marker = "[red]![/red]" if ready else "[green]OK[/green]" - if schema_ver.startswith('2') and v2_reason: + if schema_ver.startswith('2'): status_text = f"READY ({v2_reason})" if ready else "OK (v2)" console.print(f" {status_marker} {memory_type}: {status_text}") else: diff --git a/src/aipass/memory/config/.plans_processed.json b/src/aipass/memory/config/.plans_processed.json new file mode 100644 index 00000000..8e4e889a --- /dev/null +++ b/src/aipass/memory/config/.plans_processed.json @@ -0,0 +1,3 @@ +{ + "FPLAN-0025_build_status_board_per_branch_statusloca_2026-03-10.md": "2026-03-12T16:49:45.098626" +} \ No newline at end of file diff --git a/src/aipass/memory/config/memory_bank.config.json b/src/aipass/memory/config/memory_bank.config.json new file mode 100644 index 00000000..6232916a --- /dev/null +++ b/src/aipass/memory/config/memory_bank.config.json @@ -0,0 +1,15 @@ +{ + "rollover": { + "defaults": { + "max_lines": 600, + "buffer": 100 + }, + "per_branch": {} + }, + "plans": { + "enabled": true, + "path": "src/aipass/flow/processed_plans", + "supported_extensions": [".md"], + "collection_name": "flow_plans" + } +}