feat(memory): fix plan vectorization: batch subprocess, decouple from startup, add process-plans command
Co-Authored-By: @memory <memory@aipass>
This commit is contained in:
@@ -259,10 +259,16 @@ def process_plans() -> Dict[str, Any]:
|
||||
|
||||
logger.info(f"[plans] Found {len(unprocessed)} unprocessed plan files")
|
||||
|
||||
total_chunks = 0
|
||||
files_processed = 0
|
||||
errors = []
|
||||
|
||||
# -- Phase 1: Read all files, chunk them, collect texts + metadatas ----------
|
||||
all_texts: List[str] = []
|
||||
all_metadatas: List[Dict[str, str]] = []
|
||||
# Track which files produced chunks (for manifest update)
|
||||
files_with_chunks: List[Path] = []
|
||||
# Files with 0 chunks still get marked in manifest (e.g. template content)
|
||||
files_without_chunks: List[Path] = []
|
||||
|
||||
for plan_file in unprocessed:
|
||||
try:
|
||||
text = plan_file.read_text(encoding='utf-8')
|
||||
@@ -270,58 +276,65 @@ def process_plans() -> Dict[str, Any]:
|
||||
errors.append(f'{plan_file.name}: read error: {e}')
|
||||
continue
|
||||
|
||||
# Chunk the plan
|
||||
chunks = _chunk_plan_text(text, plan_file.name)
|
||||
if not chunks:
|
||||
manifest[plan_file.name] = datetime.now().isoformat()
|
||||
files_without_chunks.append(plan_file)
|
||||
continue
|
||||
|
||||
# Extract texts and build metadata
|
||||
texts = [c['text'] for c in chunks]
|
||||
metadatas = [
|
||||
{
|
||||
files_with_chunks.append(plan_file)
|
||||
for c in chunks:
|
||||
all_texts.append(c['text'])
|
||||
all_metadatas.append({
|
||||
'source_file': plan_file.name,
|
||||
'section': c['section'],
|
||||
'processed_at': datetime.now().isoformat(),
|
||||
'type': 'plan'
|
||||
}
|
||||
for c in chunks
|
||||
]
|
||||
})
|
||||
|
||||
# Embed
|
||||
embed_result = _embed_texts(texts)
|
||||
if not embed_result.get('success'):
|
||||
errors.append(f"{plan_file.name}: embed error: {embed_result.get('error')}")
|
||||
continue
|
||||
total_chunks = len(all_texts)
|
||||
files_processed = 0
|
||||
|
||||
embeddings = embed_result.get('embeddings', [])
|
||||
if not embeddings:
|
||||
errors.append(f'{plan_file.name}: no embeddings returned')
|
||||
continue
|
||||
|
||||
# Store
|
||||
store_result = _store_vectors(embeddings, texts, metadatas, collection_name)
|
||||
if not store_result.get('success'):
|
||||
errors.append(f"{plan_file.name}: store error: {store_result.get('error')}")
|
||||
continue
|
||||
|
||||
# Mark as processed
|
||||
# Mark empty-chunk files in manifest immediately (nothing to embed)
|
||||
for plan_file in files_without_chunks:
|
||||
manifest[plan_file.name] = datetime.now().isoformat()
|
||||
files_processed += 1
|
||||
total_chunks += len(chunks)
|
||||
logger.info(f"[plans] Processed {plan_file.name}: {len(chunks)} chunks vectorized")
|
||||
|
||||
# Save manifest
|
||||
# -- Phase 2: Batch embed + store (single subprocess each) ------------------
|
||||
if all_texts:
|
||||
logger.info(f"[plans] Batch embedding {total_chunks} chunks from {len(files_with_chunks)} files")
|
||||
|
||||
embed_result = _embed_texts(all_texts)
|
||||
if not embed_result.get('success'):
|
||||
error_msg = f"batch embed error: {embed_result.get('error')}"
|
||||
logger.error(f"[plans] {error_msg}")
|
||||
errors.append(error_msg)
|
||||
else:
|
||||
embeddings = embed_result.get('embeddings', [])
|
||||
if not embeddings:
|
||||
errors.append('batch embed returned no embeddings')
|
||||
else:
|
||||
store_result = _store_vectors(embeddings, all_texts, all_metadatas, collection_name)
|
||||
if not store_result.get('success'):
|
||||
error_msg = f"batch store error: {store_result.get('error')}"
|
||||
logger.error(f"[plans] {error_msg}")
|
||||
errors.append(error_msg)
|
||||
else:
|
||||
# Success — mark all chunk-producing files in manifest
|
||||
for plan_file in files_with_chunks:
|
||||
manifest[plan_file.name] = datetime.now().isoformat()
|
||||
files_processed = len(files_with_chunks)
|
||||
logger.info(f"[plans] Batch complete: {files_processed} files, {total_chunks} chunks vectorized")
|
||||
|
||||
# Save manifest (includes empty-chunk files even if embedding failed)
|
||||
_save_manifest(manifest)
|
||||
|
||||
result = {
|
||||
'success': files_processed > 0 or not errors,
|
||||
result: Dict[str, Any] = {
|
||||
'success': files_processed > 0 or (not errors and not files_with_chunks),
|
||||
'files_processed': files_processed,
|
||||
'total_chunks': total_chunks,
|
||||
'total_chunks': total_chunks if files_processed > 0 else 0,
|
||||
}
|
||||
if errors:
|
||||
result['errors'] = errors
|
||||
|
||||
json_handler.log_operation("process_plans", {"files_processed": files_processed, "total_chunks": total_chunks, "success": result['success']})
|
||||
json_handler.log_operation("process_plans", {"files_processed": files_processed, "total_chunks": result['total_chunks'], "success": result['success']})
|
||||
|
||||
return result
|
||||
|
||||
@@ -295,12 +295,14 @@ def _check_memory_pool() -> Dict[str, Any]:
|
||||
|
||||
def _check_plans() -> Dict[str, Any]:
|
||||
"""
|
||||
Check plans directory for files to vectorize.
|
||||
Check plans directory for unprocessed files (count only).
|
||||
|
||||
Processes any plan files that haven't been vectorized yet.
|
||||
Does NOT call process_plans() — that spawns heavy ML subprocesses.
|
||||
Only counts pending files and reports. Use 'drone @memory process-plans'
|
||||
to trigger actual vectorization.
|
||||
|
||||
Returns:
|
||||
Dict with processing status
|
||||
Dict with pending file count
|
||||
"""
|
||||
import json
|
||||
|
||||
@@ -308,7 +310,7 @@ def _check_plans() -> Dict[str, Any]:
|
||||
|
||||
# Load config
|
||||
try:
|
||||
with open(config_path) as f:
|
||||
with open(config_path, 'r', encoding='utf-8') as f:
|
||||
config = json.load(f)
|
||||
plans_config = config.get('plans', {})
|
||||
except Exception:
|
||||
@@ -320,7 +322,8 @@ def _check_plans() -> Dict[str, Any]:
|
||||
|
||||
# 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
|
||||
repo_root = _find_repo_root()
|
||||
plans_path = Path(plans_dir) if Path(plans_dir).is_absolute() else repo_root / plans_dir
|
||||
extensions = plans_config.get('supported_extensions', ['.md'])
|
||||
|
||||
if not plans_path.exists():
|
||||
@@ -333,21 +336,24 @@ def _check_plans() -> Dict[str, Any]:
|
||||
file_count = len(files)
|
||||
|
||||
if file_count == 0:
|
||||
return {'success': True, 'files_in_plans': 0, 'action': 'none'}
|
||||
return {'success': True, 'pending_files': 0, 'action': 'count_only'}
|
||||
|
||||
# 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'}
|
||||
# Load manifest to count unprocessed files
|
||||
manifest_path = _MEMORY_ROOT / "config" / ".plans_processed.json"
|
||||
manifest: Dict[str, str] = {}
|
||||
if manifest_path.exists():
|
||||
try:
|
||||
manifest = json.loads(manifest_path.read_text(encoding='utf-8'))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
pending = [f for f in files if f.name not in manifest]
|
||||
pending_count = len(pending)
|
||||
|
||||
if pending_count > 0:
|
||||
logger.info(f"[plans] {pending_count} plans pending vectorization. Run: drone @memory process-plans")
|
||||
|
||||
return {'success': True, 'pending_files': pending_count, 'action': 'count_only'}
|
||||
|
||||
|
||||
def _check_code_archive() -> Dict[str, Any]:
|
||||
|
||||
@@ -133,6 +133,10 @@ def handle_command(command: str, args: List[str]) -> bool:
|
||||
sync_line_counts()
|
||||
return True
|
||||
|
||||
elif command == 'process-plans':
|
||||
process_plans_command()
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
|
||||
@@ -226,6 +230,61 @@ def run_rollover() -> bool:
|
||||
return success_count > 0
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# PLAN VECTORIZATION
|
||||
# =============================================================================
|
||||
|
||||
def process_plans_command() -> None:
|
||||
"""
|
||||
Process pending plan files into vector storage.
|
||||
|
||||
Batches all chunks from all files into a single embed + store call.
|
||||
"""
|
||||
console.print()
|
||||
console.print(Panel.fit(
|
||||
"[bold cyan]Memory - Process Plans[/bold cyan]",
|
||||
border_style="cyan",
|
||||
box=box.ROUNDED
|
||||
))
|
||||
console.print()
|
||||
|
||||
console.print("[cyan]Processing plan files into vector storage...[/cyan]")
|
||||
console.print()
|
||||
|
||||
try:
|
||||
from ..handlers.intake.plans_processor import process_plans
|
||||
result = process_plans()
|
||||
except Exception as e:
|
||||
error(f"Plan processing failed: {e}")
|
||||
return
|
||||
|
||||
if not result.get('success'):
|
||||
error(result.get('error', 'Unknown error'))
|
||||
if result.get('errors'):
|
||||
for err in result['errors']:
|
||||
error(err)
|
||||
return
|
||||
|
||||
files_processed = result.get('files_processed', 0)
|
||||
total_chunks = result.get('total_chunks', 0)
|
||||
reason = result.get('reason', '')
|
||||
|
||||
if files_processed == 0 and reason:
|
||||
console.print(f"[green]>[/green] {reason}")
|
||||
elif files_processed == 0:
|
||||
console.print("[green]>[/green] No new plans to process")
|
||||
else:
|
||||
console.print(f"[green]>[/green] Processed {files_processed} files ({total_chunks} chunks vectorized)")
|
||||
|
||||
if result.get('errors'):
|
||||
console.print()
|
||||
for err in result['errors']:
|
||||
error(err)
|
||||
|
||||
console.print()
|
||||
json_handler.log_operation("process_plans_command", {"files_processed": files_processed, "total_chunks": total_chunks})
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# LINE COUNT SYNC
|
||||
# =============================================================================
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
},
|
||||
"plans": {
|
||||
"enabled": true,
|
||||
"path": "src/aipass/flow/backup/processed_plans",
|
||||
"path": "src/aipass/backup/processed_plans",
|
||||
"supported_extensions": [".md"],
|
||||
"collection_name": "flow_plans"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user