feat(system): seedgo v2 operational, full system audit, 14/15 branches at 99%

Three days of intensive work bringing seedgo to full operational status
and driving all branches through comprehensive standards compliance.

Seedgo v2.0.0:
- 22 checkers active (up from 20), standards pack fully operational
- New introspection standard researched from Dev-Pass, FPLAN-0017 open
- Bypass system for false positives (.seedgo config)
- Standards query and audit commands fully functional

System-wide audit (FPLAN-0016):
- All 14 auditable branches at 99%+ compliance
- CLI imports standardized across all branches (console from cli.apps.modules)
- handle_command(command, args) → bool contract added to all modules
- print_help() function naming fixed for checker pattern matching
- Handler extraction: large modules split, file I/O moved to handler layer
- New handlers created across ai_mail, backup, daemon, flow, skills, spawn, seedgo

Branch-specific highlights:
- ai_mail: email.py split 840→420 lines, 4 new handlers
- flow: dplan_flow.py 688→591 lines, 4 new handlers
- seedgo: massive restructure — standards moved to handlers/aipass_standards/,
  old standards/ tree removed, bypass system added, diagnostics module
- commons: database module added, CLI imports fixed
- skills: 5 handle_commands added, help function renamed
- trigger: error reporter handler, handle_command routing
- All branches: consistent architecture, clean drone routing

