feat: port Memory Bank module to AIPass repo (FPLAN-0003)

Port internal Memory Bank system into src/aipass/memory/ with adapted
imports and paths. All 17 Python files compile clean, 13 existing tests
still pass.

Ported components:
- Entry point: apps/memory.py (drone @memory command routing)
- Modules: rollover.py (orchestration), search.py (query routing)
- Handlers: detector, extractor, line_counter, json_handler, indexer,
  embedder, chroma, vector_search, normalize, manager, dashboard_push,
  central_writer, memory_watcher, chroma_subprocess

Adaptations from internal system:
- Removed all /home/aipass/ hardcoded paths
- Removed sys.path manipulation hacks
- Replaced prax logger with stdlib logging
- Replaced cli console/header with Rich (already a dep)
- Changed BRANCH_REGISTRY → AIPASS_REGISTRY references
- Made chromadb/sentence-transformers optional (try/except)
- Removed private branch registry concept
- Subprocess Python uses sys.executable with AIPASS_MEMORY_PYTHON override

Still needs: drone registration, spawn template updates, container
testing, integration tests

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
AIOSAI
2026-03-06 22:49:53 -08:00
co-authored by Claude Opus 4.6
parent 08f84233ae
commit da4fa82771
35 changed files with 7488 additions and 0 deletions
@@ -0,0 +1,308 @@
# FPLAN-0003 - [framework] Memory Bank — port to AIPass repo
**Created**: 2026-03-06
**Branch**: /home/aipass/aipass_business/AIPass
**Status**: Active
**Type**: Standard Plan
---
## What Are Flow Plans?
Flow Plans (FPLANs) are for **BUILDING** - autonomous construction of systems, features, modules.
**This is NOT for:**
- Research or exploration (use agents directly)
- Quick fixes (just do it)
- Discussion or planning (that happens before creating the FPLAN)
**This IS for:**
- Building features or modules
- Single focused construction tasks
- Sub-plans within a master plan
---
## When to Use This vs Master Plan
| This (Default) | Master Plan |
|----------------|-------------|
| Single focused task | 3+ phases, complex build |
| Self-contained | Roadmap + multiple sub-plans |
| Quick build | Multi-session project |
| One phase of a master | Entire branch/system build |
**Need a master plan?** `drone @flow create "subject" master`
---
## Branch Directory Structure
Use dedicated directories - don't scatter files:
| Directory | Purpose |
|-----------|---------|
| `apps/` | Code (modules/, handlers/) |
| `tests/` | All test files |
| `tools/` | Utility scripts |
| `artifacts/` | Agent outputs |
| `docs/` | Documentation |
---
## Critical: Branch Manager Role
**You are the ORCHESTRATOR, not the builder.**
Your 200k context is precious. Burning it on file reads and code writing risks compaction during autonomous work. Agents have clean context - use them for ALL building.
| You Do (Orchestrator) | Agents Do (Builders) |
|-----------------------|----------------------|
| Create plans | Write code |
| Give instructions | Run tests |
| Review output | Read/modify files |
| Course correct | Research/exploration |
| Update memories | Heavy lifting |
| Send status emails | Single-task execution |
**Pattern:** Instruct agent → Wait for completion → Review output → Next step
---
## Seek Branch Expertise
Don't figure everything out alone. Other branches are domain experts - ask them first.
**Before building anything that touches another branch's domain:**
```bash
ai_mail send @branch "Question: [topic]" "I'm working on X and need guidance on Y. What's the best approach?"
```
**Common examples:**
- Building something with email? Ask @ai_mail how delivery works
- Need routing or @ resolution? Ask @drone
- Unsure about standards? Ask @seed for reference code
- Need persistent storage or search? Ask @memory_bank
- Event-driven behavior? Ask @trigger about their event system
- Dashboard integration? Ask @devpulse about update_section()
They have deep memory on their systems. A 1-email question saves you hours of guessing.
---
## Notepad
Keep `notepad.md` in your branch directory as a shared scratchpad during the build. Use it for:
- **Status updates** - Quick progress lines so Patrick can glance without asking
- **Questions for Patrick** - Non-urgent questions that can wait for his next visit
- **Notes to self** - Decisions made, things to revisit, gotchas discovered
Update it as you work - lightweight, not formal. Patrick checks it when he wants to, skips it when he's busy.
---
## Command Reference
When unsure about syntax, use `--help`:
```bash
# Flow - Plan management
drone @flow create . "subject" # Create plan (. = current dir)
drone @flow close FPLAN-XXXX # Close plan
drone @flow list # List active plans
drone @flow --help # Full help
# Seed - Quality gates
drone @seed checklist <file> # 10-point check on file
drone @seed audit @branch # Full branch audit
drone @seed --help # Full help
# AI_Mail - Status updates
drone @ai_mail send @dev_central "Subject" "Message"
drone @ai_mail --help # Full help
# Discovery
drone systems # All available modules
drone list @branch # Commands for branch
```
---
## Planning Phase
### Goal
Port the internal Memory Bank system into `src/aipass/memory/` as a fully functional module. Every spawned branch should have working memory rollover, JSONL archival, and searchable archives out of the box. Vector storage (ChromaDB + sentence-transformers) as optional pip extra.
### What We're Porting
Source: `/home/aipass/MEMORY_BANK/` (~1,800 lines, 75+ files)
**Core (zero new deps — must ship):**
- `detector.py` — threshold detection via AIPASS_REGISTRY.json (600-line default)
- `extractor.py` — FIFO extraction of oldest sessions, backup-before-extract
- `line_counter.py` — metadata `current_lines` updates
- `json_handler.py` — atomic JSON read/write (temp file + rename)
- `archiver.py` — JSONL archive write + keyword search
- `rollover.py` module — orchestrates detect → extract → archive → update
- `search.py` module — keyword search over JSONL archives
**Optional (`pip install aipass[vectors]`):**
- `embedder.py` — sentence-transformers wrapper (all-MiniLM-L6-v2, 384-dim)
- `vector_store.py` — ChromaDB persistence, per-branch `.chroma/` dirs
- `vector_search.py` — semantic search with metadata filtering
- Graceful fallback to keyword search if deps not installed
**Not v1 (later):**
- Symbolic/fragmented memory (experimental, v0.3)
- Living template push system
- Dashboard integration (needs DevPulse)
- Central writer (needs AI_CENTRAL concept)
- Memory pool intake processing
### Key Decisions
- **passport.json** is the identity file name (not id.json)
- **`.trinity/`** stays as the memory directory name
- **AIPASS_REGISTRY.json** is used for branch discovery (same as internal BRANCH_REGISTRY.json)
- **No daemon** — auto-check on command execution or manual `drone @memory rollover`
- **3-layer architecture** — apps/memory.py → modules/ → handlers/
- **JSONL for base archive** — keyword search good enough for most users
- **ChromaDB optional** — subprocess isolation if Python version issues arise
### Spawn Template Updates
- Welcome session (session 0) in local.json on branch creation
- `.trinity/archive/` directory pre-created and empty
- Memory self-managing — agent updates files, rollover handles overflow
### Integration Points
- Drone: `drone @memory status|rollover|search` commands
- Spawn: post-spawn welcome memory
- Registry: reads AIPASS_REGISTRY.json for branch discovery
- Hooks: PreCompact already references memory context
### Reference Documents
- Internal Memory Bank: `/home/aipass/MEMORY_BANK/apps/` (source code)
- Internal rollover module: `/home/aipass/MEMORY_BANK/apps/modules/rollover.py` (676 lines)
- Internal detector: `/home/aipass/MEMORY_BANK/apps/handlers/monitor/detector.py`
- Internal extractor: `/home/aipass/MEMORY_BANK/apps/handlers/rollover/extractor.py`
- Spawn templates: `src/aipass/spawn/templates/agent.template/.trinity/`
- AIPass architecture: 3-layer pattern (apps → modules → handlers)
---
## Agent Preparation (Before Deploying)
Agents can't work blind. They need context before they build.
**Your Prep Work (as orchestrator):**
1. [ ] Know where agent will work (branch path, key directories)
2. [ ] Identify files agent needs to reference or modify
3. [ ] Gather any specs, planning docs, or examples to include
4. [ ] Prepare COMPLETE instructions (agents are stateless)
**Agent's First Task (context building):**
- Agent should explore/read relevant files BEFORE writing code
- "First, read X and Y to understand the current structure"
- "Look at Z for the pattern to follow"
- Context-first, build-second
**What Agents DON'T Have:**
- No prior conversation history
- No memory files loaded automatically
- No knowledge of other branches
- Only what you put in their instructions
**Your instructions determine success - be thorough and specific.**
---
## Agent Instructions Template
```
You are working at [BRANCH_PATH].
TASK: [Specific single task]
CONTEXT:
- [What they need to know]
- Reference: [planning docs, existing code to study]
- First, READ the relevant files to understand current structure
DELIVERABLES:
- [Specific file or output expected]
- Tests → tests/
- Reports/logs → artifacts/reports/ or artifacts/logs/
CONSTRAINTS:
- Follow Seed standards (3-layer architecture)
- Do NOT modify files outside your task scope
- CROSS-BRANCH: Never modify other branches' files unless explicitly authorized by DEV_CENTRAL
- 2-ATTEMPT RULE: If something fails twice, note the issue and move on
- Do NOT go down rabbit holes debugging
WHEN COMPLETE:
- Verify code runs without syntax errors
- List files created/modified
- Note any issues encountered (with what was attempted)
```
---
## Execution Log
### 2026-03-06
- [ ] Created FPLAN-0003
- [ ] Agent deployed for: [task]
- [ ] Agent completed: [outcome]
- [ ] Seed checklist passed: [file]
- [ ] Memories updated
**Log Pattern:** Task → Agent → Outcome → Quality check → Next
**If production stops (critical blocker):**
```bash
drone @ai_mail send @dev_central "PRODUCTION STOPPED: FPLAN-0003" "Issue: [description]. Attempted: [what was tried]. Awaiting guidance."
```
---
## Notes
[Working notes, issues encountered, decisions made]
---
## Completion Checklist
### Before Closing
- [ ] All goals achieved
- [ ] Agent output reviewed and verified
- [ ] Seed checklist on new code: `drone @seed checklist <file>`
- [ ] Branch memories updated:
- [ ] `BRANCH.local.json` - session/work log
- [ ] `BRANCH.observations.json` - patterns learned (if any)
- [ ] README.md updated (if build changed status/capabilities)
- [ ] Status email sent to DEV_CENTRAL:
```bash
drone @ai_mail send @dev_central "FPLAN-0003 Complete" "Summary of what was done, any issues, outcomes"
```
**Completion Order:** Memories → README → Email (README before email - don't report complete with stale docs)
### Definition of Done
- `src/aipass/memory/` module exists with 3-layer architecture
- `drone @memory status` shows all branches and their memory file line counts
- `drone @memory rollover` detects and processes oversized files
- `drone @memory search "query"` returns keyword matches from JSONL archives
- Spawn creates branches with welcome session and empty archive dir
- All 13+ existing tests still pass
- New tests for memory module (detector, extractor, archiver, search)
- `pip install aipass[vectors]` adds ChromaDB + semantic search capability
- Memory module registered in AIPASS_REGISTRY.json
---
## Close Command
When all boxes checked:
```bash
drone @flow close FPLAN-0003
```
View File
View File
@@ -0,0 +1,340 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: indexer.py - Code Archive Indexer
# Date: 2025-11-27
# 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
# =============================================
"""
Code Archive Indexer
Indexes Python files in code_archive directory:
1. Scans for .py files
2. Extracts docstrings and metadata
3. Creates/updates index.json catalog
No vectorization - just a searchable catalog.
"""
import ast
import json
import logging
from pathlib import Path
from datetime import datetime
from typing import Dict, Any, List
logger = logging.getLogger(__name__)
# Paths resolved relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[4]
CODE_ARCHIVE_PATH = _MEMORY_ROOT / "code_archive"
INDEX_PATH = CODE_ARCHIVE_PATH / "index.json"
def extract_file_info(file_path: Path) -> Dict[str, Any]:
"""
Extract metadata from a Python file.
Args:
file_path: Path to Python file
Returns:
Dict with filename, docstring, functions, classes
"""
try:
content = file_path.read_text(encoding='utf-8')
# Parse AST
try:
tree = ast.parse(content)
except SyntaxError:
return {
'filename': file_path.name,
'path': str(file_path.relative_to(CODE_ARCHIVE_PATH)),
'docstring': None,
'error': 'Syntax error - could not parse',
'size': file_path.stat().st_size,
'indexed_at': datetime.now().isoformat()
}
# Get module docstring
docstring = ast.get_docstring(tree)
# Get function and class names
functions = []
classes = []
for node in ast.walk(tree):
if isinstance(node, ast.FunctionDef):
functions.append(node.name)
elif isinstance(node, ast.ClassDef):
classes.append(node.name)
return {
'filename': file_path.name,
'path': str(file_path.relative_to(CODE_ARCHIVE_PATH)),
'docstring': docstring[:200] + '...' if docstring and len(docstring) > 200 else docstring,
'functions': functions[:10], # Limit to first 10
'classes': classes[:10],
'size': file_path.stat().st_size,
'lines': len(content.splitlines()),
'indexed_at': datetime.now().isoformat()
}
except Exception as e:
return {
'filename': file_path.name,
'path': str(file_path),
'error': str(e),
'indexed_at': datetime.now().isoformat()
}
def get_archive_files() -> List[Path]:
"""
Get all Python files in code_archive.
Returns:
List of Path objects for all .py files
"""
if not CODE_ARCHIVE_PATH.exists():
return []
files = list(CODE_ARCHIVE_PATH.rglob('*.py'))
# Exclude __init__.py files
files = [f for f in files if f.name != '__init__.py']
return sorted(files)
def load_index() -> Dict[str, Any]:
"""
Load existing index.json or create empty structure.
Returns:
Index dict with metadata and files
"""
if INDEX_PATH.exists():
try:
with open(INDEX_PATH) as f:
return json.load(f)
except Exception:
pass
return {
'metadata': {
'name': 'Code Archive Index',
'description': 'Catalog of archived Python modules from old system',
'created': datetime.now().isoformat(),
'last_updated': None,
'total_files': 0
},
'categories': {},
'files': {}
}
def save_index(index: Dict[str, Any]) -> Dict[str, Any]:
"""
Save index.json.
Args:
index: Index dict to save
Returns:
Dict with success status
"""
try:
index['metadata']['last_updated'] = datetime.now().isoformat()
index['metadata']['total_files'] = len(index['files'])
with open(INDEX_PATH, 'w') as f:
json.dump(index, f, indent=2)
return {'success': True}
except Exception as e:
return {'success': False, 'error': str(e)}
def build_index() -> Dict[str, Any]:
"""
Build complete index from scratch.
Scans all Python files and creates full index.json.
Returns:
Dict with success status and stats
"""
files = get_archive_files()
if not files:
return {
'success': True,
'message': 'No files to index',
'files_indexed': 0
}
index = load_index()
categories = {}
for file_path in files:
info = extract_file_info(file_path)
# Use relative path as key
key = info['path']
index['files'][key] = info
# Track categories (subdirectories)
category = file_path.parent.name
if category != 'code_archive':
if category not in categories:
categories[category] = []
categories[category].append(info['filename'])
index['categories'] = categories
save_result = save_index(index)
if not save_result['success']:
return save_result
return {
'success': True,
'files_indexed': len(files),
'categories': list(categories.keys())
}
def check_for_new_files() -> Dict[str, Any]:
"""
Sync index with actual files.
Adds new files, removes deleted files from index.
Returns:
Dict with changes made
"""
index = load_index()
current_files = get_archive_files()
indexed_paths = set(index.get('files', {}).keys())
current_paths = {str(f.relative_to(CODE_ARCHIVE_PATH)) for f in current_files}
new_files = current_paths - indexed_paths
deleted_files = indexed_paths - current_paths
if not new_files and not deleted_files:
return {
'success': True,
'new_files': 0,
'deleted_files': 0,
'action': 'none'
}
# Index new files
for rel_path in new_files:
file_path = CODE_ARCHIVE_PATH / rel_path
if file_path.exists():
info = extract_file_info(file_path)
index['files'][rel_path] = info
# Update category
category = file_path.parent.name
if category != 'code_archive':
if category not in index['categories']:
index['categories'][category] = []
if info['filename'] not in index['categories'][category]:
index['categories'][category].append(info['filename'])
# Remove deleted files from index
for rel_path in deleted_files:
if rel_path in index['files']:
filename = index['files'][rel_path].get('filename')
del index['files'][rel_path]
# Clean up category
for cat, files in index['categories'].items():
if filename in files:
files.remove(filename)
# Rebuild categories from current files
index['categories'] = {}
for rel_path, info in index['files'].items():
file_path = CODE_ARCHIVE_PATH / rel_path
category = file_path.parent.name
if category != 'code_archive':
if category not in index['categories']:
index['categories'][category] = []
if info['filename'] not in index['categories'][category]:
index['categories'][category].append(info['filename'])
save_index(index)
return {
'success': True,
'new_files': len(new_files),
'deleted_files': len(deleted_files),
'files_added': list(new_files) if new_files else None,
'files_removed': list(deleted_files) if deleted_files else None,
'action': 'synced'
}
def get_index_status() -> Dict[str, Any]:
"""
Get current index status.
Returns:
Dict with file counts, categories, last update
"""
index = load_index()
current_files = get_archive_files()
indexed_count = len(index.get('files', {}))
current_count = len(current_files)
return {
'indexed_files': indexed_count,
'current_files': current_count,
'unindexed': current_count - indexed_count if current_count > indexed_count else 0,
'categories': list(index.get('categories', {}).keys()),
'last_updated': index.get('metadata', {}).get('last_updated')
}
# Standalone execution
if __name__ == "__main__":
import sys
if len(sys.argv) > 1:
cmd = sys.argv[1]
if cmd == 'status':
status = get_index_status()
print(json.dumps(status, indent=2))
elif cmd == 'build':
print("Building index...")
result = build_index()
print(json.dumps(result, indent=2))
elif cmd == 'check':
result = check_for_new_files()
print(json.dumps(result, indent=2))
else:
print(f"Unknown command: {cmd}")
print("Usage: indexer.py [status|build|check]")
else:
print("Usage: indexer.py [status|build|check]")
print(" status - Show index status")
print(" build - Build complete index from scratch")
print(" check - Check for and index new files")
@@ -0,0 +1,337 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: central_writer.py - Central File Writer Handler
# Date: 2025-11-27
# 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
# =============================================
"""
Central File Writer Handler
Updates memory_bank.central.json with current statistics.
This file is Memory Bank's "API output" - used to populate dashboards.
Purpose:
Update central registry when vectors/archives change
Provide current stats to dashboard
"""
import logging
from json import load as json_load, dump as json_dump
from pathlib import Path
from datetime import datetime
from typing import Dict, Any
logger = logging.getLogger(__name__)
# No service imports - handlers are pure workers (3-tier architecture)
# =============================================================================
# CONSTANTS
# =============================================================================
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
CENTRAL_FILE = _MEMORY_ROOT / "central" / "memory_bank.central.json"
CHROMA_DB_PATH = _MEMORY_ROOT / ".chroma"
ARCHIVE_DIR = _MEMORY_ROOT / ".archive"
# =============================================================================
# STATS COLLECTION
# =============================================================================
def count_chroma_vectors() -> int:
"""
Count total vectors across all ChromaDB collections
Returns:
Total number of vectors stored (estimated from SQLite)
Raises:
Exception if database access fails
"""
try:
if not CHROMA_DB_PATH.exists():
return 0
# Read ChromaDB SQLite database directly to avoid dependency issues
import sqlite3
db_file = CHROMA_DB_PATH / "chroma.sqlite3"
if not db_file.exists():
return 0
conn = sqlite3.connect(str(db_file))
cursor = conn.cursor()
# Query embeddings table for total count
# ChromaDB stores embeddings in the 'embeddings' table
cursor.execute("SELECT COUNT(*) FROM embeddings")
total = cursor.fetchone()[0]
conn.close()
return total
except Exception:
return 0
def count_archive_files() -> int:
"""
Count total archive files in .archive directory
Returns:
Number of archived files
Raises:
Exception if directory access fails
"""
try:
if not ARCHIVE_DIR.exists():
return 0
# Count markdown files (archived memory files)
archive_files = list(ARCHIVE_DIR.glob("*.md"))
return len(archive_files)
except Exception as e:
raise Exception(f"Failed to count archive files: {e}")
def get_last_rollover_timestamp() -> str:
"""
Get timestamp of last rollover operation
Returns:
ISO timestamp string, or empty string if no rollover data
Raises:
Exception if unable to determine timestamp
"""
try:
# Check for most recent archive file
if not ARCHIVE_DIR.exists():
return ""
archive_files = list(ARCHIVE_DIR.glob("*.md"))
if not archive_files:
return ""
# Get most recent file by modification time
latest = max(archive_files, key=lambda f: f.stat().st_mtime)
mtime = datetime.fromtimestamp(latest.stat().st_mtime)
return mtime.isoformat()
except Exception as e:
raise Exception(f"Failed to get rollover timestamp: {e}")
def collect_stats() -> Dict[str, Any]:
"""
Collect all Memory Bank statistics
Returns:
Dict with stats:
- total_vectors: int
- total_archives: int
- last_rollover: ISO timestamp string
Raises:
Exception if any stat collection fails
"""
return {
"total_vectors": count_chroma_vectors(),
"total_archives": count_archive_files(),
"last_rollover": get_last_rollover_timestamp()
}
# =============================================================================
# CENTRAL FILE OPERATIONS
# =============================================================================
def read_central_file() -> Dict[str, Any]:
"""
Read current central file contents
Returns:
Dict with central file data
Raises:
Exception if file read fails
"""
try:
if not CENTRAL_FILE.exists():
# Return default structure
return {
"service": "memory_bank",
"last_updated": "",
"stats": {
"total_vectors": 0,
"total_archives": 0,
"last_rollover": ""
}
}
with open(CENTRAL_FILE, 'r', encoding='utf-8') as f:
return json_load(f)
except Exception as e:
raise Exception(f"Failed to read central file: {e}")
def write_central_file(data: Dict[str, Any]) -> None:
"""
Write data to central file
Args:
data: Dict with central file structure
Raises:
Exception if file write fails
"""
try:
# Ensure directory exists
CENTRAL_FILE.parent.mkdir(parents=True, exist_ok=True)
with open(CENTRAL_FILE, 'w', encoding='utf-8') as f:
json_dump(data, f, indent=2, ensure_ascii=False)
except Exception as e:
raise Exception(f"Failed to write central file: {e}")
# =============================================================================
# PUBLIC API
# =============================================================================
def update_central(verbose: bool = False) -> Dict[str, Any]:
"""
Update memory central.json with current statistics
Main function called by memory modules to update central registry.
Called after vector storage, rollover, or archive operations.
Args:
verbose: If True, include detailed stats in return dict
Returns:
Dict with update result:
success: bool
stats: dict (if verbose=True)
error: str (if failed)
Example:
# Update central file after rollover
result = update_central(verbose=True)
if result['success']:
print(f"Updated with {result['stats']['total_vectors']} vectors")
"""
try:
# Collect current stats
stats = collect_stats()
# Read existing file (preserves any extra fields)
central_data = read_central_file()
# Update stats and timestamp
central_data["last_updated"] = datetime.now().isoformat()
central_data["stats"] = stats
# Remove placeholder note if present
if "_note" in central_data:
del central_data["_note"]
# Write updated file
write_central_file(central_data)
result = {
"success": True,
"updated": CENTRAL_FILE.as_posix()
}
if verbose:
result["stats"] = stats
return result
except Exception as e:
return {
"success": False,
"error": str(e)
}
def get_current_stats() -> Dict[str, Any]:
"""
Get current Memory Bank statistics without writing to file
Returns:
Dict with stats or error
Example:
# Check stats without updating file
stats = get_current_stats()
if stats['success']:
print(f"Vectors: {stats['total_vectors']}")
"""
try:
stats = collect_stats()
return {
"success": True,
**stats
}
except Exception as e:
return {
"success": False,
"error": str(e)
}
# =============================================================================
# CLI ENTRY POINT (for testing)
# =============================================================================
if __name__ == "__main__":
print("\n=== Memory Bank Central Writer ===\n")
# Get current stats
print("Collecting statistics...")
stats_result = get_current_stats()
if stats_result['success']:
print(f" Total Vectors: {stats_result.get('total_vectors', 0)}")
print(f" Total Archives: {stats_result.get('total_archives', 0)}")
print(f" Last Rollover: {stats_result.get('last_rollover', 'Never')}")
else:
print(f"Error collecting stats: {stats_result.get('error')}")
print()
# Update central file
print("Updating central file...")
result = update_central(verbose=True)
if result['success']:
print(f"Updated: {result['updated']}")
if 'stats' in result:
print(f" Vectors: {result['stats']['total_vectors']}")
print(f" Archives: {result['stats']['total_archives']}")
else:
print(f"Update failed: {result.get('error')}")
print()
@@ -0,0 +1,408 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: dashboard_push.py - Memory Bank Dashboard Write-Through
# Date: 2026-02-25
# 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
# =============================================
"""
Memory Bank Dashboard Write-Through Handler
Pushes the memory_bank section to branch dashboards.
Called after rollover, pool processing, and plans processing to keep dashboards
showing accurate vector counts, collection stats, and near-rollover warnings.
This is a SYSTEM-WIDE push: Memory Bank data is global, so it pushes to ALL
branch dashboards (every branch benefits from knowing system memory health).
"""
import sys
import logging
import subprocess
from json import loads as json_loads
from pathlib import Path
from datetime import datetime
from typing import Dict, Any, List
logger = logging.getLogger(__name__)
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[3]
# =============================================================================
# CONSTANTS
# =============================================================================
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"
# Near-rollover threshold: branches with fewer than this many lines remaining
NEAR_ROLLOVER_THRESHOLD = 100
# =============================================================================
# DATA COLLECTION
# =============================================================================
def _read_central_stats() -> Dict[str, Any]:
"""
Read total_vectors and related stats from memory central.json.
Returns:
Dict with total_vectors, total_archives, last_rollover
"""
try:
if not CENTRAL_FILE.exists():
return {"total_vectors": 0, "total_archives": 0, "last_rollover": ""}
data = json_loads(CENTRAL_FILE.read_text(encoding="utf-8"))
stats = data.get("stats", {})
return {
"total_vectors": stats.get("total_vectors", 0),
"total_archives": stats.get("total_archives", 0),
"last_rollover": stats.get("last_rollover", "")
}
except Exception:
return {"total_vectors": 0, "total_archives": 0, "last_rollover": ""}
def _get_collections_count() -> int:
"""
Count ChromaDB collections by reading the SQLite database directly.
Returns:
Number of collections in the global Chroma database
"""
try:
import sqlite3
db_file = _MEMORY_ROOT / ".chroma" / "chroma.sqlite3"
if not db_file.exists():
return 0
conn = sqlite3.connect(str(db_file))
cursor = conn.cursor()
cursor.execute("SELECT COUNT(*) FROM collections")
count = cursor.fetchone()[0]
conn.close()
return count
except Exception:
return 0
def _get_rollover_config() -> Dict[str, Any]:
"""
Load rollover configuration (defaults + per-branch overrides).
Returns:
Dict with 'defaults' and 'per_branch' rollover config
"""
try:
if not CONFIG_PATH.exists():
return {"defaults": {"max_lines": 600, "buffer": 100}, "per_branch": {}}
data = json_loads(CONFIG_PATH.read_text(encoding="utf-8"))
rollover = data.get("rollover", {})
return {
"defaults": rollover.get("defaults", {"max_lines": 600, "buffer": 100}),
"per_branch": rollover.get("per_branch", {})
}
except Exception:
return {"defaults": {"max_lines": 600, "buffer": 100}, "per_branch": {}}
def _get_max_lines_for_branch(branch_name: str, rollover_config: Dict) -> int:
"""
Get the max_lines limit for a specific branch, respecting per-branch overrides.
Args:
branch_name: Uppercase branch name
rollover_config: Rollover config dict from _get_rollover_config()
Returns:
Max lines for this branch
"""
per_branch = rollover_config.get("per_branch", {})
if branch_name in per_branch:
return per_branch[branch_name].get("max_lines", rollover_config["defaults"]["max_lines"])
return rollover_config["defaults"]["max_lines"]
def _find_branches_near_rollover() -> List[Dict[str, Any]]:
"""
Scan all branches to find those near their rollover threshold.
Reads document_metadata.status.current_lines from each *.local.json
and *.observations.json, compares against max_lines limit.
Returns:
List of dicts with branch, file_type, lines_remaining
"""
near_rollover: List[Dict[str, Any]] = []
try:
if not AIPASS_REGISTRY.exists():
return near_rollover
registry = json_loads(AIPASS_REGISTRY.read_text(encoding="utf-8"))
branches = registry.get("branches", [])
rollover_config = _get_rollover_config()
for branch in branches:
branch_name = branch.get("name", "")
branch_path = Path(branch.get("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
for suffix in ["local", "observations"]:
memory_file = branch_path / f"{branch_name}.{suffix}.json"
if not memory_file.exists():
continue
try:
data = json_loads(memory_file.read_text(encoding="utf-8"))
doc_meta = data.get("document_metadata", {})
status = doc_meta.get("status", {})
current_lines = status.get("current_lines")
if current_lines is None:
continue
remaining = max_lines - current_lines
if remaining < NEAR_ROLLOVER_THRESHOLD:
near_rollover.append({
"branch": branch_name,
"file_type": suffix,
"lines_remaining": max(remaining, 0),
"current_lines": current_lines,
"max_lines": max_lines
})
except Exception:
# Skip files that can't be read
continue
except Exception:
return near_rollover # Return partial results on registry read failure
# Sort by lines_remaining ascending (most urgent first)
near_rollover.sort(key=lambda x: x["lines_remaining"])
return near_rollover
def _get_template_version() -> str:
"""
Read the current template version from .template_version.json.
Returns:
Version string (e.g., "2.0.4") or "unknown"
"""
try:
if not TEMPLATE_VERSION_FILE.exists():
return "unknown"
data = json_loads(TEMPLATE_VERSION_FILE.read_text(encoding="utf-8"))
return data.get("version", "unknown")
except Exception:
return "unknown"
def _get_last_rollover_info(central_stats: Dict) -> Dict[str, str]:
"""
Build last_rollover info from central stats.
Args:
central_stats: Stats dict from _read_central_stats()
Returns:
Dict with 'date' (and optionally other info)
"""
last_rollover_ts = central_stats.get("last_rollover", "")
if last_rollover_ts:
# Extract date portion from ISO timestamp
try:
dt = datetime.fromisoformat(last_rollover_ts)
return {"date": dt.strftime("%Y-%m-%d")}
except (ValueError, TypeError):
return {"date": last_rollover_ts}
return {"date": "never"}
def _get_all_branch_paths() -> List[Path]:
"""
Get paths for all active branches from AIPASS_REGISTRY.json.
Returns:
List of Path objects for all registered branches
"""
try:
if not AIPASS_REGISTRY.exists():
return []
registry = json_loads(AIPASS_REGISTRY.read_text(encoding="utf-8"))
paths = []
for branch in registry.get("branches", []):
branch_path = Path(branch.get("path", ""))
if branch_path.exists():
paths.append(branch_path)
return paths
except Exception:
return []
# =============================================================================
# PUBLIC API
# =============================================================================
def build_memory_bank_section() -> Dict[str, Any]:
"""
Build the memory_bank dashboard section data.
Collects all stats and returns the section dict ready for write_section().
Returns:
Dict with total_vectors, collections_count, branches_near_rollover,
last_rollover, template_version
"""
central_stats = _read_central_stats()
near_rollover = _find_branches_near_rollover()
last_rollover = _get_last_rollover_info(central_stats)
template_version = _get_template_version()
collections_count = _get_collections_count()
return {
"managed_by": "memory_bank",
"total_vectors": central_stats.get("total_vectors", 0),
"collections_count": collections_count,
"branches_near_rollover": near_rollover,
"last_rollover": last_rollover,
"template_version": template_version
}
def _write_section_to_all_branches(section_name: str, section_data: Dict,
branch_paths: List[Path]) -> int:
"""
Write a dashboard section to multiple branches via a single subprocess.
Uses one subprocess call for all branches to avoid spawning 29 processes.
The subprocess imports devpulse write_section and iterates all paths.
Args:
section_name: Dashboard section key (e.g., "memory_bank")
section_data: Section data dict to write
branch_paths: List of branch root directory paths
Returns:
Number of branches successfully updated
"""
try:
from json import dumps as json_dumps
# Single subprocess handles all branches in a loop
# Write dashboard section as JSON directly to DASHBOARD.local.json
script = (
"import sys, json\n"
"from pathlib import Path\n"
"data = json.loads(sys.stdin.read())\n"
"section_name = data['section_name']\n"
"section_data = data['section_data']\n"
"ok = 0\n"
"for bp in data['branch_paths']:\n"
" try:\n"
" dash = Path(bp) / 'DASHBOARD.local.json'\n"
" if dash.exists():\n"
" d = json.loads(dash.read_text())\n"
" d[section_name] = section_data\n"
" dash.write_text(json.dumps(d, indent=2))\n"
" ok += 1\n"
" except Exception:\n"
" continue\n"
"print(ok)\n"
)
input_data = json_dumps({
"section_name": section_name,
"section_data": section_data,
"branch_paths": [str(p) for p in branch_paths]
})
result = subprocess.run(
[sys.executable, "-c", script],
input=input_data,
capture_output=True,
text=True,
timeout=60
)
if result.returncode == 0 and result.stdout.strip().isdigit():
return int(result.stdout.strip())
return 0
except Exception:
return 0
def push_memory_bank_dashboard() -> bool:
"""
Push the memory_bank section to ALL branch dashboards.
This is the main entry point called after rollover, pool processing,
and plans processing. Memory Bank data is global so all branches
benefit from seeing system memory health.
Uses a single subprocess to call devpulse write_section() for all
branches, avoiding cross-package handler imports. Dashboard write
failures are silent - this is a best-effort operation that must not
break the calling workflow.
Returns:
True if at least one dashboard was updated, False on total failure
"""
try:
# Build the section data once
section_data = build_memory_bank_section()
# Push to all branch dashboards via single subprocess
branch_paths = _get_all_branch_paths()
success_count = _write_section_to_all_branches(
"memory_bank", section_data, branch_paths
)
return success_count > 0
except Exception:
return False
# =============================================================================
# CLI ENTRY POINT (for testing)
# =============================================================================
if __name__ == "__main__":
import json
print("Building memory_bank dashboard section...")
section = build_memory_bank_section()
print(json.dumps(section, indent=2))
print()
print("Pushing to all branch dashboards...")
result = push_memory_bank_dashboard()
print(f"Result: {'success' if result else 'failed'}")
@@ -0,0 +1,23 @@
"""
Memory File JSON Handler Package
Safe read/write operations for branch memory files.
"""
from .json_handler import (
read_memory_file,
write_memory_file,
update_metadata,
read_memory_file_data,
write_memory_file_simple,
validate_memory_file_structure
)
__all__ = [
'read_memory_file',
'write_memory_file',
'update_metadata',
'read_memory_file_data',
'write_memory_file_simple',
'validate_memory_file_structure'
]
@@ -0,0 +1,377 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: json_handler.py - Memory File Safe Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Memory File JSON Handler
Safe read/write operations for branch memory files (*.local.json, *.observations.json).
Handles the three-JSON pattern with atomic writes and proper error handling.
Purpose:
Prevent corruption of critical memory files during read/write operations.
All memory file access should use these functions instead of direct json.load/dump.
Features:
- Atomic writes (temp file + rename)
- Safe error handling
- Metadata helpers
- Preserves formatting (indent=2, ensure_ascii=False)
Usage:
from aipass.memory.apps.handlers.json.json_handler import read_memory_file, write_memory_file
"""
import json
import logging
import tempfile
from pathlib import Path
from typing import Dict, Any, Optional
from datetime import datetime
logger = logging.getLogger(__name__)
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[4]
_CONFIG_DIR = _MEMORY_ROOT / "config"
_TEMPLATES_DIR = _MEMORY_ROOT / "apps" / "json_templates"
# No service imports - handlers are pure workers (3-tier architecture)
# No module imports (handler independence)
# =============================================================================
# CORE READ/WRITE OPERATIONS
# =============================================================================
def read_memory_file(file_path: Path) -> Dict[str, Any]:
"""
Safe read of memory JSON file
Handles file not found, corrupt JSON, and other read errors gracefully.
Args:
file_path: Path to memory JSON file
Returns:
Dict with success status and data: {'success': True, 'data': {...}} or {'success': False, 'error': '...'}
Example:
result = read_memory_file(Path("path/to/BRANCH.local.json"))
if result['success']:
data = result['data']
sessions = data.get('sessions', [])
"""
if not file_path.exists():
return {
'success': False,
'error': f"File not found: {file_path}"
}
try:
with open(file_path, 'r', encoding='utf-8') as f:
data = json.load(f)
return {
'success': True,
'file': str(file_path),
'data': data
}
except json.JSONDecodeError as e:
return {
'success': False,
'error': f"Corrupt JSON in {file_path.name}: {e}"
}
except PermissionError:
return {
'success': False,
'error': f"Permission denied reading {file_path.name}"
}
except Exception as e:
return {
'success': False,
'error': f"Failed to read {file_path.name}: {e}"
}
def write_memory_file(file_path: Path, data: Dict[str, Any]) -> Dict[str, Any]:
"""
Atomic write of memory JSON file
Uses temp file + rename strategy to prevent corruption on write failures.
Preserves formatting (indent=2, ensure_ascii=False) for readability.
Args:
file_path: Path to memory JSON file
data: JSON data to write (must be dict)
Returns:
Dict with success status: {'success': True, 'file': '...'} or {'success': False, 'error': '...'}
Example:
result = read_memory_file(path)
if result['success']:
data = result['data']
data['sessions'].append(new_session)
write_result = write_memory_file(path, data)
Safety:
- Writes to temp file first
- Only renames if write succeeds
- Original file unchanged if write fails
"""
if not isinstance(data, dict):
return {
'success': False,
'error': f"Data must be dict, got {type(data).__name__}"
}
try:
# Create temp file in same directory (for atomic rename)
temp_fd, temp_path = tempfile.mkstemp(
dir=file_path.parent,
prefix=f".{file_path.name}.",
suffix=".tmp"
)
try:
# Write to temp file
with open(temp_fd, 'w', encoding='utf-8') as f:
json.dump(data, f, indent=2, ensure_ascii=False)
f.write('\n') # Add final newline
# Atomic rename (overwrites original)
Path(temp_path).rename(file_path)
return {
'success': True,
'file': str(file_path)
}
except Exception as e:
# Clean up temp file on failure
Path(temp_path).unlink(missing_ok=True)
raise e
except PermissionError:
return {
'success': False,
'error': f"Permission denied writing {file_path.name}"
}
except Exception as e:
return {
'success': False,
'error': f"Failed to write {file_path.name}: {e}"
}
# =============================================================================
# METADATA HELPERS
# =============================================================================
def update_metadata(
file_path: Path,
**updates
) -> Dict[str, Any]:
"""
Update document_metadata.status fields
Convenient helper for updating metadata without manual read-modify-write.
Only updates document_metadata.status fields, preserves rest of file.
Args:
file_path: Path to memory JSON file
**updates: Key-value pairs to update in status section
Returns:
Dict with success status: {'success': True} or {'success': False, 'error': '...'}
Example:
# Update health and line count
update_metadata(
path,
health="healthy",
current_lines=450,
last_health_check="2025-11-16"
)
Safety:
- Uses atomic write
- Creates metadata structure if missing
- Preserves all other data
"""
# Read current data
read_result = read_memory_file(file_path)
if not read_result['success']:
return read_result
data = read_result['data']
# Ensure metadata structure exists
if 'document_metadata' not in data:
data['document_metadata'] = {}
if 'status' not in data['document_metadata']:
data['document_metadata']['status'] = {}
# Apply updates
status = data['document_metadata']['status']
for key, value in updates.items():
status[key] = value
# Write back
write_result = write_memory_file(file_path, data)
return write_result
# =============================================================================
# CONVENIENCE FUNCTIONS
# =============================================================================
def read_memory_file_data(file_path: Path) -> Optional[Dict[str, Any]]:
"""
Read memory file and return data directly (no dict wrapper)
Convenience function for simple reads where you just need the data.
Returns None on any error.
Args:
file_path: Path to memory JSON file
Returns:
Parsed JSON data dict, or None on error
Example:
data = read_memory_file_data(path)
if data:
sessions = data.get('sessions', [])
"""
result = read_memory_file(file_path)
if result['success']:
return result.get('data')
return None
def write_memory_file_simple(file_path: Path, data: Dict[str, Any]) -> bool:
"""
Write memory file and return simple success/failure boolean
Convenience function for simple writes where you just need success flag.
Args:
file_path: Path to memory JSON file
data: JSON data to write
Returns:
True if successful, False otherwise
Example:
data['sessions'].append(new_session)
if write_memory_file_simple(path, data):
print("Success!")
"""
result = write_memory_file(file_path, data)
return result['success']
# =============================================================================
# VALIDATION HELPERS
# =============================================================================
def validate_memory_file_structure(data: Dict[str, Any]) -> tuple[bool, str]:
"""
Validate memory file has required structure
Checks for document_metadata presence and basic structure.
Args:
data: Parsed JSON data
Returns:
Tuple of (is_valid, error_message)
Example:
data = read_memory_file_data(path)
valid, error = validate_memory_file_structure(data)
if not valid:
logger.warning(f"Invalid structure: {error}")
"""
if not isinstance(data, dict):
return False, "Data is not a dictionary"
if 'document_metadata' not in data:
return False, "Missing 'document_metadata' field"
metadata = data['document_metadata']
if not isinstance(metadata, dict):
return False, "'document_metadata' is not a dictionary"
# Check for expected fields
expected = ['document_type', 'document_name', 'version']
missing = [field for field in expected if field not in metadata]
if missing:
return False, f"Missing metadata fields: {', '.join(missing)}"
return True, ""
# =============================================================================
# TESTING
# =============================================================================
if __name__ == "__main__":
import sys as _sys
print("\n=== MEMORY FILE JSON HANDLER - Safe Operations Test ===\n")
# Test with a file passed as argument, or show usage
if len(_sys.argv) > 1:
test_file = Path(_sys.argv[1])
else:
print("Usage: python json_handler.py <path_to_memory_file.json>")
print("No file specified, exiting.")
_sys.exit(0)
if test_file.exists():
print(f"[TEST] Reading {test_file.name}...")
result = read_memory_file(test_file)
if result['success']:
file_data = result['data']
print("+ Read successful")
print(f" Document type: {file_data.get('document_metadata', {}).get('document_type')}")
file_status = file_data.get('document_metadata', {}).get('status', {}).get('health')
print(f" Status: {file_status}")
# Validate structure
valid, error = validate_memory_file_structure(file_data)
if valid:
print("+ Structure validation passed")
else:
print(f"- Structure validation failed: {error}")
else:
print(f"- Read failed: {result['error']}")
else:
print(f"\n[TEST] Test file not found: {test_file}")
print()
@@ -0,0 +1,8 @@
"""
Key Learnings Management Handlers
Manages key_learnings section in branch .local.json files:
- Timestamp tracking for age-based pruning
- Configurable max_entries limit
- Vectorization before removal
"""
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,2 @@
# Minimal __init__.py - handler independence pattern
# No exports - handlers import directly when needed
@@ -0,0 +1,369 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: detector.py - Rollover Trigger Detection Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Rollover Trigger Detection Handler
Monitors branch memory files via AIPASS_REGISTRY.json and detects when
files exceed their max_lines threshold (typically 600 lines).
Purpose:
Detect rollover conditions without active monitoring. Called by
rollover module to check all branches for files needing rollover.
Independence:
No module imports - pure handler, transportable
"""
import json
import logging
from pathlib import Path
from typing import List, Dict, Any
from dataclasses import dataclass
logger = logging.getLogger(__name__)
# No service imports - handlers are pure workers (3-tier architecture)
# No module imports (handler independence)
# =============================================================================
# DATA STRUCTURES
# =============================================================================
@dataclass
class RolloverTrigger:
"""Represents a file that needs rollover"""
branch: str
memory_type: str # 'observations' or 'local'
file_path: Path
current_lines: int
max_lines: int
def __str__(self):
return f"{self.branch}.{self.memory_type} ({self.current_lines}/{self.max_lines} lines)"
# =============================================================================
# REGISTRY OPERATIONS
# =============================================================================
def _read_registry() -> List[Dict[str, Any]]:
"""
Read AIPASS_REGISTRY.json
Returns:
List of branch dictionaries
"""
registry_path = Path.home() / "AIPASS_REGISTRY.json"
if not registry_path.exists():
return []
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
return []
def _get_memory_file_path(branch: Dict, memory_type: str) -> Path | None:
"""
Get path to memory file for branch
Args:
branch: Branch dict from registry
memory_type: 'observations' or 'local'
Returns:
Path to memory file, or None if not found
"""
branch_path = Path(branch.get('path', ''))
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
return file_path if file_path.exists() else None
# =============================================================================
# CONFIG LOADING
# =============================================================================
def _load_config() -> Dict[str, Any]:
"""
Load memory_bank.config.json
Returns:
Config dict, or empty dict on error
"""
# Look for config relative to this handler's location
config_path = Path(__file__).resolve().parents[4] / "config" / "memory_bank.config.json"
if not config_path.exists():
return {}
try:
with open(config_path, 'r', encoding='utf-8') as f:
return json.load(f)
except Exception:
return {}
# =============================================================================
# LINE COUNTING
# =============================================================================
def _count_file_lines(file_path: Path) -> int:
"""
Count physical lines in memory file
Args:
file_path: Path to JSON file
Returns:
Number of physical lines in file
"""
try:
with open(file_path, 'r', encoding='utf-8') as f:
return len(f.readlines())
except Exception:
# Silent failure - handlers don't log (3-tier architecture)
return 0
def _get_max_lines(file_path: Path, branch_name: str | None = None) -> int:
"""
Get max_lines limit with priority: file metadata > branch config > default
Args:
file_path: Path to JSON file
branch_name: Optional branch name for config lookup
Returns:
Max lines limit (default 600)
"""
# 1. Try file-level metadata first (highest priority)
try:
with open(file_path, 'r', encoding='utf-8') as f:
data = json.load(f)
metadata = data.get('document_metadata', {})
limits = metadata.get('limits', {})
file_limit = limits.get('max_lines')
if file_limit is not None:
return file_limit
except Exception:
pass
# 2. Try branch-level config (if branch_name provided or can be extracted)
if branch_name is None:
# Extract from filename (e.g., SEED.local.json -> SEED)
parts = file_path.stem.split('.')
branch_name = parts[0] if parts else None
if branch_name:
config = _load_config()
branch_limits = config.get('rollover', {}).get('per_branch', {}).get(branch_name, {})
if 'max_lines' in branch_limits:
return branch_limits['max_lines']
# 3. Fall back to global default from config
config = _load_config()
default_limit = config.get('rollover', {}).get('defaults', {}).get('max_lines')
if default_limit is not None:
return default_limit
# 4. Final fallback to hardcoded 600
return 600
# =============================================================================
# ROLLOVER DETECTION
# =============================================================================
def _should_rollover(file_path: Path) -> tuple[bool, int, int]:
"""
Check if file should rollover
Args:
file_path: Path to memory JSON file
Returns:
Tuple of (should_rollover, current_lines, max_lines)
"""
current_lines = _count_file_lines(file_path)
max_lines = _get_max_lines(file_path)
should_trigger = current_lines >= max_lines
return (should_trigger, current_lines, max_lines)
def check_all_branches() -> Dict[str, Any]:
"""
Check all branches for rollover triggers
Scans AIPASS_REGISTRY.json and checks each branch's memory files
(observations and local) for rollover conditions.
Returns:
Dict with success status, triggers list, and count
"""
triggers = []
# Read registry
branches = _read_registry()
if not branches:
return {
'success': True,
'triggers': [],
'count': 0,
'message': 'No branches in registry'
}
# Check each branch
for branch in branches:
branch_name = branch.get('name', 'UNKNOWN')
# Check both memory types
for memory_type in ['observations', 'local']:
file_path = _get_memory_file_path(branch, memory_type)
if file_path is None:
continue # File doesn't exist, skip
should_trigger, current_lines, max_lines = _should_rollover(file_path)
if should_trigger:
trigger = RolloverTrigger(
branch=branch_name,
memory_type=memory_type,
file_path=file_path,
current_lines=current_lines,
max_lines=max_lines
)
triggers.append(trigger)
return {
'success': True,
'triggers': triggers,
'count': len(triggers),
'message': f'Found {len(triggers)} rollover triggers' if triggers else 'No rollover triggers detected'
}
def check_single_file(file_path: Path) -> Dict[str, Any]:
"""
Check single file for rollover trigger
Args:
file_path: Path to memory JSON file
Returns:
Dict with trigger status and details
"""
if not file_path.exists():
return {
'success': False,
'error': f"File not found: {file_path}"
}
should_trigger, current_lines, max_lines = _should_rollover(file_path)
if should_trigger:
# Extract branch and type from filename (e.g., SEED.observations.json)
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"
trigger = RolloverTrigger(
branch=branch_name,
memory_type=memory_type,
file_path=file_path,
current_lines=current_lines,
max_lines=max_lines
)
return {
'success': True,
'trigger': trigger,
'should_rollover': True
}
else:
return {
'success': True,
'should_rollover': False,
'current_lines': current_lines,
'max_lines': max_lines,
'remaining': max_lines - current_lines
}
# =============================================================================
# STATISTICS
# =============================================================================
def get_rollover_stats() -> Dict[str, Any]:
"""
Get rollover statistics for all branches
Returns:
Dict with statistics for all branches
"""
stats = {
'success': True,
'total_branches': 0,
'files_checked': 0,
'files_ready': 0,
'branches': {}
}
branches = _read_registry()
stats['total_branches'] = len(branches)
for branch in branches:
branch_name = branch.get('name', 'UNKNOWN')
branch_stats = {}
for memory_type in ['observations', 'local']:
file_path = _get_memory_file_path(branch, memory_type)
if file_path is None:
continue
stats['files_checked'] += 1
should_trigger, current_lines, max_lines = _should_rollover(file_path)
branch_stats[memory_type] = {
'current': current_lines,
'max': max_lines,
'ready': should_trigger,
'remaining': max_lines - current_lines
}
if should_trigger:
stats['files_ready'] += 1
if branch_stats:
stats['branches'][branch_name] = branch_stats
return stats
@@ -0,0 +1,647 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: memory_watcher.py - Memory File System Watcher
# Date: 2025-11-26
# 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
# =============================================
"""
Memory File System Watcher
Watches memory files (*.local.json, *.observations.json) for modifications.
On change, updates line counts and triggers rollover if needed.
Purpose:
Automatic memory file monitoring without polling. Responds to filesystem
events in real-time, keeping metadata accurate and triggering rollover
when thresholds are exceeded.
Independence:
Uses watchdog library for events. Delegates to line_counter and rollover
handlers for processing. No direct service dependencies.
"""
import logging
from pathlib import Path
from typing import Optional, Dict, Any
try:
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
WATCHDOG_AVAILABLE = True
except ImportError:
WATCHDOG_AVAILABLE = False
Observer = None
FileSystemEventHandler = object
# 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
logger = logging.getLogger(__name__)
# Memory root resolved relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[4]
# Global observer instance
_observer: Optional[Observer] = None
# =============================================================================
# STARTUP CHECK (runs once per command, no daemon needed)
# =============================================================================
_startup_check_done = False
def _get_rollover_threshold(branch_name: str, file_path: Path | None = None) -> int:
"""
Get rollover threshold for a memory file.
Priority: file metadata > per_branch config > defaults > hardcoded 600
Args:
branch_name: Branch name (uppercase, e.g., 'DEV_CENTRAL')
file_path: Optional path to memory file (checks file-level limits first)
Returns:
Max lines threshold for rollover
"""
import json
# 1. Check file-level metadata first (highest priority)
if file_path is not None:
try:
with open(file_path, 'r', encoding='utf-8') as f:
data = json.load(f)
file_limit = data.get('document_metadata', {}).get('limits', {}).get('max_lines')
if file_limit is not None:
return file_limit
except Exception:
pass
# 2. Check per-branch config override
config_path = _MEMORY_ROOT / "config" / "memory_bank.config.json"
try:
with open(config_path) as f:
config = json.load(f)
branch_limits = config.get('rollover', {}).get('per_branch', {}).get(branch_name, {})
if 'max_lines' in branch_limits:
return branch_limits['max_lines']
# 3. Fall back to defaults
default_limit = config.get('rollover', {}).get('defaults', {}).get('max_lines')
if default_limit is not None:
return default_limit
except Exception:
pass
# 4. Final fallback
return 600
def check_and_rollover() -> Dict[str, Any]:
"""
Check all memory files and trigger rollover if any exceed their threshold.
Also processes any new files in memory_pool.
Threshold is determined per-branch from config (defaults to 600).
This is a startup check - runs once per command, synchronous.
No daemon or file watcher needed.
Returns:
Dict with check results and any rollover actions taken
"""
global _startup_check_done
# Only run once per process
if _startup_check_done:
return {'success': True, 'skipped': True, 'reason': 'Already checked this session'}
_startup_check_done = True
results = {
'success': True,
'files_checked': 0,
'files_over_limit': [],
'rollover_triggered': False,
'memory_pool': None
}
# Get all branch paths
branch_paths = _get_branch_paths()
if not branch_paths:
results['error'] = 'No branch paths found'
return results
# Check each branch for memory files over limit
# Also sync current_lines metadata to keep it accurate
lines_synced = 0
for branch_path in branch_paths:
branch = Path(branch_path)
# 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
results['files_checked'] += 1
# Get threshold per file (file metadata > branch config > default)
threshold = _get_rollover_threshold(branch_name, memory_file)
try:
line_count = len(memory_file.read_text(encoding='utf-8').splitlines())
# Sync current_lines metadata if stale
try:
import json as _json
_data = _json.loads(memory_file.read_text(encoding='utf-8'))
meta_lines = _data.get('document_metadata', {}).get('status', {}).get('current_lines')
if meta_lines != line_count:
sync_result = update_line_count(memory_file)
if sync_result.get('success'):
lines_synced += 1
# Re-read actual line count after metadata update
line_count = len(memory_file.read_text(encoding='utf-8').splitlines())
except Exception:
pass # Non-critical - sync is best-effort
if line_count > threshold:
results['files_over_limit'].append({
'file': str(memory_file),
'lines': line_count,
'threshold': threshold
})
except Exception:
pass # Skip files we can't read
results['lines_synced'] = lines_synced
# Trigger rollover if any files are over limit
if results['files_over_limit']:
results['rollover_triggered'] = True
try:
from aipass.memory.apps.modules.rollover import execute_rollover
execute_rollover()
except ImportError:
logger.warning("Rollover module not available")
except Exception as e:
results['rollover_error'] = str(e)
results['success'] = False
# Check memory_pool for new files to process
results['memory_pool'] = _check_memory_pool()
# Check plans for new files to vectorize
results['plans'] = _check_plans()
# Check code_archive for new files to index
results['code_archive'] = _check_code_archive()
return results
def _check_memory_pool() -> Dict[str, Any]:
"""
Check memory_pool for new files and process if needed.
Compares files in pool against configured keep_recent limit.
If more files exist, runs processor to vectorize and archive.
Returns:
Dict with processing status
"""
import json
config_path = _MEMORY_ROOT / "config" / "memory_bank.config.json"
pool_path = _MEMORY_ROOT / "memory_pool"
# Load config
try:
with open(config_path) as f:
config = json.load(f)
pool_config = config.get('memory_pool', {})
except Exception:
return {'success': False, 'error': 'Could not load config'}
# Check if enabled
if not pool_config.get('enabled', False):
return {'success': True, 'skipped': True, 'reason': 'memory_pool disabled'}
# Count files in pool (excluding .archive)
extensions = pool_config.get('supported_extensions', ['.md', '.txt'])
keep_recent = pool_config.get('keep_recent', 10)
files = []
for ext in extensions:
files.extend(pool_path.glob(f'*{ext}'))
file_count = len(files)
# If under limit, nothing to do
if file_count <= keep_recent:
return {
'success': True,
'files_in_pool': file_count,
'keep_recent': keep_recent,
'action': 'none'
}
# Files exceed limit - run processor
try:
# NOTE: intake module not yet ported to aipass.memory package
from aipass.memory.apps.handlers.intake.pool_processor import process_memory_pool # type: ignore[import-not-found]
result = process_memory_pool()
return {
'success': result.get('success', False),
'files_processed': result.get('files_processed', 0),
'files_archived': result.get('archive', {}).get('archived_count', 0),
'action': 'processed'
}
except Exception as e:
return {'success': False, 'error': str(e), 'action': 'failed'}
def _check_plans() -> Dict[str, Any]:
"""
Check plans directory for files to vectorize.
Processes any plan files that haven't been vectorized yet.
Returns:
Dict with processing status
"""
import json
config_path = _MEMORY_ROOT / "config" / "memory_bank.config.json"
# Load config
try:
with open(config_path) as f:
config = json.load(f)
plans_config = config.get('plans', {})
except Exception:
return {'success': False, 'error': 'Could not load config'}
# Check if enabled
if not plans_config.get('enabled', False):
return {'success': True, 'skipped': True, 'reason': 'plans disabled'}
# Get plans path and count files (supports absolute paths)
plans_dir = plans_config.get('path', 'plans')
plans_path = Path(plans_dir) if Path(plans_dir).is_absolute() else _MEMORY_ROOT / plans_dir
extensions = plans_config.get('supported_extensions', ['.md'])
if not plans_path.exists():
return {'success': True, 'skipped': True, 'reason': 'plans directory does not exist'}
files = []
for ext in extensions:
files.extend(plans_path.glob(f'*{ext}'))
file_count = len(files)
if file_count == 0:
return {'success': True, 'files_in_plans': 0, 'action': 'none'}
# Process plans to vectors
try:
# NOTE: intake module not yet ported to aipass.memory package
from aipass.memory.apps.handlers.intake.plans_processor import process_plans # type: ignore[import-not-found]
result = process_plans()
return {
'success': result.get('success', False),
'files_processed': result.get('files_processed', 0),
'total_chunks': result.get('total_chunks', 0),
'action': 'processed'
}
except Exception as e:
return {'success': False, 'error': str(e), 'action': 'failed'}
def _check_code_archive() -> Dict[str, Any]:
"""
Check code_archive for new files and index them.
Returns:
Dict with indexing status
"""
try:
from aipass.memory.apps.handlers.archive.indexer import check_for_new_files
return check_for_new_files()
except Exception as e:
return {'success': False, 'error': str(e)}
# =============================================================================
# UTILITY FUNCTIONS
# =============================================================================
def _get_branch_paths() -> list[Path]:
"""
Get all branch paths from AIPASS_REGISTRY.json (silent - no logging)
Returns:
List of Path objects for each branch
"""
import json
registry_path = Path.home() / "AIPASS_REGISTRY.json"
if not registry_path.exists():
return []
try:
with open(registry_path, 'r', encoding='utf-8') as f:
data = json.load(f)
branches = data.get('branches', [])
paths = []
for branch in branches:
branch_path = Path(branch.get('path', ''))
if branch_path.exists():
paths.append(branch_path)
return paths
except Exception:
return []
def _is_memory_file(file_path: Path) -> bool:
"""
Check if file is a memory file (*.local.json or *.observations.json)
Args:
file_path: Path to check
Returns:
True if memory file, False otherwise
"""
name = file_path.name
return (name.endswith('.local.json') or name.endswith('.observations.json'))
# =============================================================================
# FILE SYSTEM EVENT HANDLER
# =============================================================================
class MemoryFileWatcher(FileSystemEventHandler):
"""Watch for memory file modifications"""
def __init__(self):
super().__init__()
# Track recent modifications to avoid duplicate processing
self._recent_modifications = set()
def on_modified(self, event):
"""Handle file modification events"""
# Ignore directory events
if event.is_directory:
return
file_path = Path(event.src_path)
# Only process memory files
if not _is_memory_file(file_path):
return
# Skip if we just processed this file
file_key = str(file_path)
if file_key in self._recent_modifications:
self._recent_modifications.discard(file_key)
return
logger.info(f"[memory_watcher] Detected modification: {file_path.name}")
# Step 1: Update line count metadata
update_result = update_line_count(file_path)
if not update_result['success']:
logger.error(f"[memory_watcher] Failed to update line count for {file_path.name}: {update_result.get('error')}")
return
current_lines = update_result.get('lines', 0)
logger.info(f"[memory_watcher] Updated {file_path.name}: {current_lines} lines")
# Step 2: Check if rollover needed
check_result = check_single_file(file_path)
if not check_result['success']:
logger.error(f"[memory_watcher] Failed to check rollover for {file_path.name}: {check_result.get('error')}")
return
if check_result.get('should_rollover', False):
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
# Trigger rollover
logger.info(f"[memory_watcher] Triggering rollover for {file_path.name}")
# Mark file as recently modified to avoid re-processing after rollover
self._recent_modifications.add(file_key)
# Execute rollover
try:
execute_rollover()
except Exception as e:
logger.error(f"[memory_watcher] Rollover failed: {e}")
# =============================================================================
# WATCHER CONTROL FUNCTIONS
# =============================================================================
def start_memory_watcher() -> Dict[str, Any]:
"""
Start watching memory files for modifications
Starts watchdog observer to monitor all branch directories for memory
file changes. Updates line counts and triggers rollover automatically.
Returns:
Dict with success status and watched paths
"""
global _observer
if _observer and _observer.is_alive():
return {
'success': False,
'error': 'Watcher already running'
}
# Get all branch paths
branch_paths = _get_branch_paths()
if not branch_paths:
return {
'success': False,
'error': 'No branch paths found in AIPASS_REGISTRY.json'
}
# Create watcher instance
watcher = MemoryFileWatcher()
# Create observer
new_observer = Observer()
# Schedule watcher for each branch path (silent - no logging during startup)
watched_paths = []
for branch_path in branch_paths:
try:
new_observer.schedule(watcher, str(branch_path), recursive=False)
watched_paths.append(str(branch_path))
except Exception:
pass # Skip invalid paths silently
# Start observer
new_observer.start()
_observer = new_observer
return {
'success': True,
'watched_paths': watched_paths,
'count': len(watched_paths)
}
def stop_memory_watcher() -> Dict[str, Any]:
"""
Stop the memory file watcher
Returns:
Dict with success status
"""
global _observer
if not _observer or not _observer.is_alive():
return {
'success': False,
'error': 'Watcher not running'
}
_observer.stop()
_observer.join()
_observer = None
logger.info("[memory_watcher] Stopped")
return {
'success': True,
'message': 'Memory watcher stopped'
}
def is_memory_watcher_active() -> bool:
"""
Check if memory watcher is currently active
Returns:
True if watcher is running, False otherwise
"""
return _observer is not None and _observer.is_alive()
def get_watcher_status() -> Dict[str, Any]:
"""
Get current watcher status
Returns:
Dict with watcher status and details
"""
active = is_memory_watcher_active()
if not active:
return {
'active': False,
'message': 'Watcher not running'
}
# Get watched paths
branch_paths = _get_branch_paths()
return {
'active': True,
'watched_directories': len(branch_paths),
'paths': [str(p) for p in branch_paths]
}
# =============================================================================
# STANDALONE EXECUTION
# =============================================================================
if __name__ == "__main__":
import argparse
import time
parser = argparse.ArgumentParser(
description='Memory File Watcher - Monitor memory files for rollover'
)
parser.add_argument(
'command',
choices=['start', 'stop', 'status'],
help='Command to execute'
)
args = parser.parse_args()
if args.command == 'start':
result = start_memory_watcher()
if result['success']:
print(f"Started watching {result['count']} directories")
print("Press Ctrl+C to stop...")
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
print("\nStopping...")
stop_memory_watcher()
else:
print(f"Failed to start: {result.get('error')}")
elif args.command == 'stop':
result = stop_memory_watcher()
if result['success']:
print(result['message'])
else:
print(f"Failed to stop: {result.get('error')}")
elif args.command == 'status':
status = get_watcher_status()
if status['active']:
print(f"Watcher is ACTIVE")
print(f"Watching {status['watched_directories']} directories:")
for path in status['paths']:
print(f" - {path}")
else:
print("Watcher is INACTIVE")
@@ -0,0 +1,2 @@
# Minimal __init__.py - handler independence pattern
# No exports - handlers import directly when needed
@@ -0,0 +1,435 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: extractor.py - Memory Extraction Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Memory Extraction Handler
Surgically extracts oldest items from memory files during rollover.
Understands real JSON structure (sessions, observations arrays).
Purpose:
When file exceeds 600 lines, extract oldest items from growing arrays
(sessions, observations, etc.), preserve JSON validity, update metadata.
Strategy:
- Detect which array is growing (sessions, observations, etc.)
- Calculate how many items to remove to get under limit
- Extract oldest items (FIFO)
- Update document_metadata.status
"""
import shutil
import logging
from pathlib import Path
from typing import Dict, List, Any, Tuple
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
logger = logging.getLogger(__name__)
# No module imports (handler independence)
# =============================================================================
# BACKUP OPERATIONS
# =============================================================================
def create_rollover_backup(file_path: Path) -> Dict[str, Any]:
"""
Create backup before rollover (safety net)
Creates backup in branch's .backup/ directory with timestamp.
Only keeps ONE backup - overwrites previous rollover backup.
Args:
file_path: Path to memory file to backup
Returns:
Dict with backup path and success status
"""
try:
# Create .backup directory in branch root
backup_dir = file_path.parent / '.backup'
backup_dir.mkdir(exist_ok=True)
# Backup filename: rollover_backup.json (always overwrites)
backup_name = f'rollover_backup_{file_path.name}'
backup_path = backup_dir / backup_name
# Copy file
shutil.copy2(file_path, backup_path)
return {
'success': True,
'backup_path': str(backup_path),
'message': f'Backup created: {backup_path.name}'
}
except Exception as e:
return {
'success': False,
'error': f'Backup failed: {e}'
}
def restore_from_backup(file_path: Path) -> Dict[str, Any]:
"""
Restore file from rollover backup
Called when rollover fails to restore original state.
Args:
file_path: Path to memory file to restore
Returns:
Dict with restore status
"""
try:
backup_dir = file_path.parent / '.backup'
backup_name = f'rollover_backup_{file_path.name}'
backup_path = backup_dir / backup_name
if not backup_path.exists():
return {
'success': False,
'error': 'No backup found to restore from'
}
# Restore from backup
shutil.copy2(backup_path, file_path)
return {
'success': True,
'message': f'Restored from backup: {backup_path.name}'
}
except Exception as e:
return {
'success': False,
'error': f'Restore failed: {e}'
}
# =============================================================================
# FILE OPERATIONS
# =============================================================================
def _read_memory_file(file_path: Path) -> Dict[str, Any]:
"""Read memory JSON file using json_handler"""
return read_memory_file_data(file_path)
def _write_memory_file(file_path: Path, data: Dict[str, Any]) -> None:
"""Write memory JSON file using json_handler"""
write_memory_file_simple(file_path, data)
def _count_file_lines(file_path: Path) -> int:
"""Count physical lines in file"""
with open(file_path, 'r', encoding='utf-8') as f:
return len(f.readlines())
# =============================================================================
# STRUCTURE DETECTION
# =============================================================================
def _detect_growing_array(data: Dict[str, Any]) -> str | None:
"""
Detect which array field is growing in memory file
Memory files have different structures:
- .local.json → 'sessions' array
- .observations.json → 'observations' array
- Future types → other array fields
Args:
data: Parsed JSON data
Returns:
Array field name (e.g., 'sessions'), or None if not found
"""
# Known array fields that grow over time
candidates = ['sessions', 'observations', 'recent_work', 'entries', 'items', 'records']
for field in candidates:
if field in data and isinstance(data[field], list) and len(data[field]) > 0:
return field
return None
# =============================================================================
# EXTRACTION CALCULATION
# =============================================================================
def _calculate_items_to_extract_by_lines(
data: Dict[str, Any],
array_field: str,
file_path: Path,
max_lines: int,
target_buffer: int = 100
) -> int:
"""
Calculate items to extract by SIMULATING line count (accurate)
Removes items one by one, counting actual lines after each removal,
until we reach target line count (max_lines - buffer).
Args:
data: Full memory file data
array_field: Name of array field to extract from
file_path: Path (for line counting)
max_lines: Maximum allowed lines
target_buffer: Lines of buffer to leave (default 100)
Returns:
Number of items to extract
Example:
File is 645 lines, limit is 600, buffer is 100
Target: 500 lines (600 - 100)
Simulate removing items until file is ~500 lines
"""
import tempfile
import json
target_lines = max_lines - target_buffer
total_items = len(data[array_field])
# Binary search for optimal item count
for items_to_remove in range(1, total_items + 1):
# Simulate removal (remove from END - oldest items)
test_data = data.copy()
test_data[array_field] = data[array_field][:-items_to_remove] # Keep newest
# Count lines in simulated result
with tempfile.NamedTemporaryFile(mode='w', delete=False, suffix='.json') as tmp:
json.dump(test_data, tmp, indent=2, ensure_ascii=False)
tmp_path = Path(tmp.name)
with open(tmp_path, 'r') as f:
line_count = len(f.readlines())
tmp_path.unlink()
# Check if we've reached target
if line_count <= target_lines:
return items_to_remove
# Fallback: remove 50% if simulation fails
return max(1, total_items // 2)
# =============================================================================
# EXTRACTION OPERATIONS
# =============================================================================
def extract_items(
file_path: Path,
percentage: int | None = None
) -> Dict[str, Any]:
"""
Extract items from memory file (WITH BACKUP SAFETY)
Now includes automatic backup before modification.
File is modified immediately, but backup exists for recovery.
Args:
file_path: Path to memory JSON file
percentage: Percentage of items to extract (auto-calculated if None)
Returns:
Dict with extracted items and metadata
"""
if not file_path.exists():
return {
'success': False,
'error': f"File not found: {file_path}"
}
# Read file
try:
data = _read_memory_file(file_path)
current_lines = _count_file_lines(file_path)
except Exception as e:
return {
'success': False,
'error': f"Failed to read file: {e}"
}
# Detect structure
array_field = _detect_growing_array(data)
if not array_field:
return {
'success': False,
'error': f"No growing array found in {file_path.name}"
}
# Get metadata
max_lines = data.get('document_metadata', {}).get('limits', {}).get('max_lines', 600)
# Check if under limit
if current_lines <= max_lines:
return {
'success': True,
'skipped': True,
'message': f"File under limit ({current_lines}/{max_lines} lines)"
}
# Calculate extraction amount (simulate actual line reduction)
total_items = len(data[array_field])
if percentage is None:
# Use line-based calculation for accuracy
items_to_extract = _calculate_items_to_extract_by_lines(
data, array_field, file_path, max_lines, target_buffer=100
)
else:
# Manual percentage override (for testing)
items_to_extract = max(1, int(total_items * percentage / 100))
# Extract oldest items (LAST N in array - newest first, oldest last)
extracted = data[array_field][-items_to_extract:] # Take from end (oldest)
remaining = data[array_field][:-items_to_extract] # Keep from start (newest)
# Update array
data[array_field] = remaining
# Update metadata
_update_metadata_after_extraction(data)
# Write back
try:
_write_memory_file(file_path, data)
new_line_count = _count_file_lines(file_path)
except Exception as e:
return {
'success': False,
'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"
return {
'success': True,
'file': str(file_path),
'branch': branch_name,
'type': memory_type,
'array_field': array_field,
'extracted': extracted,
'extracted_count': items_to_extract,
'remaining_count': len(remaining),
'old_lines': current_lines,
'new_lines': new_line_count
}
# =============================================================================
# METADATA OPERATIONS
# =============================================================================
def _update_metadata_after_extraction(data: Dict[str, Any]) -> None:
"""
Update document_metadata after extraction.
Updates:
- status.last_health_check
Args:
data: JSON data dict (modified in place)
"""
# Ensure metadata structure
if 'document_metadata' not in data:
data['document_metadata'] = {}
metadata = data['document_metadata']
# Update status
if 'status' not in metadata:
metadata['status'] = {}
metadata['status']['last_health_check'] = datetime.now().strftime("%Y-%m-%d")
# =============================================================================
# VECTORIZATION PREPARATION
# =============================================================================
def extract_with_metadata(
file_path: Path,
percentage: int | None = None
) -> Dict[str, Any]:
"""
Extract items with enriched metadata for vectorization
Same as extract_items but adds metadata needed for vector storage.
Args:
file_path: Path to memory JSON file
percentage: Percentage of items to extract (auto-calculated if None)
Returns:
Dict with extracted items + vectorization metadata
"""
# Do standard extraction
result = extract_items(file_path, percentage)
if not result['success']:
return result
# Enrich extracted items with metadata
extracted = result.get('extracted', [])
branch = result.get('branch')
memory_type = result.get('type')
array_field = result.get('array_field')
extraction_timestamp = datetime.now().isoformat()
enriched = []
for item in extracted:
enriched_item = {
**item, # Preserve original item data
'_metadata': {
'branch': branch,
'type': memory_type,
'array_field': array_field,
'extracted_at': extraction_timestamp,
'source_file': file_path.name
}
}
enriched.append(enriched_item)
# Return enriched version
return {
'success': True,
'file': str(file_path),
'branch': branch,
'type': memory_type,
'array_field': array_field,
'entries': enriched,
'count': len(enriched),
'old_lines': result.get('old_lines'),
'new_lines': result.get('new_lines')
}
@@ -0,0 +1,223 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: normalize.py - Memory File Schema Normalizer
# Date: 2026-01-22
# 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
# =============================================
"""
Memory File Schema Normalizer
Fixes inconsistent schema in memory JSON files:
1. Moves root-level 'limits' into document_metadata.limits
2. Removes redundant root-level 'status'
3. Removes auto_compress_at (redundant with max_lines)
4. Ensures document_metadata.status has current_lines
Target schema:
{
"document_metadata": {
"limits": { "max_lines": N, ... },
"status": { "current_lines": N, "last_health_check": "..." }
},
... (other content)
}
"""
import json
import logging
from pathlib import Path
from typing import Dict, Any
from datetime import datetime
logger = logging.getLogger(__name__)
def normalize_memory_file(file_path: Path, dry_run: bool = False) -> Dict[str, Any]:
"""
Normalize schema for a single memory file.
Args:
file_path: Path to memory JSON file
dry_run: If True, report changes without writing
Returns:
Dict with success status and changes made
"""
if not file_path.exists():
return {'success': False, 'error': f"File not found: {file_path}"}
try:
with open(file_path, 'r', encoding='utf-8') as f:
data = json.load(f)
except Exception as e:
return {'success': False, 'error': f"Failed to read: {e}"}
changes = []
# Ensure document_metadata exists
if 'document_metadata' not in data:
data['document_metadata'] = {}
changes.append("Created document_metadata")
metadata = data['document_metadata']
# 1. Move root 'limits' into document_metadata.limits
if 'limits' in data and 'limits' not in metadata:
metadata['limits'] = data.pop('limits')
changes.append("Moved root 'limits' into document_metadata")
elif 'limits' in data and 'limits' in metadata:
# Both exist - merge, preferring document_metadata values
root_limits = data.pop('limits')
for key, val in root_limits.items():
if key not in metadata['limits']:
metadata['limits'][key] = val
changes.append("Merged root 'limits' into document_metadata.limits")
# 2. Remove root 'status' (redundant)
if 'status' in data:
root_status = data.pop('status')
# If document_metadata.status doesn't have current_lines, copy it
if 'status' not in metadata:
metadata['status'] = {}
if 'current_lines' not in metadata['status'] and 'current_lines' in root_status:
metadata['status']['current_lines'] = root_status['current_lines']
changes.append("Removed redundant root 'status'")
# 3. Remove auto_compress_at from document_metadata.status (redundant with max_lines)
if 'status' in metadata and 'auto_compress_at' in metadata['status']:
del metadata['status']['auto_compress_at']
changes.append("Removed redundant 'auto_compress_at'")
# 4. Remove unused limits fields (max_word_count, max_token_count - no code uses these)
if 'limits' in metadata:
for unused_field in ['max_word_count', 'max_token_count']:
if unused_field in metadata['limits']:
del metadata['limits'][unused_field]
changes.append(f"Removed unused '{unused_field}'")
# 4. Ensure status has required fields
if 'status' not in metadata:
metadata['status'] = {}
if 'current_lines' not in metadata['status']:
# Count actual lines
try:
with open(file_path, 'r', encoding='utf-8') as f:
metadata['status']['current_lines'] = len(f.readlines())
changes.append("Added current_lines count")
except Exception:
pass
if 'last_health_check' not in metadata['status']:
metadata['status']['last_health_check'] = datetime.now().strftime("%Y-%m-%d")
changes.append("Added last_health_check")
# Write if changes made and not dry run
if changes and not dry_run:
try:
with open(file_path, 'w', encoding='utf-8') as f:
json.dump(data, f, indent=2, ensure_ascii=False)
f.write('\n')
except Exception as e:
return {'success': False, 'error': f"Failed to write: {e}"}
return {
'success': True,
'file': str(file_path),
'changes': changes,
'dry_run': dry_run
}
def normalize_all_memory_files(dry_run: bool = False) -> Dict[str, Any]:
"""
Normalize schema for all memory files in AIPASS_REGISTRY.
Args:
dry_run: If True, report changes without writing
Returns:
Dict with statistics
"""
# Read registry
registry_path = Path.home() / "AIPASS_REGISTRY.json"
if not registry_path.exists():
return {'success': False, 'error': "AIPASS_REGISTRY.json not found"}
try:
with open(registry_path, 'r', encoding='utf-8') as f:
registry = json.load(f)
branches = registry.get('branches', [])
except Exception as e:
return {'success': False, 'error': f"Failed to read registry: {e}"}
results = {
'success': True,
'files_checked': 0,
'files_modified': 0,
'dry_run': dry_run,
'details': []
}
for branch in branches:
branch_path = Path(branch.get('path', ''))
branch_name = branch.get('name', '').upper()
if not branch_path.exists():
continue
# Check both file types
for memory_type in ['local', 'observations']:
file_name = f"{branch_name}.{memory_type}.json"
file_path = branch_path / file_name
if not file_path.exists():
continue
results['files_checked'] += 1
result = normalize_memory_file(file_path, dry_run=dry_run)
if result['success'] and result.get('changes'):
results['files_modified'] += 1
results['details'].append({
'file': file_name,
'changes': result['changes']
})
return results
# CLI entry point
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="Normalize memory file schema")
parser.add_argument('--dry-run', action='store_true', help="Report changes without writing")
parser.add_argument('--file', type=str, help="Normalize single file")
args = parser.parse_args()
if args.file:
result = normalize_memory_file(Path(args.file), dry_run=args.dry_run)
print(json.dumps(result, indent=2))
else:
result = normalize_all_memory_files(dry_run=args.dry_run)
print(f"Files checked: {result['files_checked']}")
print(f"Files modified: {result['files_modified']}")
if result['details']:
print("\nChanges:")
for detail in result['details']:
print(f" {detail['file']}:")
for change in detail['changes']:
print(f" - {change}")
@@ -0,0 +1,2 @@
# Minimal __init__.py - handler independence pattern
# No exports - handlers import directly when needed
@@ -0,0 +1,509 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: vector_search.py - Vector Search Handler
# Date: 2025-11-27
# 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
# =============================================
"""
Vector Search Handler
Queries ChromaDB collections using semantic embeddings.
Uses same model as embedder.py (all-MiniLM-L6-v2) for query encoding.
Purpose:
Search memory vectors by semantic similarity.
Supports filtering by metadata and multi-collection search.
Design:
- Query encoding uses same model as storage (embedder.py)
- Singleton pattern for model efficiency
- Collection-level and database-level search
- Metadata filtering for targeted queries
Dependencies (optional):
- chromadb
- sentence-transformers
- torch
"""
import logging
from typing import List, Dict, Any
from pathlib import Path
logger = logging.getLogger(__name__)
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[4]
# Shared ChromaDB client (reuse from chroma handler)
from aipass.memory.apps.handlers.storage.chroma import get_client
# =============================================================================
# QUERY ENCODING SERVICE (Singleton)
# =============================================================================
class QueryEncoder:
"""
Query encoding service using same model as embedder.py
Ensures query embeddings are compatible with stored embeddings.
Uses all-MiniLM-L6-v2 model with same settings.
"""
def __init__(self, model_name: str = 'all-MiniLM-L6-v2'):
"""
Initialize query encoder
Args:
model_name: HuggingFace model identifier (must match embedder.py)
Raises:
ImportError: If sentence-transformers or torch are not installed
"""
# Late imports (heavy optional dependencies)
try:
import torch
from sentence_transformers import SentenceTransformer
except ImportError as e:
raise ImportError(
f"Search requires sentence-transformers and torch. "
f"Install with: pip install sentence-transformers torch. "
f"Original error: {e}"
)
self.model_name = model_name
self.model = SentenceTransformer(model_name)
# GPU optimization if available
self.use_gpu = torch.cuda.is_available()
if self.use_gpu:
self.model = self.model.to('cuda')
self.dimension = 384 # all-MiniLM-L6-v2 output dimension
def encode(self, query: str) -> List[float]:
"""
Encode query text to embedding
Args:
query: Query text string
Returns:
384-dimensional embedding as list
"""
import torch
# Encode with same settings as embedder.py
embedding = self.model.encode(
query,
convert_to_tensor=False, # Return numpy
normalize_embeddings=True, # Critical for L2 distance
show_progress_bar=False
)
# Cleanup GPU memory if used
if self.use_gpu:
torch.cuda.empty_cache()
return embedding.tolist()
# Global encoder instance (singleton pattern)
_query_encoder = None
def _get_encoder() -> QueryEncoder:
"""
Get or create query encoder singleton
Lazy initialization - model loaded on first use
"""
global _query_encoder
if _query_encoder is None:
_query_encoder = QueryEncoder()
return _query_encoder
# =============================================================================
# CHROMA SEARCH SERVICE
# =============================================================================
class SearchService:
"""
ChromaDB search service
Handles connection to ChromaDB and query execution.
Supports single collection and multi-collection search.
"""
def __init__(self, db_path: Path | None = None):
"""
Initialize search service
Args:
db_path: Path to ChromaDB database (default: memory/.chroma)
"""
if db_path is None:
db_path = _MEMORY_ROOT / ".chroma"
# Use shared singleton client (prevents write contention)
self.client = get_client(db_path)
self.db_path = db_path
def query_collection(
self,
collection_name: str,
query_embedding: List[float],
n_results: int = 5,
where: Dict[str, Any] | None = None
) -> Dict[str, Any]:
"""
Query a specific collection
Args:
collection_name: Name of collection to query
query_embedding: Query embedding vector
n_results: Number of results to return
where: Metadata filter (e.g., {"branch": "SEED"})
Returns:
Dict with query results
"""
try:
collection = self.client.get_collection(
collection_name,
embedding_function=None
)
except Exception as e:
return {
"collection": collection_name,
"exists": False,
"error": f"Collection not found: {e}"
}
# Query collection
results = collection.query(
query_embeddings=[query_embedding],
n_results=n_results,
where=where
)
return {
"collection": collection_name,
"exists": True,
"ids": results['ids'][0] if results['ids'] else [],
"documents": results['documents'][0] if results['documents'] else [],
"metadatas": results['metadatas'][0] if results['metadatas'] else [],
"distances": results['distances'][0] if results['distances'] else [],
"count": len(results['ids'][0]) if results['ids'] else 0
}
def list_collections(self) -> List[str]:
"""
List all collections in database
Returns:
List of collection names
"""
collections = self.client.list_collections()
return [col.name for col in collections]
# Global search service instance (singleton pattern)
_search_service = None
# Cache for local branch services
_local_services: Dict[str, SearchService] = {}
def _get_service(db_path: Path | None = None) -> SearchService:
"""
Get or create search service
Args:
db_path: Path to ChromaDB database (None = global default)
Returns:
SearchService instance for specified path
"""
global _search_service
# Global service (default)
if db_path is None:
if _search_service is None:
_search_service = SearchService()
return _search_service
# Local service (branch-specific)
path_str = str(db_path)
if path_str not in _local_services:
_local_services[path_str] = SearchService(db_path)
return _local_services[path_str]
# =============================================================================
# PUBLIC API
# =============================================================================
def search_collection(
query_embedding: List[float],
collection_name: str,
n_results: int = 5,
where: Dict[str, Any] | None = None,
db_path: Path | None = None
) -> Dict[str, Any]:
"""
Query a ChromaDB collection with embedding
Search a specific collection for semantically similar memories.
Args:
query_embedding: Pre-encoded query embedding (384-dim list)
collection_name: Name of collection to search
n_results: Number of results to return (default: 5)
where: Optional metadata filter (e.g., {"branch": "SEED"})
db_path: Path to ChromaDB database (None = global memory/.chroma)
Returns:
Dict with success status and search results
Example:
# Encode query first
query_result = encode_query("how does rollover work?")
# Search collection
result = search_collection(
query_embedding=query_result['embedding'],
collection_name="seed_observations",
n_results=10
)
if result['success']:
for i, doc in enumerate(result['documents']):
print(f"{i+1}. {doc[:100]}...")
"""
# Convert string db_path to Path
if db_path is not None and isinstance(db_path, str):
db_path = Path(db_path)
if not query_embedding:
return {
'success': False,
'error': 'No query embedding provided'
}
try:
service = _get_service(db_path)
result = service.query_collection(
collection_name=collection_name,
query_embedding=query_embedding,
n_results=n_results,
where=where
)
# Check if collection exists
if not result.get('exists', False):
return {
'success': False,
'error': result.get('error', 'Collection not found'),
'collection': collection_name
}
return {
'success': True,
**result
}
except Exception as e:
return {
'success': False,
'error': f"Search failed: {e}"
}
def encode_query(query: str) -> Dict[str, Any]:
"""
Encode query text to embedding using same model as storage
Uses all-MiniLM-L6-v2 model (same as embedder.py) to ensure
query embeddings are compatible with stored embeddings.
Args:
query: Query text string
Returns:
Dict with success status and embedding
Example:
result = encode_query("how does memory compression work?")
if result['success']:
embedding = result['embedding'] # 384-dim list
dimension = result['dimension'] # 384
"""
if not query or not query.strip():
return {
'success': False,
'error': 'Empty query string'
}
try:
encoder = _get_encoder()
embedding = encoder.encode(query)
return {
'success': True,
'embedding': embedding,
'dimension': len(embedding),
'model': encoder.model_name
}
except Exception as e:
return {
'success': False,
'error': f"Encoding failed: {e}"
}
def list_collections(db_path: Path | None = None) -> Dict[str, Any]:
"""
List available collections in database
Args:
db_path: Path to ChromaDB database (None = global default)
Returns:
Dict with success status and collection list
Example:
result = list_collections()
if result['success']:
for collection in result['collections']:
print(f"- {collection}")
"""
# Convert string db_path to Path
if db_path is not None and isinstance(db_path, str):
db_path = Path(db_path)
try:
service = _get_service(db_path)
collections = service.list_collections()
return {
'success': True,
'collections': collections,
'count': len(collections),
'db_path': str(service.db_path)
}
except Exception as e:
return {
'success': False,
'error': f"Failed to list collections: {e}"
}
def search_all_collections(
query_embedding: List[float],
n_results: int = 5,
where: Dict[str, Any] | None = None,
db_path: Path | None = None
) -> Dict[str, Any]:
"""
Search across all collections in database
Queries every collection and aggregates results.
Useful for cross-branch semantic search.
Args:
query_embedding: Pre-encoded query embedding
n_results: Number of results per collection
where: Optional metadata filter
db_path: Path to ChromaDB database (None = global default)
Returns:
Dict with success status and aggregated results
Example:
# Find similar memories across all branches
query_result = encode_query("deployment process")
result = search_all_collections(
query_embedding=query_result['embedding'],
n_results=3
)
if result['success']:
for coll_name, coll_results in result['results'].items():
print(f"\n{coll_name}:")
for doc in coll_results['documents']:
print(f" - {doc[:80]}...")
"""
# Convert string db_path to Path
if db_path is not None and isinstance(db_path, str):
db_path = Path(db_path)
if not query_embedding:
return {
'success': False,
'error': 'No query embedding provided'
}
try:
service = _get_service(db_path)
collections = service.list_collections()
if not collections:
return {
'success': True,
'results': {},
'message': 'No collections found'
}
# Query each collection
results = {}
for collection_name in collections:
coll_result = service.query_collection(
collection_name=collection_name,
query_embedding=query_embedding,
n_results=n_results,
where=where
)
if coll_result.get('exists', False):
results[collection_name] = {
'documents': coll_result['documents'],
'metadatas': coll_result['metadatas'],
'distances': coll_result['distances'],
'ids': coll_result['ids'],
'count': coll_result['count']
}
return {
'success': True,
'results': results,
'collections_searched': len(results),
'total_results': sum(r['count'] for r in results.values())
}
except Exception as e:
return {
'success': False,
'error': f"Multi-collection search failed: {e}"
}
@@ -0,0 +1,2 @@
# Minimal __init__.py - handler independence pattern
# No exports - handlers import directly when needed
@@ -0,0 +1,516 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: chroma.py - Chroma Vector Storage Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Chroma Vector Storage Handler
Manages Chroma vector database collections for Memory Bank.
Stores embeddings with metadata for semantic search.
Purpose:
Store 384-dim embeddings from rollover in branch-specific collections.
Enables fast local search per branch and global search across all branches.
Best Practices Applied:
- PersistentClient for local-first architecture
- Collection-per-branch strategy (fast local search)
- Batch inserts (100-150 optimal)
- Metadata for filtering (branch, type, date)
- Singleton pattern (client initialized once)
Dependencies (optional):
- chromadb
"""
import logging
from typing import List, Dict, Any
from pathlib import Path
from datetime import datetime
logger = logging.getLogger(__name__)
# Resolve paths relative to handler location
_MEMORY_ROOT = Path(__file__).resolve().parents[4]
# ChromaDB client - create inline since the old symbolic.chroma_client
# was an internal singleton wrapper
_chroma_clients: Dict[str, Any] = {}
def get_client(db_path: Path):
"""
Get or create a ChromaDB PersistentClient for the given path.
Singleton per path to prevent write contention.
Args:
db_path: Path to ChromaDB database directory
Returns:
chromadb.PersistentClient instance
Raises:
ImportError: If chromadb is not installed
"""
path_str = str(db_path)
if path_str not in _chroma_clients:
try:
import chromadb
except ImportError:
raise ImportError(
"chromadb is required for vector storage. "
"Install with: pip install chromadb"
)
db_path.mkdir(parents=True, exist_ok=True)
_chroma_clients[path_str] = chromadb.PersistentClient(path=str(db_path))
return _chroma_clients[path_str]
# =============================================================================
# CHROMA SERVICE (Singleton)
# =============================================================================
class ChromaService:
"""
Chroma vector database service
Implements best practices:
- PersistentClient for local-first
- Collection-per-branch strategy
- Batch inserts (optimal for 100-vector rollover)
- Metadata structure for fast filtering
"""
def __init__(self, db_path: Path | None = None):
"""
Initialize Chroma service
Args:
db_path: Path to Chroma database (default: memory/.chroma)
"""
if db_path is None:
db_path = _MEMORY_ROOT / ".chroma"
# Use shared singleton client (prevents write contention)
self.client = get_client(db_path)
self.db_path = db_path
def get_collection_name(self, branch: str, memory_type: str) -> str:
"""
Generate collection name using pattern: {branch}_{type}
Args:
branch: Branch name (SEED, CLI, etc.)
memory_type: Memory type (observations, local)
Returns:
Collection name (e.g., 'seed_observations')
"""
return f"{branch.lower()}_{memory_type.lower()}"
def store_vectors(
self,
branch: str,
memory_type: str,
embeddings: List,
documents: List[str],
metadatas: List[Dict[str, Any]]
) -> Dict[str, Any]:
"""
Store vectors in branch-specific collection
Args:
branch: Branch name
memory_type: Memory type
embeddings: List of embedding vectors (numpy arrays)
documents: List of text documents
metadatas: List of metadata dicts
Returns:
Dict with storage details
"""
collection_name = self.get_collection_name(branch, memory_type)
# Get or create collection
collection = self.client.get_or_create_collection(
name=collection_name,
metadata={"hnsw:space": "cosine", "branch": branch, "type": memory_type},
embedding_function=None
)
# Get existing count for ID generation
existing_count = collection.count()
# Generate unique IDs
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))
]
# Convert embeddings to list format (Chroma requirement)
embeddings_list = [emb.tolist() if hasattr(emb, 'tolist') else emb for emb in embeddings]
# Batch insert (optimal size: 100-150)
collection.add(
embeddings=embeddings_list,
documents=documents,
metadatas=metadatas,
ids=ids
)
new_count = collection.count()
return {
"collection": collection_name,
"count": len(embeddings),
"total_vectors": new_count,
"ids": ids
}
def get_collection_stats(self, branch: str, memory_type: str) -> Dict[str, Any]:
"""
Get statistics for collection
Args:
branch: Branch name
memory_type: Memory type
Returns:
Dict with collection statistics
"""
collection_name = self.get_collection_name(branch, memory_type)
try:
collection = self.client.get_collection(
collection_name,
embedding_function=None
)
count = collection.count()
return {
"collection": collection_name,
"exists": True,
"vector_count": count
}
except Exception:
return {
"collection": collection_name,
"exists": False,
"vector_count": 0
}
def list_all_collections(self) -> List[str]:
"""
List all collections in database
Returns:
List of collection names
"""
collections = self.client.list_collections()
return [col.name for col in collections]
# Global service instance (singleton pattern for default/global)
_chroma_service = None
# Cache for local branch services (one per branch)
_local_services: Dict[str, ChromaService] = {}
def _get_service(db_path: Path | None = None) -> ChromaService:
"""
Get or create Chroma service
Args:
db_path: Path to Chroma database (None = global default)
Returns:
ChromaService instance for specified path
Behavior:
- If db_path is None: Returns global singleton (memory/.chroma)
- If db_path specified: Returns cached service for that path or creates new one
"""
global _chroma_service
# Global service (default)
if db_path is None:
if _chroma_service is None:
_chroma_service = ChromaService()
return _chroma_service
# Local service (branch-specific)
path_str = str(db_path)
if path_str not in _local_services:
_local_services[path_str] = ChromaService(db_path)
return _local_services[path_str]
# =============================================================================
# PUBLIC API
# =============================================================================
def store_vectors(
branch: str,
memory_type: str,
embeddings: List,
documents: List[str],
metadatas: List[Dict[str, Any]],
db_path: Path | None = None
) -> Dict[str, Any]:
"""
Store vectors in Chroma collection
This is the main public API for storing vectors after rollover.
Creates or uses existing collection for branch + type combination.
Args:
branch: Branch name (SEED, CLI, etc.)
memory_type: Memory type (observations, local)
embeddings: List of embedding vectors
documents: List of original text documents
metadatas: List of metadata dicts for each entry
db_path: Path to Chroma database (None = global, or specify local branch path)
Returns:
Dict with storage details
Example:
# Store in global Memory Bank
result = store_vectors("SEED", "observations", embeddings, texts, metadatas)
# Store in SEED's local Chroma
local_path = Path("path/to/seed/.chroma")
result = store_vectors("SEED", "observations", embeddings, texts, metadatas, local_path)
"""
# Convert string db_path to Path (subprocess passes strings via JSON)
if db_path is not None and isinstance(db_path, str):
db_path = Path(db_path)
if not embeddings:
return {
'success': True,
'message': 'No vectors provided',
'count': 0
}
if len(embeddings) != len(documents) or len(embeddings) != len(metadatas):
return {
'success': False,
'error': f"Length mismatch: {len(embeddings)} embeddings, "
f"{len(documents)} documents, {len(metadatas)} metadatas"
}
try:
service = _get_service(db_path)
result = service.store_vectors(branch, memory_type, embeddings, documents, metadatas)
return {
'success': True,
**result
}
except Exception as e:
return {
'success': False,
'error': f"Storage failed: {e}"
}
def get_collection_stats(branch: str, memory_type: str) -> Dict[str, Any]:
"""
Get statistics for collection
Args:
branch: Branch name
memory_type: Memory type
Returns:
Dict with collection statistics
"""
try:
service = _get_service()
stats = service.get_collection_stats(branch, memory_type)
return {
'success': True,
**stats
}
except Exception as e:
return {
'success': False,
'error': f"Failed to get stats: {e}"
}
def list_all_collections() -> Dict[str, Any]:
"""
List all collections in database
Returns:
Dict with collection list
"""
try:
service = _get_service()
collections = service.list_all_collections()
return {
'success': True,
'collections': collections,
'count': len(collections)
}
except Exception as e:
return {
'success': False,
'error': f"Failed to list collections: {e}"
}
def get_database_info() -> Dict[str, Any]:
"""
Get database information
Returns:
Dict with database metadata
"""
try:
service = _get_service()
collections = service.list_all_collections()
return {
'success': True,
'db_path': str(service.db_path),
'collections_count': len(collections),
'collections': collections
}
except Exception as e:
return {
'success': False,
'error': f"Failed to get database info: {e}"
}
def search_vectors(
query_embedding: List[float],
branch: str | None = None,
memory_type: str | None = None,
n_results: int = 5,
db_path: Path | None = None
) -> Dict[str, Any]:
"""
Search for similar vectors in Chroma collections
Args:
query_embedding: Query vector (384-dim for all-MiniLM-L6-v2)
branch: Optional branch filter (if None, searches all collections)
memory_type: Optional memory type filter (observations, local)
n_results: Number of results to return per collection
db_path: Path to Chroma database (None = global)
Returns:
Dict with search results grouped by collection
Example:
# Search specific branch
results = search_vectors(query_emb, branch="SEED", memory_type="observations")
# Global search across all branches
results = search_vectors(query_emb, n_results=10)
"""
# Convert string db_path to Path (subprocess passes strings via JSON)
if db_path is not None and isinstance(db_path, str):
db_path = Path(db_path)
try:
service = _get_service(db_path)
# Determine which collections to search
if branch and memory_type:
# Search specific collection
collection_names = [service.get_collection_name(branch, memory_type)]
else:
# Search all collections
collection_names = service.list_all_collections()
# Apply filters
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 found'
}
# Search each collection
all_results = []
for collection_name in collection_names:
try:
collection = service.client.get_collection(
collection_name,
embedding_function=None
)
# Query collection
results = collection.query(
query_embeddings=[query_embedding],
n_results=n_results
)
# Format results
if results['documents'] and results['documents'][0]:
for i, doc in enumerate(results['documents'][0]):
all_results.append({
'collection': collection_name,
'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 as e:
# Collection might not exist - skip it
continue
# Sort by distance (lower is better)
all_results.sort(key=lambda x: x['distance'] if x['distance'] is not None else float('inf'))
return {
'success': True,
'results': all_results[:n_results * len(collection_names)] if all_results else [],
'collections_searched': len(collection_names),
'total_results': len(all_results)
}
except Exception as e:
return {
'success': False,
'error': f"Search failed: {e}"
}
@@ -0,0 +1,73 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: chroma_subprocess.py - ChromaDB Subprocess Handler
# Date: 2025-11-27
# 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
# =============================================
"""
ChromaDB Subprocess Handler
Called via subprocess from rollover module to ensure ChromaDB operations
run in an isolated process.
Input: JSON on stdin with operation and parameters
Output: JSON on stdout with result
"""
import sys
import json
from aipass.memory.apps.handlers.storage.chroma import store_vectors, list_all_collections, search_vectors
def main():
"""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(
branch=input_data.get('branch'),
memory_type=input_data.get('memory_type'),
embeddings=input_data.get('embeddings'),
documents=input_data.get('documents'),
metadatas=input_data.get('metadatas'),
db_path=input_data.get('db_path')
)
elif operation == 'list_collections':
result = list_all_collections()
elif operation == 'search_vectors':
result = search_vectors(
query_embedding=input_data.get('query_embedding'),
branch=input_data.get('branch'),
memory_type=input_data.get('memory_type'),
n_results=input_data.get('n_results', 5),
db_path=input_data.get('db_path')
)
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)
if __name__ == '__main__':
main()
@@ -0,0 +1 @@
"""Tracking handlers - Monitor and update memory file metadata"""
@@ -0,0 +1,151 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: line_counter.py - Memory File Line Counter Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Line Counter Handler
Updates document_metadata.status.current_lines field in memory files.
Called after any edit to keep metadata accurate.
Purpose:
Keep memory file metadata in sync with actual file state.
Provides accurate current_lines count for rollover detection.
Independence:
Uses json_handler for safe, atomic metadata updates
"""
import logging
from pathlib import Path
from typing import Dict, Any
from datetime import datetime
# Handler imports (relative within package)
from aipass.memory.apps.handlers.json.json_handler import update_metadata
logger = logging.getLogger(__name__)
# =============================================================================
# LINE COUNTING
# =============================================================================
def _count_physical_lines(file_path: Path) -> int:
"""
Count physical lines in file
Args:
file_path: Path to file
Returns:
Number of lines
"""
try:
with open(file_path, 'r', encoding='utf-8') as f:
return len(f.readlines())
except Exception:
return 0
# =============================================================================
# METADATA UPDATE
# =============================================================================
def update_line_count(file_path: Path) -> Dict[str, Any]:
"""
Update current_lines in document_metadata.status
Reads file, counts lines, updates metadata field using safe json_handler.
Args:
file_path: Path to memory JSON file
Returns:
Dict with success status and updated line count
"""
if not file_path.exists():
return {
'success': False,
'error': f"File not found: {file_path}"
}
# Count lines
line_count = _count_physical_lines(file_path)
# Update metadata using safe handler (atomic write)
result = update_metadata(
file_path,
current_lines=line_count,
last_health_check=datetime.now().strftime("%Y-%m-%d")
)
if not result['success']:
return {
'success': False,
'error': f"Failed to update metadata: {result['error']}"
}
return {
'success': True,
'file': str(file_path),
'lines': line_count
}
def update_all_memory_files() -> Dict[str, Any]:
"""
Update line counts for all memory files in AIPASS_REGISTRY
Returns:
Dict with update statistics
"""
from aipass.memory.apps.handlers.monitor.detector import _read_registry, _get_memory_file_path
branches = _read_registry()
if not branches:
return {
'success': True,
'updated': 0,
'failed': 0,
'message': 'No branches in registry'
}
updated = 0
failed = []
for branch in branches:
branch_name = branch.get('name', 'UNKNOWN')
for memory_type in ['observations', 'local']:
file_path = _get_memory_file_path(branch, memory_type)
if file_path is None:
continue
result = update_line_count(file_path)
if result['success']:
updated += 1
else:
failed.append((branch_name, memory_type, result.get('error')))
return {
'success': True,
'updated': updated,
'failed': len(failed),
'failures': failed
}
@@ -0,0 +1,2 @@
# Minimal __init__.py - handler independence pattern
# No exports - handlers import directly when needed
@@ -0,0 +1,295 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: embedder.py - Vector Embedding Handler
# Date: 2025-11-16
# 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
# =============================================
"""
Vector Embedding Handler
Generates semantic embeddings using sentence-transformers/all-MiniLM-L6-v2.
Implements production best practices from research.
Purpose:
Convert text memories into 384-dimensional vectors for semantic search.
Optimized for batch processing (100 lines during rollover).
Best Practices Applied:
- Pre-sort by length (30% padding reduction)
- Built-in normalization (L2 distance requirement)
- GPU memory cleanup (prevent VRAM leaks)
- Batch size optimization (64 GPU, 16 CPU)
- Singleton pattern (model loaded once)
Dependencies (optional):
- sentence-transformers
- torch
"""
import logging
from typing import List, Dict, Any
from pathlib import Path
logger = logging.getLogger(__name__)
# No service imports - handlers are pure workers (3-tier architecture)
# No module imports (handler independence)
# =============================================================================
# EMBEDDING SERVICE (Singleton)
# =============================================================================
class EmbeddingService:
"""
Production-ready embedding service
Implements best practices:
- Batch size optimization (64 GPU, 16 CPU)
- Pre-sorting by length (reduces padding waste 30%)
- Built-in normalization (critical for L2 distance)
- GPU memory cleanup (prevents VRAM leaks)
"""
def __init__(self, model_name: str = 'all-MiniLM-L6-v2'):
"""
Initialize embedding service
Args:
model_name: HuggingFace model identifier
Raises:
ImportError: If sentence-transformers or torch are not installed
"""
# Late imports (heavy optional dependencies)
try:
import torch
from sentence_transformers import SentenceTransformer
except ImportError as e:
raise ImportError(
f"Embedding requires sentence-transformers and torch. "
f"Install with: pip install sentence-transformers torch. "
f"Original error: {e}"
)
self.model_name = model_name
self.model = SentenceTransformer(model_name)
# GPU optimization if available
self.use_gpu = torch.cuda.is_available()
if self.use_gpu:
self.model = self.model.to('cuda')
self.batch_size = 64
else:
self.batch_size = 16
self.dimension = 384 # all-MiniLM-L6-v2 output dimension
def encode_batch(self, texts: List[str]) -> Dict[str, Any]:
"""
Encode batch of texts with all optimizations
Best practices applied:
1. Pre-sort by length (reduces padding waste)
2. Batch processing (optimal batch size)
3. Built-in normalization (L2 distance requirement)
4. GPU cleanup (prevent VRAM leaks)
Args:
texts: List of text strings to encode
Returns:
Dict with embeddings and metadata
"""
import torch
if not texts:
return {
"embeddings": [],
"count": 0,
"dimension": self.dimension
}
# Pre-sort by length (reduces padding waste by 30%)
sorted_pairs = sorted(enumerate(texts), key=lambda x: len(x[1]))
sorted_indices, sorted_texts = zip(*sorted_pairs)
# Encode with optimal settings
embeddings = self.model.encode(
sorted_texts,
batch_size=self.batch_size,
convert_to_tensor=False, # Return numpy for Chroma
normalize_embeddings=True, # Critical for L2 distance
show_progress_bar=False
)
# Restore original order
ordered_embeddings = [None] * len(texts)
for original_idx, sorted_idx in enumerate(sorted_indices):
ordered_embeddings[sorted_idx] = embeddings[original_idx]
# Cleanup GPU memory if used
if self.use_gpu:
torch.cuda.empty_cache()
return {
"embeddings": ordered_embeddings,
"count": len(ordered_embeddings),
"dimension": self.dimension
}
# Global service instance (singleton pattern)
_embedding_service = None
def _get_service() -> EmbeddingService:
"""
Get or create embedding service singleton
Lazy initialization - model loaded on first use
"""
global _embedding_service
if _embedding_service is None:
_embedding_service = EmbeddingService()
return _embedding_service
# =============================================================================
# PUBLIC API
# =============================================================================
def encode_batch(texts: List[str]) -> Dict[str, Any]:
"""
Encode batch of texts to embeddings
This is the main public API. Delegates to singleton service
to avoid reloading the model.
Args:
texts: List of text strings to encode
Returns:
Dict with embeddings and metadata
Example:
result = encode_batch(["memory 1", "memory 2"])
if result['success']:
embeddings = result['embeddings']
# Each embedding is 384-dim numpy array
"""
if not texts:
return {
'success': True,
'embeddings': [],
'count': 0,
'message': 'No texts provided'
}
try:
service = _get_service()
result = service.encode_batch(texts)
return {
'success': True,
**result
}
except Exception as e:
return {
'success': False,
'error': f"Encoding failed: {e}"
}
def encode_memories(memories: List[Dict[str, Any]]) -> Dict[str, Any]:
"""
Encode memory entries to embeddings
Extracts text from memory entries and encodes them.
Preserves original memory structure for metadata.
Args:
memories: List of memory entry dicts (from extraction)
Returns:
Dict with embeddings and original memories
Example:
memories = [{"content": "...", "timestamp": "..."}]
result = encode_memories(memories)
embeddings = result['embeddings']
original = result['memories']
"""
if not memories:
return {
'success': True,
'embeddings': [],
'memories': [],
'count': 0,
'message': 'No memories provided'
}
# Extract text content from memories
texts = []
for memory in memories:
# Try common fields for text content
text = (
memory.get('content') or
memory.get('text') or
memory.get('message') or
str(memory) # Fallback to string representation
)
texts.append(text)
# Encode texts
encode_result = encode_batch(texts)
if not encode_result['success']:
return encode_result
# Combine embeddings with original memories
return {
'success': True,
'embeddings': encode_result['embeddings'],
'memories': memories,
'count': len(memories),
'dimension': encode_result['dimension']
}
def get_model_info() -> Dict[str, Any]:
"""
Get embedding model information
Returns:
Dict with model metadata
"""
try:
service = _get_service()
return {
'success': True,
'model_name': service.model_name,
'dimension': service.dimension,
'batch_size': service.batch_size,
'gpu_enabled': service.use_gpu
}
except Exception as e:
return {
'success': False,
'error': f"Failed to get model info: {e}"
}
+368
View File
@@ -0,0 +1,368 @@
# ===================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
# =============================================
"""
Memory System - Central Memory Archive & Rollover System
PURPOSE: Manages memory rollover, archival, and retrieval across all AIPass branches.
ARCHITECTURE:
- Entry point auto-discovers modules
- Modules orchestrate workflow by calling handlers
- Handlers implement domain-specific business logic
"""
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
# =============================================================================
# INFRASTRUCTURE SETUP
# =============================================================================
logger = logging.getLogger(__name__)
console = Console()
# Package name for module discovery (relative imports within this package)
_PACKAGE_BASE = "aipass.memory.apps.modules"
# =============================================================================
# INTROSPECTION DISPLAY
# =============================================================================
def print_introspection():
"""Display discovered modules and status"""
console.print()
console.print("[bold cyan]Memory - Central Memory Archive System[/bold cyan]")
console.print()
console.print("[dim]Manages memory rollover and archival across AIPass branches[/dim]")
console.print()
# Discover modules
modules = discover_modules()
console.print(f"[yellow]Discovered Modules:[/yellow] {len(modules)}")
console.print()
for module in modules:
module_name = module.__name__.split('.')[-1]
console.print(f" [cyan]*[/cyan] {module_name}")
if not modules:
console.print(" [dim]No modules discovered yet[/dim]")
console.print()
console.print("[dim]Run 'python3 -m aipass.memory.apps.memory --help' for usage information[/dim]")
console.print()
# =============================================================================
# HELP SYSTEM
# =============================================================================
def print_help():
"""Display Rich-formatted help"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Central Memory Archive System[/bold cyan]\n[dim]Vector search, memory rollover, and fragmented memory for AIPass[/dim]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
# What is Memory section
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"
)
console.print(Panel(
what_content,
title="[bold cyan]What is Memory?[/bold cyan]",
border_style="dim",
box=box.ROUNDED
))
console.print()
console.print("[bold cyan]AVAILABLE COMMANDS:[/bold cyan]")
console.print()
# Commands table showing actual drone commands
table = Table(show_header=True, header_style="bold cyan", border_style="dim")
table.add_column("Command", style="green")
table.add_column("Description", style="dim")
# Core commands
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("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)
console.print()
console.print("-" * 70)
console.print()
console.print("[bold cyan]USAGE:[/bold cyan]")
console.print()
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()
console.print(" [yellow]Direct execution:[/yellow]")
console.print(" [dim]python3 -m aipass.memory.apps.memory search \"query\"[/dim]")
console.print(" [dim]python3 -m aipass.memory.apps.memory rollover[/dim]")
console.print()
console.print("-" * 70)
console.print()
console.print("[bold cyan]SEARCH OPTIONS:[/bold cyan]")
console.print()
console.print(" [cyan]--branch BRANCH[/cyan] Filter by branch (e.g., SEED, CLI)")
console.print(" [cyan]--type TYPE[/cyan] Filter by memory type (observations, local)")
console.print(" [cyan]--n N[/cyan] Number of results (default: 5)")
console.print()
console.print(" [yellow]Example:[/yellow]")
console.print(" [dim]drone @memory search \"registry bugs\" --branch SEED --n 10[/dim]")
console.print()
console.print("-" * 70)
console.print()
console.print("[bold cyan]WATCH MODE:[/bold cyan]")
console.print()
console.print(" [yellow]Start memory file watcher:[/yellow]")
console.print(" [dim]drone @memory watch[/dim]")
console.print(" [dim]Monitors all branches, auto-rolls when limit exceeded[/dim]")
console.print(" [dim]Press Ctrl+C to stop[/dim]")
console.print()
console.print("-" * 70)
console.print()
console.print("Commands: search, rollover, status, check, watch, sync-lines, push-templates, diff-templates, template-status, symbolic")
console.print()
# =============================================================================
# MODULE DISCOVERY
# =============================================================================
MODULES_DIR = Path(__file__).parent / "modules"
def discover_modules() -> List[Any]:
"""
Auto-discover modules in modules/ directory
Pattern: Any .py file in modules/ that implements handle_command()
gets automatically discovered and registered.
Returns:
List of module objects with handle_command() method
"""
modules = []
if not MODULES_DIR.exists():
return modules
for file_path in MODULES_DIR.glob("*.py"):
if file_path.name.startswith("_"): # Skip __init__.py, __pycache__, etc.
continue
module_name = f"{_PACKAGE_BASE}.{file_path.stem}"
try:
module = importlib.import_module(module_name)
# Duck typing: If it has handle_command(), it's a module
if hasattr(module, 'handle_command'):
modules.append(module)
logger.info(f"[memory] Discovered module: {module_name}")
except Exception as e:
logger.error(f"[memory] Failed to load module {module_name}: {e}")
return modules
def route_command(command: str, args: List[str], modules: List[Any]) -> bool:
"""
Route command to appropriate module
Pattern: Each module's handle_command() returns True if it handled the command
Args:
command: Command name (e.g., 'rollover', 'search', 'status')
args: Additional arguments
modules: List of discovered modules
Returns:
True if command was handled, False otherwise
"""
# Built-in commands handled by entry point
if command == 'watch':
start_watch()
return True
for module in modules:
try:
if module.handle_command(command, args):
return True
except Exception as e:
logger.error(f"[memory] Module {module.__name__} error: {e}")
return False
# =============================================================================
# WATCH MODE
# =============================================================================
def start_watch() -> None:
"""
Start memory watcher - monitors branch memory files for auto-rollover
Watches all branches from AIPASS_REGISTRY.json. When a memory file
exceeds 600 lines, automatically triggers rollover.
Press Ctrl+C to stop.
"""
from ..handlers.monitor.memory_watcher import (
start_memory_watcher,
stop_memory_watcher,
is_memory_watcher_active,
get_watcher_status
)
from .modules.rollover import get_rollover_stats
# Signal handler for graceful shutdown
def signal_handler(sig, frame):
"""Handle SIGINT for graceful watcher shutdown."""
console.print("\n")
console.print("[dim]Stopping watcher...[/dim]")
stop_memory_watcher()
console.print("[green]>[/green] Watcher stopped")
sys.exit(0)
signal.signal(signal.SIGINT, signal_handler)
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Watch Mode[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
# Start the watcher
result = start_memory_watcher()
if not result.get('success'):
console.print(f"[red]x Failed to start watcher: {result.get('error')}[/red]")
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]Press Ctrl+C to stop[/dim]")
console.print()
# Show initial status
stats = get_rollover_stats()
if stats.get('success'):
ready = stats.get('files_ready', 0)
total = stats.get('files_checked', 0)
status_marker = "[red]![/red]" if ready > 0 else "[green]OK[/green]"
console.print(f"{status_marker} Current: {total} files monitored, {ready} ready for rollover")
console.print()
# Keep running until Ctrl+C
console.print("[dim]Watcher active. Waiting for file changes...[/dim]")
while True:
time.sleep(1)
# =============================================================================
# MAIN ENTRY POINT
# =============================================================================
def main():
"""Main entry point - routes commands or shows help"""
# Parse arguments
args = sys.argv[1:]
# Show introspection when run without arguments
if len(args) == 0:
print_introspection()
return
# Version flag
if args[0] in ['--version', '-V']:
console.print("memory v1.0.0")
return
# Show help only for explicit help flags
if args[0] in ['--help', '-h', 'help']:
print_help()
return
# Command provided - try to route to modules
modules = discover_modules()
command = args[0]
remaining_args = args[1:] if len(args) > 1 else []
if route_command(command, remaining_args, modules):
return # Module handled it successfully
else:
console.print()
console.print(f"[red]Unknown command: {command}[/red]")
console.print()
console.print("Run [dim]python3 -m aipass.memory.apps.memory --help[/dim] for available commands")
console.print()
return
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
console.print("\n\nOperation cancelled by user")
sys.exit(0)
except Exception as e:
logger.error(f"[memory] Entry point error: {e}", exc_info=True)
console.print(f"\nError: {e}")
sys.exit(1)
+645
View File
@@ -0,0 +1,645 @@
# ===================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
# =============================================
"""
Rollover Orchestration Module
Coordinates the memory rollover workflow by calling handlers in sequence:
1. Detect rollover triggers (monitor/detector)
2. Extract oldest memories (rollover/extractor)
3. Generate embeddings (vector/embedder)
4. Store in Chroma (storage/chroma)
Purpose:
Thin orchestration layer - no business logic implementation.
All domain logic lives in handlers.
"""
import sys
import logging
import subprocess
import json
from pathlib import Path
from typing import List, Dict
from rich.console import Console
from rich.panel import Panel
from rich import box
# =============================================================================
# INFRASTRUCTURE SETUP
# =============================================================================
logger = logging.getLogger(__name__)
console = Console()
# Handler imports (relative within the memory package)
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)}
# =============================================================================
# COMMAND HANDLERS
# =============================================================================
def handle_command(command: str, args: List[str]) -> bool: # noqa: ARG001
"""
Handle rollover commands
Commands supported:
- rollover: Execute rollover for triggered branches
- status: Show rollover statistics
- check: Check which branches need rollover
- sync-lines: Update line count metadata for all branches
Args:
command: Command name
args: Additional arguments
Returns:
True if command handled, False otherwise
"""
if command in ('--help', '-h', 'help'):
print_help()
return True
if command == 'rollover':
execute_rollover()
return True
elif command == 'status':
show_status()
return True
elif command == 'check':
check_triggers()
return True
elif command == 'sync-lines':
sync_line_counts()
return True
return False
def print_help() -> None:
"""Display rollover module help"""
console.print()
console.print(Panel.fit(
"[bold cyan]Rollover Module - Memory Rollover Orchestration[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
console.print("[bold]USAGE:[/bold]")
console.print(" python3 -m aipass.memory.apps.modules.rollover <command>")
console.print()
console.print("[bold]COMMANDS:[/bold]")
console.print(" [cyan]rollover[/cyan] Execute rollover for files over 600 lines")
console.print(" [cyan]status[/cyan] Show rollover statistics for all branches")
console.print(" [cyan]check[/cyan] Check which files need rollover (dry run)")
console.print(" [cyan]sync-lines[/cyan] Update line count metadata for all branches")
console.print(" [cyan]help[/cyan] Show this help message")
console.print()
console.print("[bold]WORKFLOW:[/bold]")
console.print(" 1. Detect files over 600 lines")
console.print(" 2. Extract oldest entries (target ~500 lines)")
console.print(" 3. Generate embeddings via sentence-transformers")
console.print(" 4. Store vectors in local + global ChromaDB")
console.print()
# =============================================================================
# ROLLOVER ORCHESTRATION
# =============================================================================
def execute_rollover() -> bool:
"""
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
"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Rollover Execution[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
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")
return False
triggers = triggers_result.get('triggers', [])
if not triggers:
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()
# 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"
console.print(
f" [green]>[/green] Rolled over {len(memories)} items -> {global_collection} "
f"({old_lines} -> {new_lines} lines, global: {global_total} vectors, {local_status})"
)
logger.info(f"[rollover] Successfully rolled over {trigger}: {len(memories)} items, {old_lines} -> {new_lines} lines")
# Report results
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")
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")
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
# =============================================================================
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.
"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Sync Line Counts[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
console.print("[cyan]Updating line counts for all memory files...[/cyan]")
console.print()
result = line_counter.update_all_memory_files()
if result['success']:
console.print(f"[green]>[/green] Updated {result['updated']} files")
if result['failed'] > 0:
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()
# =============================================================================
# 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
Displays:
- Files checked
- Files ready for rollover
- Per-branch status (current/max lines)
"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Rollover Status[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
# Get stats from detector
stats_result = detector.get_rollover_stats()
if not stats_result['success']:
console.print(f"[red]x[/red] Failed to get status: {stats_result.get('error', 'Unknown error')}")
logger.error(f"[rollover] Failed to get status: {stats_result.get('error')}")
return
stats = stats_result
# Summary
console.print(f"[cyan]Branches:[/cyan] {stats['total_branches']}")
console.print(f"[cyan]Files checked:[/cyan] {stats['files_checked']}")
console.print(f"[cyan]Ready for rollover:[/cyan] {stats['files_ready']}")
console.print()
# Per-branch details
if stats['branches']:
console.print("[yellow]Branch Details:[/yellow]")
console.print()
for branch_name, branch_stats in stats['branches'].items():
console.print(f" [bold]{branch_name}[/bold]")
for memory_type, file_stats in branch_stats.items():
current = file_stats['current']
max_lines = file_stats['max']
ready = file_stats['ready']
remaining = file_stats['remaining']
status_marker = "[red]![/red]" if ready else "[green]OK[/green]"
status_text = "READY" if ready else f"{remaining} remaining"
console.print(
f" {status_marker} {memory_type}: {current}/{max_lines} lines ({status_text})"
)
console.print()
def check_triggers() -> None:
"""
Check which branches need rollover (without executing)
Displays list of files that hit rollover threshold
"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Rollover Check[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
triggers_result = detector.check_all_branches()
if not triggers_result['success']:
console.print(f"[red]x[/red] Failed to check triggers: {triggers_result.get('error', 'Unknown error')}")
logger.error(f"[rollover] Failed to check triggers: {triggers_result.get('error')}")
return
triggers = triggers_result.get('triggers', [])
if not triggers:
console.print("[green]>[/green] No files need rollover")
return
console.print(f"[yellow]Found {len(triggers)} files ready for rollover:[/yellow]")
console.print()
for trigger in triggers:
console.print(f" * {trigger}")
console.print()
console.print("[dim]Run 'drone @memory rollover' to process these files[/dim]")
console.print()
# =============================================================================
# STANDALONE EXECUTION
# =============================================================================
if __name__ == "__main__":
import sys
# Handle --help before argparse (module standard)
if len(sys.argv) < 2 or sys.argv[1] in ('--help', '-h', 'help'):
handle_command('help', [])
sys.exit(0)
# Execute command via handle_command
command = sys.argv[1]
if not handle_command(command, sys.argv[2:]):
console.print(f"[red]Unknown command:[/red] {command}")
console.print("Run with [cyan]help[/cyan] for available commands")
sys.exit(1)
+383
View File
@@ -0,0 +1,383 @@
# ===================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
# =============================================
"""
Search Orchestration Module
Coordinates semantic search workflow by calling handlers in sequence:
1. Encode query text to embedding (vector/embedder)
2. Search Chroma collections (storage/chroma via subprocess)
3. Format and display results (Rich panels)
Purpose:
Thin orchestration layer - no business logic implementation.
All domain logic lives in handlers.
"""
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
# =============================================================================
# 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)}
# =============================================================================
# COMMAND HANDLERS
# =============================================================================
def handle_command(command: str, args: List[str]) -> bool:
"""
Handle search commands
Commands supported:
- search <query>: Execute semantic search across all branches
- help: Show search help
Args:
command: Command name
args: Additional arguments (query text, options)
Returns:
True if command handled, False otherwise
"""
if command in ('--help', '-h', 'help'):
print_help()
return True
if command == 'search':
if not args:
console.print("[red]Error:[/red] Search query required")
console.print("Usage: search <query> [--branch BRANCH] [--type TYPE] [--n N]")
return True
# Parse arguments
query_parts = []
branch = None
memory_type = None
n_results = 5
i = 0
while i < len(args):
if args[i] == '--branch' and i + 1 < len(args):
branch = args[i + 1]
i += 2
elif args[i] == '--type' and i + 1 < len(args):
memory_type = args[i + 1]
i += 2
elif args[i] == '--n' and i + 1 < len(args):
try:
n_results = int(args[i + 1])
except ValueError:
console.print(f"[red]Error:[/red] Invalid number: {args[i + 1]}")
return True
i += 2
else:
query_parts.append(args[i])
i += 1
query = ' '.join(query_parts)
if not query:
console.print("[red]Error:[/red] Search query required")
return True
execute_search(query, branch=branch, memory_type=memory_type, n_results=n_results)
return True
return False
def print_help() -> None:
"""Display search module help"""
console.print()
console.print(Panel.fit(
"[bold cyan]Search Module - Semantic Memory Search[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
console.print("[bold]USAGE:[/bold]")
console.print(" python3 -m aipass.memory.apps.modules.search search <query> [options]")
console.print()
console.print("[bold]COMMANDS:[/bold]")
console.print(" [cyan]search <query>[/cyan] Search across all memory collections")
console.print(" [cyan]help[/cyan] Show this help message")
console.print()
console.print("[bold]OPTIONS:[/bold]")
console.print(" [cyan]--branch BRANCH[/cyan] Filter by branch (e.g., SEED, CLI)")
console.print(" [cyan]--type TYPE[/cyan] Filter by memory type (observations, local)")
console.print(" [cyan]--n N[/cyan] Number of results (default: 5)")
console.print()
console.print("[bold]EXAMPLES:[/bold]")
console.print(" # Search all branches")
console.print(" [dim]drone @memory search \"error handling patterns\"[/dim]")
console.print()
console.print(" # Search specific branch")
console.print(" [dim]drone @memory search \"registry bugs\" --branch SEED[/dim]")
console.print()
console.print(" # Search specific memory type")
console.print(" [dim]drone @memory search \"collaboration\" --type observations --n 10[/dim]")
console.print()
console.print("[bold]HOW IT WORKS:[/bold]")
console.print(" 1. Convert query to 384-dim embedding (all-MiniLM-L6-v2)")
console.print(" 2. Search ChromaDB collections for similar vectors")
console.print(" 3. Display top N most relevant memories")
console.print()
# =============================================================================
# SEARCH ORCHESTRATION
# =============================================================================
def execute_search(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
Args:
query: Search query text
branch: Optional branch filter
memory_type: Optional memory type filter
n_results: Number of results to return
Returns:
True if search successful, False otherwise
"""
console.print()
console.print(Panel.fit(
"[bold cyan]Memory - Semantic Search[/bold cyan]",
border_style="cyan",
box=box.ROUNDED
))
console.print()
# Step 1: Encode query
console.print(f"[cyan]Query:[/cyan] {query}")
if branch:
console.print(f"[cyan]Branch:[/cyan] {branch}")
if memory_type:
console.print(f"[cyan]Type:[/cyan] {memory_type}")
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,
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}")
return False
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: Display results
console.print(f"[green]>[/green] Found {total_results} results in {collections_searched} collections")
console.print()
if not results:
console.print("[yellow]No matching memories found[/yellow]")
console.print()
console.print("[dim]Try:[/dim]")
console.print(" * Different search terms")
console.print(" * Broader query without filters")
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()
console.print("[dim]The search found some results but none were relevant enough (>40% similarity).[/dim]")
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)
# Parse collection name
parts = collection.split('_')
branch_name = parts[0].upper() if parts else 'UNKNOWN'
mem_type = parts[1] if len(parts) > 1 else 'unknown'
# Build metadata display
meta_lines = []
if 'timestamp' in metadata:
meta_lines.append(f"[dim]Time:[/dim] {metadata['timestamp']}")
if 'source' in metadata:
meta_lines.append(f"[dim]Source:[/dim] {metadata['source']}")
meta_text = " | ".join(meta_lines) if meta_lines else ""
# Create panel for each result
panel_title = f"Result {i} - {branch_name} ({mem_type}) - Similarity: {similarity:.2%}"
panel_content = document
if meta_text:
panel_content += f"\n\n{meta_text}"
console.print(Panel(
panel_content,
title=panel_title,
title_align="left",
border_style="cyan" if similarity > 0.7 else "blue" if similarity > 0.5 else "dim"
))
console.print()
logger.info(f"[search] Displayed {len(filtered_results)} results")
return True
# =============================================================================
# STANDALONE EXECUTION
# =============================================================================
if __name__ == "__main__":
# Handle --help before argparse (module standard)
if len(sys.argv) < 2 or sys.argv[1] in ('--help', '-h', 'help'):
handle_command('help', [])
sys.exit(0)
# Execute command via handle_command
command = sys.argv[1]
if not handle_command(command, sys.argv[2:]):
console.print(f"[red]Unknown command:[/red] {command}")
console.print("Run with [cyan]help[/cyan] for available commands")
sys.exit(1)
View File