fix(memory): unwedge plan vectorization — salt vector IDs with source file, per-file batches with incremental manifest (DPLAN-0245). Backlog drained 229/229, 990 tests green
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user