feat(system): hook containment + memory bank operational (FPLAN-0026) (#35)

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 <noreply@anthropic.com>
This commit is contained in:
AIPass
2026-03-12 19:29:06 -07:00
committed by GitHub
co-authored by Claude Opus 4.6
parent 1c5e3ee561
commit 717acf078a
16 changed files with 889 additions and 112 deletions
+59 -1
View File
@@ -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
}
]
}
]
}
}
+1
View File
@@ -2,6 +2,7 @@ __pycache__/
*.pyc
*.pyo
.env
.venv/
*.egg-info/
.coverage
htmlcov/
@@ -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
@@ -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
@@ -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
@@ -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')
# =============================================================================
@@ -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,
@@ -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')
@@ -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
@@ -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)
@@ -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()
+6 -13
View File
@@ -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 <query>", "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 <q>", "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()
+1 -1
View File
@@ -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:
@@ -0,0 +1,3 @@
{
"FPLAN-0025_build_status_board_per_branch_statusloca_2026-03-10.md": "2026-03-12T16:49:45.098626"
}
@@ -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"
}
}