Culture doc (CLAUDE.md) added — documents AIPass philosophy, identity,
memory system, and collaboration principles.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
AIOSAI
2026-03-10 01:26:42 -07:00
co-authored by Claude Opus 4.6
parent 09e759a8a4
commit babedd9c64
589 changed files with 17140 additions and 24881 deletions
+102 -3
View File
@@ -1,15 +1,110 @@
# MEMORY
**Purpose:** Vector memory archive with ChromaDB semantic search
**Purpose:** Central memory archive with semantic search, rollover, and archival across all AIPass branches.
**Module:** `aipass.memory`
**Created:** 2026-03-07
**Citizen Class:** birthright
**Last Updated:** 2026-03-08
**Citizen Class:** builder
---
## Overview
Birthright citizen — minimal presence with identity and memory.
Memory is the central memory archive system that:
- Provides semantic search across all branch memories
- Archives memories when branches hit rollover limits (600 lines)
- Extracts symbolic dimensions from conversations
- Manages template distribution and line-count tracking across branches
---
## Commands / Usage
**Via Drone (recommended):**
```bash
drone @memory search "error handling" # Semantic search across all branch memories
drone @memory search "query" --branch SEED # Filter search by branch
drone @memory search "query" --n 10 # Limit number of results
drone @memory rollover # Execute memory rollover for files over 600 lines
drone @memory status # Show rollover statistics for all branches
drone @memory check # Dry run — check which files need rollover
drone @memory watch # Start auto-rollover watcher (Ctrl+C to stop)
drone @memory sync-lines # Update line count metadata for all branches
drone @memory push-templates # Push template updates to all branches
drone @memory push-templates --dry-run # Preview template changes without writing
drone @memory diff-templates # Show template differences per branch
drone @memory template-status # Show template version and push status
drone @memory symbolic demo # Run fragmented memory demonstration
drone @memory symbolic fragments "query" # Search symbolic fragments
```
**Direct execution:**
```bash
python3 -m aipass.memory.apps.memory search "query"
python3 -m aipass.memory.apps.memory rollover
```
---
## Architecture
```
memory/
├── __init__.py # Package init
├── README.md # This file
├── DASHBOARD.local.json # System status dashboard
├── pytest.ini # Test configuration
├── apps/
│ ├── __init__.py
│ ├── memory.py # Entry point (CLI) — auto-discovers modules
│ ├── modules/
│ │ ├── __init__.py
│ │ ├── rollover.py # Rollover orchestrator — line checks, archival triggers
│ │ └── search.py # Search orchestrator — semantic query routing
│ ├── handlers/
│ │ ├── __init__.py
│ │ ├── central_writer.py # Central memory write operations
│ │ ├── dashboard_push.py # Dashboard status push
│ │ ├── archive/ # Memory archival indexing
│ │ ├── json/ # JSON handler operations
│ │ ├── learnings/ # Learning extraction
│ │ ├── monitor/ # File watcher for auto-rollover
│ │ ├── rollover/ # Rollover implementation logic
│ │ ├── schema/ # Memory schema definitions
│ │ ├── search/ # Search implementation (vector/semantic)
│ │ ├── storage/ # Storage backend operations
│ │ ├── tracking/ # Line count and metadata tracking
│ │ └── vector/ # Vector DB (ChromaDB) operations
│ ├── extensions/ # Extension plugins
│ ├── json_templates/ # JSON template files
│ └── plugins/ # Plugin system
├── artifacts/ # Build/output artifacts
├── docs/ # Documentation
├── memory_json/ # Memory JSON data store
├── tests/ # Test suite
└── tools/ # Utility scripts
```
---
## Integration Points
### Depends On
- `rich` — Console output, panels, and tables
- Python stdlib (`sys`, `time`, `signal`, `logging`, `pathlib`, `importlib`)
### Provides To
- All branches — memory rollover, archival, and retrieval services
- All branches — semantic search across branch memories
- All branches — template distribution via `push-templates`
- All branches — line count metadata via `sync-lines`
---
## Key Modules
### rollover
The rollover module monitors memory files across all branches registered in `AIPASS_REGISTRY.json`. When files exceed 600 lines, it triggers archival — splitting the file, preserving recent context, and indexing the archived portion. The `watch` command runs a persistent file watcher that auto-triggers rollover on changes.
---
@@ -19,3 +114,7 @@ Birthright citizen — minimal presence with identity and memory.
- **Session History:** `.trinity/local.json`
- **Observations:** `.trinity/observations.json`
- **Branch Prompt:** `.aipass/branch_system_prompt.md`
---
*Last Updated: 2026-03-08*
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: indexer.py - Code Archive Indexer
# Date: 2025-11-27
# =================== AIPass ====================
# Name: indexer.py
# Description: Code Archive Indexer
# Version: 0.2.0
# Category: memory/handlers/archive
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2025-11-27): Initial version - auto-index code archive files
#
# CODE STANDARDS:
# - Pure handler: No orchestration, just indexing logic
# - Stateless functions
# - Returns dict with success/error
# Created: 2025-11-27
# Modified: 2026-03-06
# =============================================
"""
@@ -33,7 +24,9 @@ from pathlib import Path
from datetime import datetime
from typing import Dict, Any, List
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Paths resolved relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: central_writer.py - Central File Writer Handler
# Date: 2025-11-27
# =================== AIPass ====================
# Name: central_writer.py
# Description: Central File Writer Handler
# Version: 0.2.0
# Category: memory/handlers
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2025-11-27): Initial implementation - central file writer
#
# CODE STANDARDS:
# - Handler tier 3: Pure functions, raises exceptions
# - No CLI imports, no Prax imports
# - Returns success/error dicts
# Created: 2025-11-27
# Modified: 2026-03-06
# =============================================
"""
@@ -32,7 +23,9 @@ from pathlib import Path
from datetime import datetime
from typing import Dict, Any
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# No service imports - handlers are pure workers (3-tier architecture)
@@ -1,19 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: dashboard_push.py - Memory Bank Dashboard Write-Through
# Date: 2026-02-25
# =================== AIPass ====================
# Name: dashboard_push.py
# Description: Memory Bank Dashboard Write-Through
# Version: 0.2.0
# Category: memory/handlers
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2026-02-25): Dashboard write-through for Memory Bank
#
# CODE STANDARDS:
# - Handler tier 3: Pure functions, no CLI/Prax imports
# - Dashboard write failures are silent (return False, never raise)
# - Uses subprocess for cross-branch dashboard writes (no cross-package imports)
# - BYPASS: Direct json reads required for central.json, registry, config files
# Created: 2026-02-25
# Modified: 2026-03-06
# =============================================
"""
@@ -35,7 +25,9 @@ from pathlib import Path
from datetime import datetime
from typing import Dict, Any, List
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: json_handler.py - Memory File Safe Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: json_handler.py
# Description: Memory File Safe Handler
# Version: 0.2.0
# Category: memory/handlers/json
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2025-11-16): Initial version - safe memory file operations
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <300 lines target
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -42,7 +33,9 @@ from pathlib import Path
from typing import Dict, Any, Optional
from datetime import datetime
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -1,19 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: manager.py - Memory Sections Management Handler
# Date: 2026-02-04
# =================== AIPass ====================
# Name: manager.py
# Description: Memory Sections Management Handler
# Version: 1.2.0
# Category: memory/handlers/learnings
#
# CHANGELOG (Max 5 entries):
# - v1.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v1.1.0 (2026-02-04): Added recently_completed management + status count updates
# - v1.0.0 (2026-02-04): Initial version - timestamp tracking, max_entries, vectorization
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <600 lines target
# Created: 2026-02-04
# Modified: 2026-03-06
# =============================================
"""
@@ -43,7 +33,9 @@ from pathlib import Path
from typing import Dict, Any, List, Tuple
from datetime import datetime
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Handler imports (relative within package)
from aipass.memory.apps.handlers.json.json_handler import (
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: detector.py - Rollover Trigger Detection Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: detector.py
# Description: Rollover Trigger Detection Handler
# Version: 0.2.0
# Category: memory/handlers/monitor
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2025-11-16): Initial version - detect 600-line rollover triggers
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Returns status dicts (3-tier architecture)
# - File size: <200 lines target
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -35,7 +26,9 @@ from pathlib import Path
from typing import List, Dict, Any
from dataclasses import dataclass
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# No service imports - handlers are pure workers (3-tier architecture)
# No module imports (handler independence)
@@ -1,21 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: memory_watcher.py - Memory File System Watcher
# Date: 2025-11-26
# =================== AIPass ====================
# Name: memory_watcher.py
# Description: Memory File System Watcher
# Version: 1.1.0
# Category: memory/handlers/monitor
#
# CHANGELOG (Max 5 entries):
# - v1.1.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v1.0.0 (2025-11-26): Initial version
# * Watches memory files for modifications
# * Auto-updates line counts on file changes
# * Triggers rollover when files exceed 600 lines
#
# CODE STANDARDS:
# - Handler independence: Uses watchdog library
# - Error handling: Returns status dicts (3-tier architecture)
# - Pure handler - no CLI display imports
# Created: 2025-11-26
# Modified: 2026-03-06
# =============================================
"""
@@ -50,8 +38,9 @@ except ImportError:
# Handler imports (relative within package)
from aipass.memory.apps.handlers.tracking.line_counter import update_line_count
from aipass.memory.apps.handlers.monitor.detector import check_single_file
from aipass.prax.apps.modules.logger import get_system_logger
logger = logging.getLogger(__name__)
logger = get_system_logger()
# Memory root resolved relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -204,10 +193,10 @@ def check_and_rollover() -> Dict[str, Any]:
results['rollover_triggered'] = True
try:
from aipass.memory.apps.modules.rollover import execute_rollover
from aipass.memory.apps.handlers.rollover.orchestrator import execute_rollover
execute_rollover()
except ImportError:
logger.warning("Rollover module not available")
logger.warning("Rollover handler not available")
except Exception as e:
results['rollover_error'] = str(e)
results['success'] = False
@@ -457,8 +446,8 @@ class MemoryFileWatcher(FileSystemEventHandler):
trigger = check_result.get('trigger')
logger.warning(f"[memory_watcher] ROLLOVER TRIGGERED: {trigger}")
# Import rollover module here to avoid circular imports
from aipass.memory.apps.modules.rollover import execute_rollover
# Import rollover handler here to avoid circular imports
from aipass.memory.apps.handlers.rollover.orchestrator import execute_rollover
# Trigger rollover
logger.info(f"[memory_watcher] Triggering rollover for {file_path.name}")
@@ -1,20 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: extractor.py - Memory Extraction Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: extractor.py
# Description: Memory Extraction Handler
# Version: 0.4.0
# Category: memory/handlers/rollover
#
# CHANGELOG (Max 5 entries):
# - v0.4.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.3.0 (2025-11-16): Refactored to use json_handler instead of direct json ops
# - v0.2.0 (2025-11-16): Complete redesign - work with real memory file structure
# - v0.1.0 (2025-11-16): Initial version (broken - assumed 'entries' array)
#
# CODE STANDARDS:
# - Handler independence: Uses json_handler only
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <400 lines target
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -42,8 +31,9 @@ from datetime import datetime
# Handler imports (relative within package)
from aipass.memory.apps.handlers.json.json_handler import read_memory_file_data, write_memory_file_simple
from aipass.prax.apps.modules.logger import get_system_logger
logger = logging.getLogger(__name__)
logger = get_system_logger()
# No module imports (handler independence)
@@ -0,0 +1,433 @@
# =================== AIPass ====================
# Name: orchestrator.py
# Description: Rollover Orchestration Handler
# Version: 1.0.0
# Created: 2026-03-08
# Modified: 2026-03-08
# =============================================
"""
Rollover Orchestration Handler
Contains the core rollover execution logic: trigger detection, extraction,
embedding, vector storage, and line count sync. Called by the rollover
module which handles display/CLI concerns.
Purpose:
Implementation logic for rollover workflow, separated from CLI/display
layer to satisfy thin-module standard.
"""
import subprocess
import json
import os
import sys
from pathlib import Path
from typing import List, Dict, Any
from aipass.prax import logger
# logger imported from aipass.prax
# 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
_HANDLERS_DIR = Path(__file__).resolve().parent.parent
CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py"
# Use system python by default; can be overridden via environment variable
MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable)
# =============================================================================
# REPO ROOT DISCOVERY
# =============================================================================
def _find_repo_root() -> Path:
"""Walk up from this file to find the 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()
# =============================================================================
# VECTOR STORAGE (SUBPROCESS)
# =============================================================================
def store_vectors_subprocess(branch: str, memory_type: str, embeddings: list,
documents: list, metadatas: list, db_path: str | Path | None = None) -> dict:
"""
Store vectors via subprocess.
This ensures ChromaDB compatibility regardless of calling Python version.
Args:
branch: Branch name
memory_type: Type of memory (e.g., 'sessions', 'observations')
embeddings: List of embedding vectors
documents: List of text documents
metadatas: List of metadata dicts
db_path: Path to Chroma database (None for global)
Returns:
Dict with success status and storage details
"""
# Convert numpy arrays to lists for JSON serialization
embeddings_serializable = [
emb.tolist() if hasattr(emb, 'tolist') else emb
for emb in embeddings
]
input_data = {
'operation': 'store_vectors',
'branch': branch,
'memory_type': memory_type,
'embeddings': embeddings_serializable,
'documents': documents,
'metadatas': metadatas,
'db_path': str(db_path) if db_path else None
}
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 'Subprocess failed'}
return json.loads(result.stdout)
except subprocess.TimeoutExpired:
return {'success': False, 'error': 'Storage operation timed out'}
except json.JSONDecodeError as e:
return {'success': False, 'error': f'Invalid JSON response: {e}'}
except Exception as e:
return {'success': False, 'error': str(e)}
# =============================================================================
# PATH HELPERS
# =============================================================================
def get_branch_local_chroma_path(branch_name: str) -> Path | None:
"""
Get local .chroma path for branch
Args:
branch_name: Branch name (e.g., "SEED", "AIPASS")
Returns:
Path to branch's local .chroma directory, or None if branch not found
"""
if not branch_name:
return None
registry = detector._read_registry()
for branch in registry:
if branch.get('name', '').upper() == branch_name.upper():
branch_path = Path(branch.get('path', ''))
if branch_path.exists():
chroma_path = branch_path / '.chroma'
# Auto-create .chroma directory if missing
if not chroma_path.exists():
chroma_path.mkdir(parents=True, exist_ok=True)
logger.info(f"[rollover] Created local .chroma directory for {branch_name}")
return chroma_path
logger.warning(f"[rollover] Branch {branch_name} not found in registry")
return None
# =============================================================================
# TEXT EXTRACTION HELPERS
# =============================================================================
def extract_text_from_memories(memories: List[Dict]) -> List[str]:
"""
Extract text content from memory items for vectorization
Memory items have different structures:
- sessions: 'activities' array (join into text)
- observations: might have 'content' or 'text' field
- generic: convert to JSON string
Args:
memories: List of memory items
Returns:
List of text strings for embedding
"""
texts = []
for memory in memories:
# Try common text fields
if 'activities' in memory and isinstance(memory['activities'], list):
# Sessions type - join activities
text = '\n'.join(str(a) for a in memory['activities'])
elif 'content' in memory:
text = str(memory['content'])
elif 'text' in memory:
text = str(memory['text'])
elif 'message' in memory:
text = str(memory['message'])
else:
# Fallback - convert to string representation
text = str(memory)
texts.append(text)
return texts
# =============================================================================
# ROLLOVER EXECUTION
# =============================================================================
def execute_rollover() -> Dict[str, Any]:
"""
Execute rollover workflow for all triggered branches.
Workflow:
1. Check all branches for triggers
2. For each trigger:
- Create backup
- Extract oldest entries
- Generate embeddings
- Store in local + global Chroma
- Update line counts
3. Return results
Returns:
Dict with success status, counts, and details for each trigger
"""
# Step 1: Detect triggers
triggers_result = detector.check_all_branches()
if not triggers_result['success']:
error = triggers_result.get('error', 'Unknown error')
logger.error(f"[rollover] Failed to check branches: {error}")
return {
'success': False,
'error': f'Failed to check for rollover triggers: {error}',
'triggers_count': 0,
'success_count': 0,
'failed': [],
}
triggers = triggers_result.get('triggers', [])
if not triggers:
logger.info("[rollover] No rollover triggers detected")
return {
'success': True,
'triggers_count': 0,
'success_count': 0,
'failed': [],
'results': [],
}
logger.info(f"[rollover] Found {len(triggers)} files ready for rollover")
# Process each trigger
success_count = 0
failed = []
results = []
for trigger in triggers:
# Step 1: CREATE BACKUP (safety net)
backup_result = extractor.create_rollover_backup(trigger.file_path)
if not backup_result['success']:
error_msg = backup_result.get('error', 'Backup failed')
logger.error(f"[rollover] Backup failed for {trigger}: {error_msg}")
failed.append({'trigger': str(trigger), 'stage': 'backup', 'error': error_msg})
continue # Don't proceed without backup
logger.info(f"[rollover] {backup_result.get('message')}")
# Step 2: Extract memories (auto-calculates percentage)
extract_result = extractor.extract_with_metadata(trigger.file_path)
if not extract_result['success']:
error_msg = extract_result.get('error', 'Unknown error')
logger.error(f"[rollover] Extraction failed for {trigger}: {error_msg}")
# RESTORE from backup
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after extraction failure")
failed.append({'trigger': str(trigger), 'stage': 'extraction', 'error': error_msg})
continue
memories = extract_result.get('entries', [])
branch = extract_result.get('branch', '')
memory_type = extract_result.get('type', 'unknown')
old_lines = extract_result.get('old_lines', 0)
new_lines = extract_result.get('new_lines', 0)
if not branch:
logger.error(f"[rollover] No branch found in extraction result for {trigger}")
failed.append({'trigger': str(trigger), 'stage': 'extraction', 'error': 'No branch in result'})
continue
logger.info(f"[rollover] Extracted {len(memories)} items from {trigger} ({old_lines} -> {new_lines} lines)")
# Convert memory items to text for vectorization
texts = extract_text_from_memories(memories)
# Step 3: Generate embeddings
embed_result = embedder.encode_batch(texts)
if not embed_result['success']:
error_msg = embed_result.get('error', 'Unknown error')
logger.error(f"[rollover] Embedding failed for {trigger}: {error_msg}")
# RESTORE from backup
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after embedding failure")
failed.append({'trigger': str(trigger), 'stage': 'embedding', 'error': error_msg})
continue
embeddings = embed_result.get('embeddings', [])
if not embeddings:
logger.error(f"[rollover] No embeddings generated for {trigger}")
failed.append({'trigger': str(trigger), 'stage': 'embedding', 'error': 'No embeddings in result'})
continue
logger.info(f"[rollover] Generated {len(embeddings)} embeddings for {trigger}")
# Step 4: Prepare metadata for vectorization
metadatas = []
for memory in memories:
metadata = memory.get('_metadata', {})
metadata['timestamp'] = memory.get('timestamp', '')
metadatas.append(metadata)
# Step 5: Store in LOCAL branch Chroma (via subprocess)
branch_str: str = branch
memory_type_str: str = memory_type
embeddings_list: list = embeddings
local_chroma_path = get_branch_local_chroma_path(branch_str)
local_store_result = None
if local_chroma_path:
local_store_result = store_vectors_subprocess(
branch=branch_str,
memory_type=memory_type_str,
embeddings=embeddings_list,
documents=texts,
metadatas=metadatas,
db_path=str(local_chroma_path)
)
if not local_store_result['success']:
logger.warning(f"[rollover] Local storage failed for {branch}: {local_store_result.get('error')}")
# Continue anyway - global storage is primary
else:
logger.info(f"[rollover] Stored {len(embeddings)} vectors in local Chroma for {branch}")
# Step 6: Store in GLOBAL Memory Chroma (via subprocess)
global_store_result = store_vectors_subprocess(
branch=branch_str,
memory_type=memory_type_str,
embeddings=embeddings_list,
documents=texts,
metadatas=metadatas
# db_path=None means global
)
if not global_store_result['success']:
error_msg = global_store_result.get('error', 'Unknown error')
logger.error(f"[rollover] Global storage failed for {trigger}: {error_msg}")
# RESTORE from backup (CRITICAL - file was modified but storage failed)
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after storage failure")
else:
logger.error(f"[rollover] CRITICAL: Failed to restore from backup: {restore_result.get('error')}")
failed.append({'trigger': str(trigger), 'stage': 'global_storage', 'error': error_msg})
continue
logger.info(f"[rollover] Stored {len(embeddings)} vectors in global Chroma for {branch}")
# Step 7: Update line count metadata
update_result = line_counter.update_line_count(trigger.file_path)
if update_result['success']:
logger.info(f"[rollover] Updated line count metadata for {trigger.file_path.name}")
else:
logger.warning(f"[rollover] Failed to update line count for {trigger.file_path.name}: {update_result.get('error')}")
# Success!
success_count += 1
global_collection = global_store_result.get('collection')
global_total = global_store_result.get('total_vectors')
local_ok = local_store_result and local_store_result['success']
results.append({
'trigger': str(trigger),
'memories_count': len(memories),
'old_lines': old_lines,
'new_lines': new_lines,
'global_collection': global_collection,
'global_total': global_total,
'local_stored': local_ok,
})
logger.info(f"[rollover] Successfully rolled over {trigger}: {len(memories)} items, {old_lines} -> {new_lines} lines")
# Summary logging
if success_count > 0:
logger.info(f"[rollover] Rollover complete: {success_count}/{len(triggers)} successful")
if failed:
logger.error(f"[rollover] {len(failed)} operations failed")
return {
'success': success_count > 0 or len(triggers) == 0,
'triggers_count': len(triggers),
'success_count': success_count,
'failed': failed,
'results': results,
}
# =============================================================================
# LINE COUNT SYNC
# =============================================================================
def sync_line_counts() -> Dict[str, Any]:
"""
Update line count metadata for all branch memory files.
Reads actual line counts and updates document_metadata.status.current_lines
for all *.local.json and *.observations.json files in AIPASS_REGISTRY.
Returns:
Dict with success status, updated count, and failures
"""
result = line_counter.update_all_memory_files()
if result['success']:
logger.info(f"[rollover] Synced line counts: {result['updated']} updated, {result['failed']} failed")
else:
logger.error("[rollover] Failed to sync line counts")
return result
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: normalize.py - Memory File Schema Normalizer
# Date: 2026-01-22
# =================== AIPass ====================
# Name: normalize.py
# Description: Memory File Schema Normalizer
# Version: 0.2.0
# Category: memory/handlers/schema
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2026-01-22): Initial version - normalize metadata schema
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Returns status dicts (3-tier architecture)
# - File size: <200 lines target
# Created: 2026-01-22
# Modified: 2026-03-06
# =============================================
"""
@@ -40,7 +31,9 @@ from pathlib import Path
from typing import Dict, Any
from datetime import datetime
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
def normalize_memory_file(file_path: Path, dry_run: bool = False) -> Dict[str, Any]:
@@ -0,0 +1,250 @@
# =================== AIPass ====================
# Name: query_executor.py
# Description: Search Query Execution Handler
# Version: 1.0.0
# Created: 2026-03-08
# Modified: 2026-03-08
# =============================================
"""
Search Query Execution Handler
Contains the core search execution logic: subprocess-based vector search,
query encoding via embedder, similarity calculation, and result filtering.
Called by the search module which handles display/CLI concerns.
Purpose:
Implementation logic for semantic search, separated from CLI/display
layer to satisfy thin-module standard.
"""
import subprocess
import json
import os
import sys
from pathlib import Path
from typing import Dict, Any
from aipass.prax import logger
# Handler imports
from aipass.memory.apps.handlers.vector import embedder
# ChromaDB search via subprocess
_HANDLERS_DIR = Path(__file__).resolve().parent.parent
CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py"
# Use system python by default; can be overridden via environment variable
MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable)
# Minimum similarity threshold - filter out irrelevant results
MIN_SIMILARITY_THRESHOLD = 0.40 # 40% minimum relevance
# =============================================================================
# SUBPROCESS VECTOR SEARCH
# =============================================================================
def search_vectors_subprocess(
query_embedding: list,
branch: str | None = None,
memory_type: str | None = None,
n_results: int = 5,
db_path: str | Path | None = None
) -> dict:
"""
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
memory_type: Optional memory type filter
n_results: Number of results to return
db_path: Path to Chroma database (None for global)
Returns:
Dict with success status and search results
"""
input_data = {
'operation': 'search_vectors',
'query_embedding': query_embedding,
'branch': branch,
'memory_type': memory_type,
'n_results': n_results,
'db_path': str(db_path) if db_path else None
}
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 'Subprocess failed'}
return json.loads(result.stdout)
except subprocess.TimeoutExpired:
return {'success': False, 'error': 'Search operation timed out'}
except json.JSONDecodeError as e:
return {'success': False, 'error': f'Invalid JSON response: {e}'}
except Exception as e:
return {'success': False, 'error': str(e)}
# =============================================================================
# RESULT PROCESSING
# =============================================================================
def _calculate_similarity(distance: float) -> float:
"""
Calculate similarity from ChromaDB L2 distance.
ChromaDB L2 distance: 0=identical, ~2=very different.
Args:
distance: L2 distance from ChromaDB
Returns:
Similarity score between 0 and 1
"""
return max(0, 1 - (distance / 2))
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
Returns:
List of filtered result dicts with similarity scores added
"""
filtered = []
for result in results[:n_results]:
document = result.get('document', '')
distance = result.get('distance', 0)
similarity = _calculate_similarity(distance)
# Skip empty documents and low-relevance results
if not document or not document.strip():
continue
if similarity < MIN_SIMILARITY_THRESHOLD:
continue
result['similarity'] = similarity
filtered.append(result)
return filtered
# =============================================================================
# PUBLIC API
# =============================================================================
def execute_search(
query: str,
branch: str | None = None,
memory_type: str | None = None,
n_results: int = 5
) -> Dict[str, Any]:
"""
Execute semantic search: encode query, search vectors, filter results.
Workflow:
1. Encode query to embedding vector via embedder handler
2. Search ChromaDB via subprocess
3. Filter and score results by similarity
Args:
query: Search query text
branch: Optional branch filter
memory_type: Optional memory type filter
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)
"""
# Step 1: Encode query
embed_result = embedder.encode_batch([query])
if not embed_result['success']:
error_msg = embed_result.get('error', 'Unknown error')
logger.error(f"[search] Failed to encode query: {error_msg}")
return {
'success': False,
'error': f'Failed to encode query: {error_msg}',
'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()
logger.info(f"[search] Encoded query to {len(query_embedding)}-dim vector")
# Step 2: Search via subprocess
search_result = search_vectors_subprocess(
query_embedding=query_embedding,
branch=branch,
memory_type=memory_type,
n_results=n_results
)
if not search_result['success']:
error_msg = search_result.get('error', 'Unknown error')
logger.error(f"[search] Search failed: {error_msg}")
return {
'success': False,
'error': f'Search failed: {error_msg}',
'query': query,
}
raw_results = search_result.get('results', [])
collections_searched = search_result.get('collections_searched', 0)
total_results = search_result.get('total_results', 0)
logger.info(f"[search] Found {total_results} results across {collections_searched} collections")
# Step 3: Filter and score results
filtered_results = _filter_results(raw_results, n_results)
logger.info(f"[search] Filtered to {len(filtered_results)} relevant results")
return {
'success': True,
'query': query,
'branch': branch,
'memory_type': memory_type,
'results': filtered_results,
'collections_searched': collections_searched,
'total_results': total_results,
'filtered_count': len(filtered_results),
}
@@ -1,21 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: vector_search.py - Vector Search Handler
# Date: 2025-11-27
# =================== AIPass ====================
# Name: vector_search.py
# Description: Vector Search Handler
# Version: 0.3.0
# Category: memory/handlers/search
#
# CHANGELOG (Max 5 entries):
# - v0.3.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.2.0 (2026-02-15): Use shared singleton client, embedding_function=None
# on all collection access
# - v0.1.0 (2025-11-27): Initial version - ChromaDB semantic search
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <300 lines target
# - Best practices: Singleton pattern, same model as embedder
# Created: 2025-11-27
# Modified: 2026-03-06
# =============================================
"""
@@ -44,7 +32,9 @@ import logging
from typing import List, Dict, Any
from pathlib import Path
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -1,21 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: chroma.py - Chroma Vector Storage Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: chroma.py
# Description: Chroma Vector Storage Handler
# Version: 0.3.0
# Category: memory/handlers/storage
#
# CHANGELOG (Max 5 entries):
# - v0.3.0 (2026-03-06): Adapted for AIPass public repo - optional chromadb
# - v0.2.0 (2026-02-15): Use shared singleton client, cosine distance,
# embedding_function=None on all collection access
# - v0.1.0 (2025-11-16): Initial version - Chroma collection management
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <300 lines target
# - Best practices: Collection-per-branch, batch inserts
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -44,7 +32,9 @@ from typing import List, Dict, Any
from pathlib import Path
from datetime import datetime
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
@@ -1,17 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: chroma_subprocess.py - ChromaDB Subprocess Handler
# Date: 2025-11-27
# =================== AIPass ====================
# Name: chroma_subprocess.py
# Description: ChromaDB Subprocess Handler
# Version: 1.1.0
# Category: memory/handlers/storage
#
# CHANGELOG (Max 5 entries):
# - v1.1.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v1.0.0 (2025-11-27): Initial version - subprocess wrapper for ChromaDB
#
# CODE STANDARDS:
# - Accepts JSON input via stdin, outputs JSON to stdout
# - Handler pattern - pure worker, no logging
# Created: 2025-11-27
# Modified: 2026-03-06
# =============================================
"""
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: line_counter.py - Memory File Line Counter Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: line_counter.py
# Description: Memory File Line Counter Handler
# Version: 0.2.0
# Category: memory/handlers/tracking
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed hardcoded paths
# - v0.1.0 (2025-11-16): Initial version - count lines and update metadata
#
# CODE STANDARDS:
# - Handler independence: Uses json_handler for safe operations
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <200 lines target
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -36,8 +27,9 @@ from datetime import datetime
# Handler imports (relative within package)
from aipass.memory.apps.handlers.json.json_handler import update_metadata
from aipass.prax.apps.modules.logger import get_system_logger
logger = logging.getLogger(__name__)
logger = get_system_logger()
# =============================================================================
@@ -1,19 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: embedder.py - Vector Embedding Handler
# Date: 2025-11-16
# =================== AIPass ====================
# Name: embedder.py
# Description: Vector Embedding Handler
# Version: 0.2.0
# Category: memory/handlers/vector
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - optional deps, no hardcoded paths
# - v0.1.0 (2025-11-16): Initial version - sentence-transformers integration
#
# CODE STANDARDS:
# - Handler independence: No module imports
# - Error handling: Return status dicts (3-tier architecture)
# - File size: <300 lines target
# - Best practices: Sorting, normalization, GPU cleanup
# Created: 2025-11-16
# Modified: 2026-03-06
# =============================================
"""
@@ -42,7 +32,9 @@ import logging
from typing import List, Dict, Any
from pathlib import Path
logger = logging.getLogger(__name__)
from aipass.prax.apps.modules.logger import get_system_logger
logger = get_system_logger()
# No service imports - handlers are pure workers (3-tier architecture)
# No module imports (handler independence)
+11 -21
View File
@@ -1,18 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: memory.py - Memory System Entry Point
# Date: 2025-11-15
# Version: 0.2.0
# Category: memory
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed internal deps
# - v0.1.0 (2025-11-15): Initial version - modular architecture, module discovery
#
# CODE STANDARDS:
# - Follows AIPass architecture patterns (learned from Seed)
# - Module auto-discovery with handle_command() interface
# =================== AIPass ====================
# Name: memory.py
# Description: Entry point CLI for drone @memory
# Version: 1.0.0
# Created: 2026-03-08
# Modified: 2026-03-08
# =============================================
"""
@@ -29,23 +20,22 @@ ARCHITECTURE:
import sys
import time
import signal
import logging
from pathlib import Path
from typing import List, Any
import importlib
from rich.console import Console
from rich.panel import Panel
from rich import box
from rich.table import Table
from aipass.prax import logger
from aipass.cli.apps.modules import console
# =============================================================================
# INFRASTRUCTURE SETUP
# =============================================================================
logger = logging.getLogger(__name__)
console = Console()
# Package name for module discovery (relative imports within this package)
_PACKAGE_BASE = "aipass.memory.apps.modules"
@@ -267,7 +257,7 @@ def start_watch() -> None:
is_memory_watcher_active,
get_watcher_status
)
from .modules.rollover import get_rollover_stats
from ..handlers.monitor.detector import get_rollover_stats
# Signal handler for graceful shutdown
def signal_handler(sig, frame):
+59 -359
View File
@@ -1,19 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: rollover.py - Rollover Orchestration Module
# Date: 2025-11-16
# Version: 0.2.0
# Category: memory/modules
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed internal deps
# - v0.1.0 (2025-11-16): Initial version - orchestrate rollover workflow
#
# CODE STANDARDS:
# - Thin orchestration: Delegate all logic to handlers
# - No business logic: Only coordinate workflow
# - handle_command() pattern
# =================== AIPass ====================
# Name: rollover.py
# Description: Rollover Orchestration Module
# Version: 0.5.0
# Created: 2025-11-16
# Modified: 2026-03-08
# =============================================
"""
@@ -31,100 +21,24 @@ Purpose:
"""
import sys
import logging
import subprocess
import json
from pathlib import Path
from typing import List, Dict
from typing import List
from rich.console import Console
from rich.panel import Panel
from rich import box
from aipass.prax import logger
from aipass.cli.apps.modules import console
# =============================================================================
# INFRASTRUCTURE SETUP
# =============================================================================
logger = logging.getLogger(__name__)
console = Console()
# Handler imports (relative within the memory package)
# Handler imports
from ..handlers.monitor import detector
from ..handlers.rollover import extractor
from ..handlers.vector import embedder
from ..handlers.tracking import line_counter
# ChromaDB storage via subprocess
# Resolve paths relative to the handlers directory
_HANDLERS_DIR = Path(__file__).resolve().parent.parent / "handlers"
CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py"
# Use system python by default; can be overridden via environment variable
import os
MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable)
# =============================================================================
# REPO ROOT DISCOVERY
# =============================================================================
def _find_repo_root() -> Path:
"""Walk up from this file to find the 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()
# No other module imports (modules don't import modules)
def _store_vectors_subprocess(branch: str, memory_type: str, embeddings: list,
documents: list, metadatas: list, db_path: str | Path | None = None) -> dict:
"""
Store vectors via subprocess.
This ensures ChromaDB compatibility regardless of calling Python version.
"""
# Convert numpy arrays to lists for JSON serialization
embeddings_serializable = [
emb.tolist() if hasattr(emb, 'tolist') else emb
for emb in embeddings
]
input_data = {
'operation': 'store_vectors',
'branch': branch,
'memory_type': memory_type,
'embeddings': embeddings_serializable,
'documents': documents,
'metadatas': metadatas,
'db_path': str(db_path) if db_path else None
}
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 'Subprocess failed'}
return json.loads(result.stdout)
except subprocess.TimeoutExpired:
return {'success': False, 'error': 'Storage operation timed out'}
except json.JSONDecodeError as e:
return {'success': False, 'error': f'Invalid JSON response: {e}'}
except Exception as e:
return {'success': False, 'error': str(e)}
from ..handlers.rollover.orchestrator import (
execute_rollover as _handler_execute_rollover,
sync_line_counts as _handler_sync_line_counts,
)
# =============================================================================
@@ -153,7 +67,7 @@ def handle_command(command: str, args: List[str]) -> bool: # noqa: ARG001
return True
if command == 'rollover':
execute_rollover()
run_rollover()
return True
elif command == 'status':
@@ -202,17 +116,11 @@ def print_help() -> None:
# ROLLOVER ORCHESTRATION
# =============================================================================
def execute_rollover() -> bool:
def run_rollover() -> bool:
"""
Execute rollover workflow for all triggered branches
Execute rollover workflow for all triggered branches.
Workflow:
1. Check all branches for triggers
2. For each trigger:
- Extract oldest 100 entries
- Generate embeddings
- Store in Chroma
3. Report results
Delegates to handler and renders results with Rich.
"""
console.print()
console.print(Panel.fit(
@@ -222,267 +130,47 @@ def execute_rollover() -> bool:
))
console.print()
# Step 1: Detect triggers
console.print("[cyan]Checking for rollover triggers...[/cyan]")
triggers_result = detector.check_all_branches()
if not triggers_result['success']:
logger.error(f"[rollover] Failed to check branches: {triggers_result.get('error', 'Unknown error')}")
console.print("[red]x[/red] Failed to check for rollover triggers")
result = _handler_execute_rollover()
if not result.get('success') and result.get('error'):
console.print(f"[red]x[/red] {result['error']}")
return False
triggers = triggers_result.get('triggers', [])
if not triggers:
triggers_count = result.get('triggers_count', 0)
if triggers_count == 0:
console.print("[green]>[/green] No files need rollover")
logger.info("[rollover] No rollover triggers detected")
return True
console.print(f"[green]>[/green] Found {len(triggers)} files ready for rollover")
logger.info(f"[rollover] Found {len(triggers)} files ready for rollover")
console.print(f"[green]>[/green] Found {triggers_count} files ready for rollover")
console.print()
# Process each trigger
success_count = 0
failed = []
for trigger in triggers:
console.print(f"[yellow]Processing:[/yellow] {trigger}")
# Step 1: CREATE BACKUP (safety net)
backup_result = extractor.create_rollover_backup(trigger.file_path)
if not backup_result['success']:
error_msg = backup_result.get('error', 'Backup failed')
logger.error(f"[rollover] Backup failed for {trigger}: {error_msg}")
failed.append((trigger, "backup", error_msg))
continue # Don't proceed without backup
logger.info(f"[rollover] {backup_result.get('message')}")
# Step 2: Extract memories (auto-calculates percentage)
extract_result = extractor.extract_with_metadata(trigger.file_path)
if not extract_result['success']:
error_msg = extract_result.get('error', 'Unknown error')
logger.error(f"[rollover] Extraction failed for {trigger}: {error_msg}")
# RESTORE from backup
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after extraction failure")
failed.append((trigger, "extraction", error_msg))
continue
memories = extract_result.get('entries', [])
branch = extract_result.get('branch', '')
memory_type = extract_result.get('type', 'unknown')
old_lines = extract_result.get('old_lines', 0)
new_lines = extract_result.get('new_lines', 0)
if not branch:
logger.error(f"[rollover] No branch found in extraction result for {trigger}")
failed.append((trigger, "extraction", "No branch in result"))
continue
logger.info(f"[rollover] Extracted {len(memories)} items from {trigger} ({old_lines} -> {new_lines} lines)")
# Convert memory items to text for vectorization
texts = _extract_text_from_memories(memories)
# Step 3: Generate embeddings
embed_result = embedder.encode_batch(texts)
if not embed_result['success']:
error_msg = embed_result.get('error', 'Unknown error')
logger.error(f"[rollover] Embedding failed for {trigger}: {error_msg}")
# RESTORE from backup
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after embedding failure")
failed.append((trigger, "embedding", error_msg))
continue
embeddings = embed_result.get('embeddings', [])
if not embeddings:
logger.error(f"[rollover] No embeddings generated for {trigger}")
failed.append((trigger, "embedding", "No embeddings in result"))
continue
logger.info(f"[rollover] Generated {len(embeddings)} embeddings for {trigger}")
# Step 4: Prepare metadata for vectorization
metadatas = []
for memory in memories:
metadata = memory.get('_metadata', {})
metadata['timestamp'] = memory.get('timestamp', '')
metadatas.append(metadata)
# Step 5: Store in LOCAL branch Chroma (via subprocess)
# Type assertions for Pylance (validated above with early returns)
branch_str: str = branch
memory_type_str: str = memory_type
embeddings_list: list = embeddings
local_chroma_path = _get_branch_local_chroma_path(branch_str)
local_store_result = None
if local_chroma_path:
local_store_result = _store_vectors_subprocess(
branch=branch_str,
memory_type=memory_type_str,
embeddings=embeddings_list,
documents=texts,
metadatas=metadatas,
db_path=str(local_chroma_path)
)
if not local_store_result['success']:
logger.warning(f"[rollover] Local storage failed for {branch}: {local_store_result.get('error')}")
# Continue anyway - global storage is primary
else:
logger.info(f"[rollover] Stored {len(embeddings)} vectors in local Chroma for {branch}")
# Step 6: Store in GLOBAL Memory Chroma (via subprocess)
global_store_result = _store_vectors_subprocess(
branch=branch_str,
memory_type=memory_type_str,
embeddings=embeddings_list,
documents=texts,
metadatas=metadatas
# db_path=None means global
)
if not global_store_result['success']:
error_msg = global_store_result.get('error', 'Unknown error')
logger.error(f"[rollover] Global storage failed for {trigger}: {error_msg}")
# RESTORE from backup (CRITICAL - file was modified but storage failed)
restore_result = extractor.restore_from_backup(trigger.file_path)
if restore_result['success']:
logger.info("[rollover] Restored from backup after storage failure")
else:
logger.error(f"[rollover] CRITICAL: Failed to restore from backup: {restore_result.get('error')}")
failed.append((trigger, "global_storage", error_msg))
continue
logger.info(f"[rollover] Stored {len(embeddings)} vectors in global Chroma for {branch}")
# Step 7: Update line count metadata
update_result = line_counter.update_line_count(trigger.file_path)
if update_result['success']:
logger.info(f"[rollover] Updated line count metadata for {trigger.file_path.name}")
else:
logger.warning(f"[rollover] Failed to update line count for {trigger.file_path.name}: {update_result.get('error')}")
# Success!
success_count += 1
global_collection = global_store_result.get('collection')
global_total = global_store_result.get('total_vectors')
# Report both local and global storage
local_status = "> local" if local_store_result and local_store_result['success'] else "x local"
# Display individual results
for item in result.get('results', []):
local_status = "> local" if item.get('local_stored') else "x local"
console.print(
f" [green]>[/green] Rolled over {len(memories)} items -> {global_collection} "
f"({old_lines} -> {new_lines} lines, global: {global_total} vectors, {local_status})"
f" [green]>[/green] Rolled over {item['memories_count']} items -> {item['global_collection']} "
f"({item['old_lines']} -> {item['new_lines']} lines, global: {item['global_total']} vectors, {local_status})"
)
logger.info(f"[rollover] Successfully rolled over {trigger}: {len(memories)} items, {old_lines} -> {new_lines} lines")
# Report results
success_count = result.get('success_count', 0)
failed = result.get('failed', [])
console.print()
if success_count > 0:
console.print(f"[green]>[/green] Rollover complete: {success_count}/{len(triggers)} successful")
logger.info(f"[rollover] Rollover complete: {success_count}/{len(triggers)} successful")
console.print(f"[green]>[/green] Rollover complete: {success_count}/{triggers_count} successful")
if failed:
console.print()
console.print("[red]Failed operations:[/red]")
for trigger, stage, err in failed:
console.print(f" [red]x[/red] {trigger} - {stage}: {err}")
logger.error(f"[rollover] {len(failed)} operations failed")
for fail in failed:
console.print(f" [red]x[/red] {fail['trigger']} - {fail['stage']}: {fail['error']}")
return success_count > 0
# =============================================================================
# PATH HELPERS
# =============================================================================
def _get_branch_local_chroma_path(branch_name: str) -> Path | None:
"""
Get local .chroma path for branch
Args:
branch_name: Branch name (e.g., "SEED", "AIPASS")
Returns:
Path to branch's local .chroma directory, or None if branch not found
"""
# Read registry to get branch path
if not branch_name:
return None
registry = detector._read_registry()
for branch in registry:
if branch.get('name', '').upper() == branch_name.upper():
branch_path = Path(branch.get('path', ''))
if branch_path.exists():
chroma_path = branch_path / '.chroma'
# Auto-create .chroma directory if missing
if not chroma_path.exists():
chroma_path.mkdir(parents=True, exist_ok=True)
logger.info(f"[rollover] Created local .chroma directory for {branch_name}")
return chroma_path
logger.warning(f"[rollover] Branch {branch_name} not found in registry")
return None
# =============================================================================
# TEXT EXTRACTION HELPERS
# =============================================================================
def _extract_text_from_memories(memories: List[Dict]) -> List[str]:
"""
Extract text content from memory items for vectorization
Memory items have different structures:
- sessions: 'activities' array (join into text)
- observations: might have 'content' or 'text' field
- generic: convert to JSON string
Args:
memories: List of memory items
Returns:
List of text strings for embedding
"""
texts = []
for memory in memories:
# Try common text fields
if 'activities' in memory and isinstance(memory['activities'], list):
# Sessions type - join activities
text = '\n'.join(str(a) for a in memory['activities'])
elif 'content' in memory:
text = str(memory['content'])
elif 'text' in memory:
text = str(memory['text'])
elif 'message' in memory:
text = str(memory['message'])
else:
# Fallback - convert to string representation
text = str(memory)
texts.append(text)
return texts
# =============================================================================
# LINE COUNT SYNC
# =============================================================================
@@ -491,8 +179,7 @@ def sync_line_counts() -> None:
"""
Update line count metadata for all branch memory files.
Reads actual line counts and updates document_metadata.status.current_lines
for all *.local.json and *.observations.json files in AIPASS_REGISTRY.
Delegates to handler and renders results with Rich.
"""
console.print()
console.print(Panel.fit(
@@ -505,7 +192,7 @@ def sync_line_counts() -> None:
console.print("[cyan]Updating line counts for all memory files...[/cyan]")
console.print()
result = line_counter.update_all_memory_files()
result = _handler_sync_line_counts()
if result['success']:
console.print(f"[green]>[/green] Updated {result['updated']} files")
@@ -513,10 +200,8 @@ def sync_line_counts() -> None:
console.print(f"[yellow]![/yellow] {result['failed']} files failed:")
for branch, mem_type, error in result.get('failures', []):
console.print(f" [red]x[/red] {branch}.{mem_type}: {error}")
logger.info(f"[rollover] Synced line counts: {result['updated']} updated, {result['failed']} failed")
else:
console.print("[red]x[/red] Failed to sync line counts")
logger.error("[rollover] Failed to sync line counts")
console.print()
@@ -525,11 +210,6 @@ def sync_line_counts() -> None:
# STATUS & CHECKING
# =============================================================================
def get_rollover_stats() -> dict:
"""Return rollover stats from detector (module-layer passthrough)."""
return detector.get_rollover_stats()
def show_status() -> None:
"""
Show rollover statistics for all branches
@@ -625,6 +305,26 @@ def check_triggers() -> None:
console.print()
# =============================================================================
# INTROSPECTION
# =============================================================================
def print_introspection():
"""Display module introspection info."""
console.print()
console.print("rollover Module")
console.print("Orchestrates memory rollover workflow: trigger detection, extraction, embedding, and vector storage")
console.print()
console.print("Connected Handlers:")
console.print(" handlers/monitor/")
console.print(" - detector.py (check_all_branches — detect branches exceeding rollover threshold)")
console.print(" - detector.py (get_rollover_stats — retrieve rollover statistics for all branches)")
console.print(" handlers/rollover/")
console.print(" - orchestrator.py (execute_rollover — run full rollover pipeline for triggered branches)")
console.print(" - orchestrator.py (sync_line_counts — update line count metadata for all memory files)")
console.print()
# =============================================================================
# STANDALONE EXECUTION
# =============================================================================
+55 -148
View File
@@ -1,19 +1,9 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: search.py - Search Orchestration Module
# Date: 2025-11-27
# Version: 0.2.0
# Category: memory/modules
#
# CHANGELOG (Max 5 entries):
# - v0.2.0 (2026-03-06): Adapted for AIPass public repo - removed internal deps
# - v0.1.0 (2025-11-27): Initial version - orchestrate semantic search
#
# CODE STANDARDS:
# - Thin orchestration: Delegate all logic to handlers
# - No business logic: Only coordinate workflow
# - handle_command() pattern
# =================== AIPass ====================
# Name: search.py
# Description: Search Orchestration Module
# Version: 0.4.0
# Created: 2025-11-27
# Modified: 2026-03-08
# =============================================
"""
@@ -30,75 +20,22 @@ Purpose:
"""
import sys
import os
import logging
import subprocess
import json
from pathlib import Path
from typing import List
from rich.console import Console
from rich.panel import Panel
from rich import box
from aipass.prax import logger
from aipass.cli.apps.modules import console
# =============================================================================
# INFRASTRUCTURE SETUP
# =============================================================================
logger = logging.getLogger(__name__)
console = Console()
# Handler imports (relative within the memory package)
from ..handlers.vector import embedder
# ChromaDB search via subprocess
_HANDLERS_DIR = Path(__file__).resolve().parent.parent / "handlers"
CHROMA_SUBPROCESS_SCRIPT = _HANDLERS_DIR / "storage" / "chroma_subprocess.py"
# Use system python by default; can be overridden via environment variable
MEMORY_PYTHON = os.environ.get("AIPASS_MEMORY_PYTHON", sys.executable)
def _search_vectors_subprocess(
query_embedding: list,
branch: str | None = None,
memory_type: str | None = None,
n_results: int = 5,
db_path: str | Path | None = None
) -> dict:
"""
Search vectors via subprocess.
This ensures ChromaDB compatibility regardless of calling Python version.
"""
input_data = {
'operation': 'search_vectors',
'query_embedding': query_embedding,
'branch': branch,
'memory_type': memory_type,
'n_results': n_results,
'db_path': str(db_path) if db_path else None
}
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 'Subprocess failed'}
return json.loads(result.stdout)
except subprocess.TimeoutExpired:
return {'success': False, 'error': 'Search operation timed out'}
except json.JSONDecodeError as e:
return {'success': False, 'error': f'Invalid JSON response: {e}'}
except Exception as e:
return {'success': False, 'error': str(e)}
# Handler imports
from aipass.memory.apps.handlers.search.query_executor import (
execute_search as _handler_execute_search,
)
# =============================================================================
@@ -160,7 +97,7 @@ def handle_command(command: str, args: List[str]) -> bool:
console.print("[red]Error:[/red] Search query required")
return True
execute_search(query, branch=branch, memory_type=memory_type, n_results=n_results)
show_search_results(query, branch=branch, memory_type=memory_type, n_results=n_results)
return True
return False
@@ -205,17 +142,17 @@ def print_help() -> None:
# =============================================================================
# SEARCH ORCHESTRATION
# SEARCH RESULTS DISPLAY
# =============================================================================
def execute_search(query: str, branch: str | None = None, memory_type: str | None = None, n_results: int = 5) -> bool:
def show_search_results(
query: str,
branch: str | None = None,
memory_type: str | None = None,
n_results: int = 5
) -> bool:
"""
Execute semantic search and display results
Workflow:
1. Encode query to embedding vector
2. Search ChromaDB via subprocess
3. Format and display results with Rich
Execute semantic search via handler and display results with Rich.
Args:
query: Search query text
@@ -234,7 +171,7 @@ def execute_search(query: str, branch: str | None = None, memory_type: str | Non
))
console.print()
# Step 1: Encode query
# Display query info
console.print(f"[cyan]Query:[/cyan] {query}")
if branch:
console.print(f"[cyan]Branch:[/cyan] {branch}")
@@ -243,52 +180,29 @@ def execute_search(query: str, branch: str | None = None, memory_type: str | Non
console.print()
console.print("[dim]Encoding query...[/dim]")
embed_result = embedder.encode_batch([query])
if not embed_result['success']:
error_msg = embed_result.get('error', 'Unknown error')
logger.error(f"[search] Failed to encode query: {error_msg}")
console.print(f"[red]x[/red] Failed to encode query: {error_msg}")
return False
embeddings = embed_result.get('embeddings', [])
if not embeddings:
console.print("[red]x[/red] No embedding generated")
return False
query_embedding = embeddings[0]
# Convert numpy array to list for JSON serialization
if hasattr(query_embedding, 'tolist'):
query_embedding = query_embedding.tolist()
logger.info(f"[search] Encoded query to {len(query_embedding)}-dim vector")
# Step 2: Search via subprocess
console.print("[dim]Searching collections...[/dim]")
search_result = _search_vectors_subprocess(
query_embedding=query_embedding,
# Delegate to handler
result = _handler_execute_search(
query=query,
branch=branch,
memory_type=memory_type,
n_results=n_results
)
if not search_result['success']:
error_msg = search_result.get('error', 'Unknown error')
logger.error(f"[search] Search failed: {error_msg}")
console.print(f"[red]x[/red] Search failed: {error_msg}")
if not result['success']:
console.print(f"[red]x[/red] {result.get('error', 'Unknown error')}")
return False
results = search_result.get('results', [])
collections_searched = search_result.get('collections_searched', 0)
total_results = search_result.get('total_results', 0)
collections_searched = result.get('collections_searched', 0)
total_results = result.get('total_results', 0)
filtered_results = result.get('results', [])
logger.info(f"[search] Found {total_results} results across {collections_searched} collections")
# Step 3: Display results
# Display summary
console.print(f"[green]>[/green] Found {total_results} results in {collections_searched} collections")
console.print()
if not results:
if not filtered_results and total_results == 0:
console.print("[yellow]No matching memories found[/yellow]")
console.print()
console.print("[dim]Try:[/dim]")
@@ -297,27 +211,6 @@ def execute_search(query: str, branch: str | None = None, memory_type: str | Non
console.print(" * Check if memories have been rolled over (drone @memory status)")
return True
# Minimum similarity threshold - filter out irrelevant results
MIN_SIMILARITY_THRESHOLD = 0.40 # 40% minimum relevance
# Filter and process results
filtered_results = []
for result in results[:n_results]:
document = result.get('document', '')
distance = result.get('distance', 0)
# Calculate similarity (ChromaDB L2 distance: 0=identical, ~2=very different)
similarity = max(0, 1 - (distance / 2))
# Skip empty documents and low-relevance results
if not document or not document.strip():
continue
if similarity < MIN_SIMILARITY_THRESHOLD:
continue
result['similarity'] = similarity
filtered_results.append(result)
if not filtered_results:
console.print("[yellow]No relevant memories found[/yellow]")
console.print()
@@ -325,11 +218,11 @@ def execute_search(query: str, branch: str | None = None, memory_type: str | Non
console.print("[dim]Try more specific search terms related to your AIPass work.[/dim]")
return True
for i, result in enumerate(filtered_results, 1):
collection = result.get('collection', 'unknown')
document = result.get('document', '')
metadata = result.get('metadata', {})
similarity = result.get('similarity', 0)
for i, item in enumerate(filtered_results, 1):
collection = item.get('collection', 'unknown')
document = item.get('document', '')
metadata = item.get('metadata', {})
similarity = item.get('similarity', 0)
# Parse collection name
parts = collection.split('_')
@@ -360,11 +253,25 @@ def execute_search(query: str, branch: str | None = None, memory_type: str | Non
))
console.print()
logger.info(f"[search] Displayed {len(filtered_results)} results")
return True
# =============================================================================
# INTROSPECTION
# =============================================================================
def print_introspection():
"""Display module introspection info."""
console.print()
console.print("search Module")
console.print("Orchestrates semantic search across memory collections via vector embeddings and ChromaDB")
console.print()
console.print("Connected Handlers:")
console.print(" handlers/search/")
console.print(" - query_executor.py (execute_search — encode query, search collections, filter results by similarity)")
console.print()
# =============================================================================
# STANDALONE EXECUTION
# =============================================================================