diff --git a/src/aipass/memory/apps/handlers/intake/plans_processor.py b/src/aipass/memory/apps/handlers/intake/plans_processor.py index cb11e0cb..c3417c80 100644 --- a/src/aipass/memory/apps/handlers/intake/plans_processor.py +++ b/src/aipass/memory/apps/handlers/intake/plans_processor.py @@ -164,7 +164,7 @@ def _save_manifest(manifest: Dict[str, str]) -> None: # ============================================================================= -def _embed_texts(texts: List[str]) -> dict: +def _embed_texts(texts: List[str], timeout: int = 120) -> dict: """Encode texts via subprocess.""" input_data = json.dumps({"texts": texts}) try: @@ -173,7 +173,7 @@ def _embed_texts(texts: List[str]) -> dict: input=input_data, capture_output=True, text=True, - timeout=120, + timeout=timeout, ) if result.returncode != 0: return {"success": False, "error": result.stderr or "Embedding failed"} @@ -219,24 +219,17 @@ def process_plans() -> Dict[str, Any]: """ Process plan files from flow/processed_plans/ into vector storage. - Workflow: - 1. Load config to find plans directory - 2. Scan for unprocessed .md files - 3. Chunk each file into sections - 4. Embed all chunks via subprocess - 5. Store vectors in ChromaDB - 6. Update processed manifest + Processes each file independently so partial failure makes partial + progress — manifest is saved after every successful file. Returns: Dict with success, files_processed, total_chunks """ - # Load config plans_config = config_loader.section("plans") if not plans_config.get("enabled", False): return {"success": True, "skipped": True, "reason": "plans disabled"} - # Resolve plans directory (relative to repo root) plans_dir = plans_config.get("path", ".backup/processed_plans") repo_root = _find_repo_root() plans_path = Path(plans_dir) if Path(plans_dir).is_absolute() else repo_root / plans_dir @@ -246,7 +239,6 @@ def process_plans() -> Dict[str, Any]: if not plans_path.exists(): return {"success": True, "files_processed": 0, "total_chunks": 0, "reason": "plans dir not found"} - # Get plan files files = [] for ext in extensions: files.extend(plans_path.glob(f"*{ext}")) @@ -254,7 +246,6 @@ def process_plans() -> Dict[str, Any]: if not files: return {"success": True, "files_processed": 0, "total_chunks": 0} - # Load manifest to skip already-processed files manifest = _load_manifest() unprocessed = [f for f in files if f.name not in manifest] @@ -263,88 +254,73 @@ def process_plans() -> Dict[str, Any]: logger.info(f"[plans] Found {len(unprocessed)} unprocessed plan files") - 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] = [] + errors: List[str] = [] + files_processed = 0 + total_chunks = 0 for plan_file in unprocessed: try: text = plan_file.read_text(encoding="utf-8") except Exception as e: - logger.warning(f"[plans_processor] Failed to read plan file {plan_file.name}: {e}") + logger.warning(f"[plans] Failed to read {plan_file.name}: {e}") errors.append(f"{plan_file.name}: read error: {e}") continue chunks = _chunk_plan_text(text, plan_file.name) if not chunks: - files_without_chunks.append(plan_file) + manifest[plan_file.name] = datetime.now().isoformat() + _save_manifest(manifest) continue - 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", - } - ) + texts = [c["text"] for c in chunks] + metadatas = [ + { + "source_file": plan_file.name, + "section": c["section"], + "processed_at": datetime.now().isoformat(), + "type": "plan", + } + for c in chunks + ] - total_chunks = len(all_texts) - files_processed = 0 - - # Mark empty-chunk files in manifest immediately (nothing to embed) - for plan_file in files_without_chunks: - manifest[plan_file.name] = datetime.now().isoformat() - - # -- 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) + timeout = max(30, len(texts) * 3) + embed_result = _embed_texts(texts, timeout=timeout) 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") + logger.warning(f"[plans] Embed failed for {plan_file.name}: {embed_result.get('error')}") + errors.append(f"{plan_file.name}: embed error: {embed_result.get('error')}") + continue - # Save manifest (includes empty-chunk files even if embedding failed) - _save_manifest(manifest) + embeddings = embed_result.get("embeddings", []) + if not embeddings: + errors.append(f"{plan_file.name}: embed returned no embeddings") + continue + + store_result = _store_vectors(embeddings, texts, metadatas, collection_name) + if not store_result.get("success"): + logger.warning(f"[plans] Store failed for {plan_file.name}: {store_result.get('error')}") + errors.append(f"{plan_file.name}: store error: {store_result.get('error')}") + continue + + manifest[plan_file.name] = datetime.now().isoformat() + _save_manifest(manifest) + files_processed += 1 + total_chunks += len(texts) + logger.info(f"[plans] {plan_file.name}: {len(texts)} chunks vectorized") + + if files_processed > 0: + logger.info(f"[plans] Complete: {files_processed} files, {total_chunks} chunks vectorized") result: Dict[str, Any] = { - "success": files_processed > 0 or (not errors and not files_with_chunks), + "success": files_processed > 0 or not errors, "files_processed": files_processed, - "total_chunks": total_chunks if files_processed > 0 else 0, + "total_chunks": total_chunks, } if errors: result["errors"] = errors json_handler.log_operation( "process_plans", - {"files_processed": files_processed, "total_chunks": result["total_chunks"], "success": result["success"]}, + {"files_processed": files_processed, "total_chunks": total_chunks, "success": result["success"]}, ) return result diff --git a/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py b/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py index 6a04dcc8..18b867e9 100755 --- a/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py +++ b/src/aipass/memory/apps/handlers/storage/chroma_subprocess.py @@ -67,12 +67,30 @@ def _store_vectors(branch, memory_type, embeddings, documents, metadatas, db_pat embedding_function=None, ) - # Content-hash IDs prevent duplicates across rollover runs - ids = [f"{branch}_{memory_type}_{hashlib.sha256(doc.encode()).hexdigest()[:16]}" for doc in documents] + # Content-hash IDs — idempotent across runs. When metadata carries + # source_file (e.g. plans intake), salt the hash so identical boilerplate + # from different files gets distinct IDs and per-file provenance survives. + ids = [] + for doc, meta in zip(documents, metadatas): + salt = meta.get("source_file", "") if isinstance(meta, dict) else "" + hash_input = f"{salt}:{doc}" if salt else doc + ids.append(f"{branch}_{memory_type}_{hashlib.sha256(hash_input.encode()).hexdigest()[:16]}") # Chroma expects lists, not numpy arrays embeddings_list = [emb.tolist() if hasattr(emb, "tolist") else emb for emb in embeddings] + # Safety net: deduplicate within batch — ChromaDB rejects non-unique IDs + # in a single upsert call. + seen = {} + for i, doc_id in enumerate(ids): + seen[doc_id] = i + if len(seen) < len(ids): + unique_indices = sorted(seen.values()) + ids = [ids[i] for i in unique_indices] + embeddings_list = [embeddings_list[i] for i in unique_indices] + documents = [documents[i] for i in unique_indices] + metadatas = [metadatas[i] for i in unique_indices] + # Upsert: idempotent — same content gets same ID, no duplicates collection.upsert(embeddings=embeddings_list, documents=documents, metadatas=metadatas, ids=ids) diff --git a/src/aipass/memory/tests/test_plans_processor.py b/src/aipass/memory/tests/test_plans_processor.py index d2728e7c..3803e348 100644 --- a/src/aipass/memory/tests/test_plans_processor.py +++ b/src/aipass/memory/tests/test_plans_processor.py @@ -501,7 +501,7 @@ class TestProcessPlans: monkeypatch.setattr( mod, "_embed_texts", - lambda texts: {"success": True, "embeddings": [[0.1, 0.2]] * len(texts)}, + lambda texts, timeout=120: {"success": True, "embeddings": [[0.1, 0.2]] * len(texts)}, ) monkeypatch.setattr( mod, @@ -547,7 +547,7 @@ class TestProcessPlans: monkeypatch.setattr( mod, "_embed_texts", - lambda texts: {"success": False, "error": "GPU out of memory"}, + lambda texts, timeout=120: {"success": False, "error": "GPU out of memory"}, ) mock_jh = MagicMock() @@ -584,7 +584,7 @@ class TestProcessPlans: monkeypatch.setattr( mod, "_embed_texts", - lambda texts: {"success": True, "embeddings": []}, + lambda texts, timeout=120: {"success": True, "embeddings": []}, ) mock_jh = MagicMock() @@ -620,7 +620,7 @@ class TestProcessPlans: monkeypatch.setattr( mod, "_embed_texts", - lambda texts: {"success": True, "embeddings": [[0.1, 0.2]] * len(texts)}, + lambda texts, timeout=120: {"success": True, "embeddings": [[0.1, 0.2]] * len(texts)}, ) monkeypatch.setattr( mod,