feat(memory): auto-process memory pool + rollover on session-start/pre-compact (TDPLAN-0005)
This commit is contained in:
@@ -0,0 +1,160 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: auto_process.py
|
||||
# Description: Automated pool + rollover entry point
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-06-06
|
||||
# Modified: 2026-06-06
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Auto-process handler — session-start pool + rollover entry point.
|
||||
|
||||
Single callable the hook engine fires each session to:
|
||||
1. Process any files dropped into memory_pool/ (vectorize + archive)
|
||||
2. Check/run rollover for .trinity/ files exceeding limits
|
||||
|
||||
Idempotent: safe to call every session. Fast no-op when nothing to do.
|
||||
Pool uses upsert with content-hash IDs — re-processing same files is a no-op.
|
||||
|
||||
HOOK ENGINE CONTRACT:
|
||||
Module: aipass.memory.apps.handlers.intake.auto_process
|
||||
Function: auto_process()
|
||||
Invocation: importlib.import_module('aipass.memory.apps.handlers.intake.auto_process').auto_process()
|
||||
Returns: dict with success, pool, and rollover results
|
||||
"""
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict
|
||||
|
||||
from aipass.prax import logger
|
||||
from aipass.memory.apps.handlers.json import json_handler
|
||||
|
||||
_MEMORY_ROOT = Path(__file__).resolve().parent.parent.parent.parent
|
||||
CONFIG_PATH = _MEMORY_ROOT / "config" / "memory.config.json"
|
||||
|
||||
|
||||
def _load_pool_enabled() -> bool:
|
||||
try:
|
||||
with open(CONFIG_PATH, encoding="utf-8") as f:
|
||||
config = json.load(f)
|
||||
return config.get("memory_pool", {}).get("enabled", False)
|
||||
except Exception as e:
|
||||
logger.warning(f"[auto_process] Failed to load config: {e}")
|
||||
return False
|
||||
|
||||
|
||||
def run_pool_processing() -> Dict[str, Any]:
|
||||
"""
|
||||
Process memory pool files if enabled.
|
||||
|
||||
Checks config, calls process_memory_pool(), returns summary.
|
||||
Fast no-op when pool is empty or disabled.
|
||||
|
||||
Returns:
|
||||
dict with success/skipped, files_processed, total_chunks
|
||||
"""
|
||||
if not _load_pool_enabled():
|
||||
return {"skipped": True, "reason": "memory_pool disabled in config"}
|
||||
|
||||
try:
|
||||
from aipass.memory.apps.handlers.intake.pool_processor import process_memory_pool
|
||||
|
||||
pool_result = process_memory_pool()
|
||||
result = {
|
||||
"success": pool_result.get("success", False),
|
||||
"files_processed": pool_result.get("files_processed", 0),
|
||||
"total_chunks": pool_result.get("total_chunks", 0),
|
||||
}
|
||||
if pool_result.get("files_processed", 0) > 0:
|
||||
logger.info(
|
||||
f"[auto_process] Pool: {pool_result['files_processed']} files, "
|
||||
f"{pool_result.get('total_chunks', 0)} chunks"
|
||||
)
|
||||
|
||||
json_handler.log_operation(
|
||||
"run_pool_processing",
|
||||
{
|
||||
"files_processed": result.get("files_processed", 0),
|
||||
"success": result.get("success", False),
|
||||
},
|
||||
)
|
||||
|
||||
return result
|
||||
except Exception as e:
|
||||
logger.warning(f"[auto_process] Pool processing failed: {e}")
|
||||
return {"success": False, "error": str(e)}
|
||||
|
||||
|
||||
def _run_rollover_check() -> Dict[str, Any]:
|
||||
"""
|
||||
Check all branches for rollover triggers and execute if needed.
|
||||
|
||||
Returns:
|
||||
dict with success/skipped and rollover details
|
||||
"""
|
||||
try:
|
||||
from aipass.memory.apps.handlers.monitor.detector import check_all_branches
|
||||
|
||||
check_result = check_all_branches()
|
||||
triggers = check_result.get("triggers", []) if check_result else []
|
||||
|
||||
if not triggers:
|
||||
return {"skipped": True, "reason": "no rollover triggers"}
|
||||
|
||||
from aipass.memory.apps.handlers.rollover.orchestrator import execute_rollover
|
||||
|
||||
rollover_result = execute_rollover()
|
||||
result = {
|
||||
"success": rollover_result.get("success", False),
|
||||
"triggers": rollover_result.get("triggers_count", 0),
|
||||
"processed": rollover_result.get("success_count", 0),
|
||||
}
|
||||
logger.info(f"[auto_process] Rollover: {result['processed']}/{result['triggers']} triggers processed")
|
||||
return result
|
||||
except Exception as e:
|
||||
logger.warning(f"[auto_process] Rollover check failed: {e}")
|
||||
return {"success": False, "error": str(e)}
|
||||
|
||||
|
||||
def auto_process() -> Dict[str, Any]:
|
||||
"""
|
||||
Single idempotent entry point for session-start auto-processing.
|
||||
|
||||
Processes memory pool files and checks/runs rollover if needed.
|
||||
Fast no-op when pool is empty and no rollover triggers.
|
||||
Safe to call every session.
|
||||
|
||||
Returns:
|
||||
dict with success, pool, and rollover results
|
||||
"""
|
||||
result: Dict[str, Any] = {"success": True, "pool": None, "rollover": None}
|
||||
|
||||
if not _load_pool_enabled():
|
||||
result["pool"] = {"skipped": True, "reason": "memory_pool disabled in config"}
|
||||
result["rollover"] = {"skipped": True}
|
||||
logger.info("[auto_process] Skipped — memory_pool disabled in config")
|
||||
return result
|
||||
|
||||
# 1. Process pool files
|
||||
pool_result = run_pool_processing()
|
||||
result["pool"] = pool_result
|
||||
if pool_result.get("success") is False:
|
||||
result["success"] = False
|
||||
|
||||
# 2. Check/run rollover
|
||||
rollover_result = _run_rollover_check()
|
||||
result["rollover"] = rollover_result
|
||||
if rollover_result.get("success") is False:
|
||||
result["success"] = False
|
||||
|
||||
json_handler.log_operation(
|
||||
"auto_process",
|
||||
{
|
||||
"pool_files": result.get("pool", {}).get("files_processed", 0),
|
||||
"rollover_triggered": not result.get("rollover", {}).get("skipped", False),
|
||||
"success": result["success"],
|
||||
},
|
||||
)
|
||||
|
||||
return result
|
||||
@@ -117,6 +117,8 @@ def print_help():
|
||||
table.add_row("search <query>", "Semantic search across all branch memories")
|
||||
table.add_row("symbolic <subcommand>", "Symbolic/fragmented memory extraction and search")
|
||||
table.add_row("templates <subcommand>", "Living template push, diff, and status")
|
||||
table.add_row("pool process", "Process pool files + check/run rollover")
|
||||
table.add_row("pool status", "Show pool file count, config, vector stats")
|
||||
table.add_row("verify <plan_label>", "Check if a plan is vectorized in ChromaDB")
|
||||
table.add_row("watch", "Start memory watcher (auto-rollover on changes)")
|
||||
|
||||
@@ -163,7 +165,9 @@ def print_help():
|
||||
console.print("-" * 70)
|
||||
console.print()
|
||||
|
||||
console.print("Commands: search, rollover [run|status|check|sync-lines], symbolic, templates, verify, watch")
|
||||
console.print(
|
||||
"Commands: search, rollover [run|status|check|sync-lines], pool [process|status], symbolic, templates, verify, watch"
|
||||
)
|
||||
console.print()
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,192 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: pool.py
|
||||
# Description: Pool Module — drone CLI for pool commands
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-06-06
|
||||
# Modified: 2026-06-06
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Pool Module — drone CLI routing for memory pool commands.
|
||||
|
||||
Thin delegation layer. All implementation lives in handlers/intake/auto_process.py.
|
||||
"""
|
||||
|
||||
from typing import List, Any
|
||||
|
||||
from rich.panel import Panel
|
||||
from rich import box
|
||||
|
||||
from aipass.prax import logger # noqa: F401
|
||||
from aipass.cli.apps.modules import console, error
|
||||
from aipass.memory.apps.handlers.json import json_handler
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# COMMAND HANDLERS
|
||||
# =============================================================================
|
||||
|
||||
_SUBCOMMANDS = {
|
||||
"process": "Process memory pool files (vectorize + archive)",
|
||||
"status": "Show memory pool status",
|
||||
}
|
||||
|
||||
|
||||
def handle_command(command: str, args: List[Any]) -> bool:
|
||||
"""
|
||||
Handle pool commands.
|
||||
|
||||
Routing:
|
||||
pool (no args) -> print_introspection()
|
||||
pool --help/-h/help -> print_help()
|
||||
pool process -> run auto_process()
|
||||
pool status -> show pool status
|
||||
|
||||
Args:
|
||||
command: Command name
|
||||
args: Additional arguments
|
||||
|
||||
Returns:
|
||||
True if command handled, False otherwise
|
||||
"""
|
||||
if command == "pool":
|
||||
if not args:
|
||||
print_introspection()
|
||||
return True
|
||||
|
||||
if args[0] in ("--help", "-h", "help"):
|
||||
print_help()
|
||||
return True
|
||||
|
||||
sub = args[0]
|
||||
|
||||
if sub == "process":
|
||||
_run_process_command()
|
||||
return True
|
||||
|
||||
if sub == "status":
|
||||
_run_status_command()
|
||||
return True
|
||||
|
||||
error(
|
||||
f"Unknown subcommand: '{sub}'",
|
||||
suggestion="Available: " + ", ".join(_SUBCOMMANDS.keys()),
|
||||
)
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# CLI DISPLAY
|
||||
# =============================================================================
|
||||
|
||||
|
||||
def _run_process_command() -> None:
|
||||
"""Execute pool processing + rollover check and display results."""
|
||||
from ..handlers.intake.auto_process import auto_process
|
||||
|
||||
console.print()
|
||||
console.print("[bold cyan]Processing memory pool...[/bold cyan]")
|
||||
console.print()
|
||||
|
||||
result = auto_process()
|
||||
|
||||
json_handler.log_operation(
|
||||
"pool_process_command",
|
||||
{"success": result.get("success", False)},
|
||||
)
|
||||
|
||||
# Pool results
|
||||
pool = result.get("pool", {})
|
||||
if pool.get("skipped"):
|
||||
console.print(f"[dim]Pool: skipped — {pool.get('reason', 'unknown')}[/dim]")
|
||||
elif pool.get("success") is False:
|
||||
console.print(f"[red]Pool: failed — {pool.get('error', 'unknown')}[/red]")
|
||||
else:
|
||||
files = pool.get("files_processed", 0)
|
||||
chunks = pool.get("total_chunks", 0)
|
||||
if files > 0:
|
||||
console.print(f"[green]>[/green] Pool: {files} files processed, {chunks} chunks vectorized")
|
||||
else:
|
||||
console.print("[dim]Pool: no files to process[/dim]")
|
||||
|
||||
# Rollover results
|
||||
rollover = result.get("rollover", {})
|
||||
if rollover.get("skipped"):
|
||||
console.print("[dim]Rollover: no triggers[/dim]")
|
||||
elif rollover.get("success") is False:
|
||||
console.print(f"[red]Rollover: failed — {rollover.get('error', 'unknown')}[/red]")
|
||||
else:
|
||||
processed = rollover.get("processed", 0)
|
||||
total = rollover.get("triggers", 0)
|
||||
console.print(f"[green]>[/green] Rollover: {processed}/{total} triggers processed")
|
||||
|
||||
console.print()
|
||||
|
||||
|
||||
def _run_status_command() -> None:
|
||||
"""Display memory pool status."""
|
||||
from ..handlers.intake.pool_processor import get_pool_status
|
||||
|
||||
console.print()
|
||||
|
||||
status = get_pool_status()
|
||||
|
||||
json_handler.log_operation(
|
||||
"pool_status_command",
|
||||
{"files_in_pool": status.get("files_in_pool", 0)},
|
||||
)
|
||||
|
||||
enabled = "[green]enabled[/green]" if status.get("enabled") else "[red]disabled[/red]"
|
||||
console.print(f"[bold cyan]Memory Pool Status[/bold cyan] ({enabled})")
|
||||
console.print()
|
||||
console.print(f" Files in pool: {status.get('files_in_pool', 0)}")
|
||||
console.print(f" Keep recent: {status.get('keep_recent', 0)}")
|
||||
console.print(f" Vectors stored: {status.get('vectors_stored', 0)}")
|
||||
console.print(f" Collection: {status.get('collection_name', 'unknown')}")
|
||||
|
||||
newest = status.get("newest_file")
|
||||
oldest = status.get("oldest_file")
|
||||
if newest:
|
||||
console.print(f" Newest file: {newest}")
|
||||
if oldest and oldest != newest:
|
||||
console.print(f" Oldest file: {oldest}")
|
||||
|
||||
console.print()
|
||||
|
||||
|
||||
def print_introspection() -> None:
|
||||
"""Display pool module introspection."""
|
||||
console.print()
|
||||
console.print("[bold cyan]Pool Module - Memory Pool Processing[/bold cyan]")
|
||||
console.print()
|
||||
console.print("[dim]Processes memory_pool/ files and checks rollover triggers[/dim]")
|
||||
console.print()
|
||||
for sub, desc in _SUBCOMMANDS.items():
|
||||
console.print(f" [cyan]*[/cyan] {sub} — {desc}")
|
||||
console.print()
|
||||
|
||||
|
||||
def print_help() -> None:
|
||||
"""Display pool module help."""
|
||||
console.print()
|
||||
console.print(
|
||||
Panel.fit(
|
||||
"[bold cyan]Pool Module - Memory Pool & Auto-Processing[/bold cyan]\n"
|
||||
"[dim]Vectorize pool files, check rollover, manual or hook-driven[/dim]",
|
||||
border_style="cyan",
|
||||
box=box.ROUNDED,
|
||||
)
|
||||
)
|
||||
console.print()
|
||||
console.print("[bold cyan]COMMANDS:[/bold cyan]")
|
||||
console.print()
|
||||
console.print(" [green]pool process[/green] Process pool files + check/run rollover")
|
||||
console.print(" [green]pool status[/green] Show pool file count, config, vector stats")
|
||||
console.print()
|
||||
console.print("[bold cyan]USAGE:[/bold cyan]")
|
||||
console.print()
|
||||
console.print(" [dim]drone @memory pool process[/dim]")
|
||||
console.print(" [dim]drone @memory pool status[/dim]")
|
||||
console.print()
|
||||
Reference in New Issue
Block a user