From 973e8ba21973725ab872802b5deb5c65e7c0f010 Mon Sep 17 00:00:00 2001 From: Victor Giers Date: Mon, 15 Jun 2026 02:03:29 +0200 Subject: [PATCH] Enhance RAG pipeline by implementing robust metadata persistence, extracting detailed information (summary, keywords, entities) into the enhanced corpus record. --- backend/local_rag.py | 134 +++++++++++++++++++++++++++---------------- 1 file changed, 83 insertions(+), 51 deletions(-) diff --git a/backend/local_rag.py b/backend/local_rag.py index 052fcf4..755b55b 100644 --- a/backend/local_rag.py +++ b/backend/local_rag.py @@ -528,6 +528,9 @@ def _mark_entry_pending(entry: Dict[str, Any]) -> None: entry["sync_status"] = "pending" entry.pop("sync_error", None) entry.pop("synced_at", None) + metadata = entry.get("metadata") if isinstance(entry.get("metadata"), dict) else {} + entry["metadata"] = {**metadata, "status": "pending"} + entry["metadata"].pop("error", None) async def _save_library_change(slug: str, data: Dict[str, Any]) -> Optional[str]: @@ -585,6 +588,62 @@ def _apply_source_metadata(corpus_path: Path, files: List[Dict[str, Any]]) -> No tmp.replace(corpus_path) +def _persist_item_metadata(slug: str, enhanced_path: Path) -> None: + if not enhanced_path.exists(): + return + + path = lib_json(slug) + if not path.exists(): + return + data = _read_json(path) + by_path: Dict[str, Dict[str, Any]] = {} + for entry in data.get("files", []): + raw_path = str(entry.get("path") or "") + try: + key = str(Path(raw_path).expanduser().resolve()) + except Exception: + key = raw_path + if key: + by_path[key] = entry + + changed = False + with enhanced_path.open("r", encoding="utf-8", errors="ignore") as handle: + for line in handle: + if not line.strip(): + continue + try: + record = json.loads(line) + except Exception: + continue + raw_source = str(record.get("source_path") or "") + try: + source_key = str(Path(raw_source).expanduser().resolve()) + except Exception: + source_key = raw_source + entry = by_path.get(source_key) + if not entry: + continue + enrichment_meta = record.get("enrichment_meta") if isinstance(record.get("enrichment_meta"), dict) else {} + ok = bool(enrichment_meta.get("ok", True)) + entry["metadata"] = { + "status": "ready" if ok else "fallback", + "headline": str(record.get("headline") or ""), + "summary": str(record.get("summary") or ""), + "keywords": list(record.get("keywords") or [])[:12], + "entities": list(record.get("entities") or [])[:12], + "language": record.get("lang"), + "level": str(enrichment_meta.get("level") or "standard"), + "model": enrichment_meta.get("model"), + "strategy": enrichment_meta.get("strategy"), + "qa_count": len(record.get("qa") or []), + "updated_at": now_iso(), + "error": enrichment_meta.get("error") if not ok else None, + } + changed = True + if changed: + write_library(slug, data) + + def _stage_name(sha: str, path: Path) -> str: safe_name = re.sub(r"[^A-Za-z0-9._-]+", "_", path.name).strip("._") or "file" return f"{sha}--{safe_name}" @@ -907,33 +966,15 @@ def _build_merged_shadow_corpus(corpus_path: Path, partial_shadow_path: Path, ou def _run_selected_enrichment(slug: str, on_progress=None, **opts) -> Dict[str, Any]: data = read_library(slug) files = list(data.get("files", [])) - selected_paths = _selected_enrichment_paths(files) + deep_paths = _selected_enrichment_paths(files) paths = _collect_library_paths(slug) - if not selected_paths: - _cleanup_enrichment_artifacts(slug) - if on_progress: - on_progress("enrich", 1.0, "Skipping enrichment. All files are raw-indexed.") - return { - "status": "skipped", - "selected_files": 0, - "selected_records": 0, - } - - selected_records = _write_selected_corpus(paths["corpus"], paths["enrich_input"], selected_paths) - if selected_records == 0: - _cleanup_enrichment_artifacts(slug) - if on_progress: - on_progress("enrich", 1.0, "No selected records were eligible for enrichment.") - return { - "status": "skipped", - "selected_files": len(selected_paths), - "selected_records": 0, - } + if not paths["corpus"].exists(): + raise RuntimeError("Build the corpus before generating metadata.") enrich_runner = _load_pipeline_fn("corpus_enricher", "run_enrich") result = enrich_runner( - inp=paths["enrich_input"], + inp=paths["corpus"], out=paths["enhanced"], shadow_out=paths["shadow_partial"], on_progress=on_progress, @@ -944,9 +985,11 @@ def _run_selected_enrichment(slug: str, on_progress=None, **opts) -> Dict[str, A min_chars=opts.get("min_chars", DEFAULT_ENRICH_MIN_CHARS), max_text=opts.get("max_text", DEFAULT_ENRICH_MAX_TEXT), cache_dir=opts.get("cache_dir", str(paths["base"] / ".rag_cache")), + deep_source_paths=deep_paths, ) - result["selected_files"] = len(selected_paths) - result["selected_records"] = selected_records + _persist_item_metadata(slug, paths["enhanced"]) + result["metadata_records"] = _line_count(paths["enhanced"]) + result["deep_files"] = len(deep_paths) result["shadow_records"] = _build_merged_shadow_corpus(paths["corpus"], paths["shadow_partial"], paths["shadow"]) return result @@ -965,7 +1008,6 @@ def _run_prepare_pipeline(slug: str, on_progress=None, **opts): paths = _collect_library_paths(slug) states = dict(payload.get("states") or {}) results: Dict[str, Any] = {} - enrichment_enabled_files = int(states.get("enrichment_enabled_files") or 0) build_runner = _load_pipeline_fn("corpus_builder", "run_build") index_runner = _load_pipeline_fn("index_builder", "run_index") @@ -994,32 +1036,22 @@ def _run_prepare_pipeline(slug: str, on_progress=None, **opts): states["is_enriched"] = False states["is_indexed"] = False - if enrichment_enabled_files > 0: - if not states.get("is_enriched"): - enrich_progress = _scaled_progress(on_progress, 0.34, 0.69, "Enriching selected files") if on_progress else None - results["enrich"] = _run_selected_enrichment( - slug, - on_progress=enrich_progress, - ollama=_resolve_ollama_url(opts.get("ollama")), - enrich_model=opts.get("enrich_model", DEFAULT_ENRICH_MODEL), - summary_lang=opts.get("summary_lang", "auto"), - enrich_concurrency=opts.get("enrich_concurrency", DEFAULT_ENRICH_CONCURRENCY), - min_chars=opts.get("min_chars", DEFAULT_ENRICH_MIN_CHARS), - max_text=opts.get("max_text", DEFAULT_ENRICH_MAX_TEXT), - ) - _mark_pipeline_stage(slug, "enrich", prepare_signature) - states["is_enriched"] = True - states["is_indexed"] = False - else: - _cleanup_enrichment_artifacts(slug) - results["enrich"] = { - "status": "skipped", - "selected_files": 0, - "selected_records": 0, - } - if on_progress: - on_progress("enrich", 0.69, "Skipping enrichment. Files will stay raw-indexed.") - states["is_enriched"] = False + if not states.get("has_metadata"): + enrich_progress = _scaled_progress(on_progress, 0.34, 0.69, "Generating metadata") if on_progress else None + results["enrich"] = _run_selected_enrichment( + slug, + on_progress=enrich_progress, + ollama=_resolve_ollama_url(opts.get("ollama")), + enrich_model=opts.get("enrich_model", DEFAULT_ENRICH_MODEL), + summary_lang=opts.get("summary_lang", "auto"), + enrich_concurrency=opts.get("enrich_concurrency", DEFAULT_ENRICH_CONCURRENCY), + min_chars=opts.get("min_chars", DEFAULT_ENRICH_MIN_CHARS), + max_text=opts.get("max_text", DEFAULT_ENRICH_MAX_TEXT), + ) + _mark_pipeline_stage(slug, "enrich", prepare_signature) + states["has_metadata"] = True + states["is_enriched"] = int(states.get("enrichment_enabled_files") or 0) > 0 + states["is_indexed"] = False if not states.get("is_indexed"): index_progress = _scaled_progress(on_progress, 0.69, 1.0, "Building search indexes") if on_progress else None