"""Local standards knowledge base ingestion. Indexes PDFs/DOCX from configured folders into the vector store under a reserved tenant id. This allows the agentic pipeline to cite RICS standards/exemplar material without requiring uploads through the UI. Files are NOT committed to git; the ingestion simply reads from disk. """ from __future__ import annotations import hashlib import json import logging from pathlib import Path from langchain_core.documents import Document from app.config import settings from app.ingest.hierarchy import build_hierarchical_documents from app.ingest.pipeline import load_raw_documents, stamp_chunks_for_index from app.retrieval.chunk_role import classify_chunk_role from app.vectorstore.factory import get_vectorstore logger = logging.getLogger(__name__) _KB_MANIFEST = "kb_manifest.json" def _kb_manifest_path() -> Path: return Path(settings.faiss_index_path) / _KB_MANIFEST def _file_fingerprint(fp: Path) -> str: st = fp.stat() return f"{int(st.st_mtime_ns)}:{st.st_size}" def _load_kb_manifest() -> dict[str, str]: path = _kb_manifest_path() if not path.is_file(): return {} try: data = json.loads(path.read_text(encoding="utf-8")) if isinstance(data, dict): return {str(k): str(v) for k, v in data.items()} except Exception as exc: # noqa: BLE001 logger.warning("Could not read KB manifest %s: %s", path, exc) return {} def _save_kb_manifest(manifest: dict[str, str]) -> None: path = _kb_manifest_path() path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps(manifest, indent=0, sort_keys=True), encoding="utf-8") def _stable_kb_doc_id(path: Path) -> str: # Stable across restarts; changes if file path changes. h = hashlib.sha1(str(path).encode("utf-8"), usedforsecurity=False).hexdigest()[:20] return f"kb-{h}" def iter_knowledge_base_roots() -> list[Path]: """Resolved directory (or single-file parent) roots listed in ``knowledge_base_dirs``.""" raw = (settings.knowledge_base_dirs or "").strip() if not raw: return [] parts = [p.strip() for p in raw.split(",") if p.strip()] root = Path(__file__).resolve().parents[2] # repo root out: list[Path] = [] for part in parts: p = Path(part) if not p.is_absolute(): p = root / p if not p.exists(): continue out.append(p if p.is_dir() else p.parent) # de-dup while preserving order seen: set[str] = set() uniq: list[Path] = [] for p in out: k = str(p.resolve()).lower() if k in seen: continue seen.add(k) uniq.append(p) return uniq def iter_kb_files() -> list[Path]: """List all candidate KB documents from configured dirs.""" raw = (settings.knowledge_base_dirs or "").strip() if not raw: return [] parts = [p.strip() for p in raw.split(",") if p.strip()] root = Path(__file__).resolve().parents[2] # repo root out: list[Path] = [] for part in parts: p = Path(part) if not p.is_absolute(): p = root / p if not p.exists(): continue if p.is_file(): out.append(p) continue for ext in (".pdf", ".docx"): out.extend(sorted(p.rglob(f"*{ext}"))) # de-dup seen: set[str] = set() uniq: list[Path] = [] for p in out: k = str(p.resolve()).lower() if k in seen: continue seen.add(k) uniq.append(p) return uniq def upsert_knowledge_base() -> dict[str, int]: """(Re)index all KB files into the vector store. Returns counts for observability; safe to call at startup. """ if not settings.knowledge_base_enabled: return {"files": 0, "chunks": 0, "skipped": 0} files = iter_kb_files() if not files: logger.info("KB ingest: no knowledge base files found") return {"files": 0, "chunks": 0, "skipped": 0} vs = get_vectorstore() tenant_id = settings.knowledge_base_tenant_id chunks_total = 0 skipped = 0 skipped_unchanged = 0 manifest = _load_kb_manifest() if settings.knowledge_base_skip_unchanged else {} new_manifest: dict[str, str] = dict(manifest) for fp in files: suf = fp.suffix.lower() if suf not in (".pdf", ".docx"): skipped += 1 continue doc_id = _stable_kb_doc_id(fp) fp_key = str(fp.resolve()) fingerprint = _file_fingerprint(fp) if ( settings.knowledge_base_skip_unchanged and manifest.get(fp_key) == fingerprint and vs.count_for_doc(doc_id) > 0 ): skipped_unchanged += 1 continue try: # Best-effort upsert: delete prior chunks for this kb doc_id. vs.delete_document(doc_id) raw_docs = load_raw_documents(fp) # Ensure a reasonable source label even without DB Document rows. label = fp.name for d in raw_docs: d.metadata = {**(d.metadata or {}), "source": label} hier = build_hierarchical_documents( raw_docs, doc_id=doc_id, filename=label, chunk_size=settings.chunk_size, chunk_overlap=settings.chunk_overlap, ) stamped = stamp_chunks_for_index(hier, doc_id=doc_id, tenant_id=tenant_id) # Mark provenance as knowledge base for downstream agents and # classify each chunk so retrieval can bias toward boilerplate at # assembly tier. for d in stamped: d.metadata["kb"] = True d.metadata["kb_path"] = str(fp) d.metadata["chunk_role"] = classify_chunk_role(d.page_content or "") vs.add_documents(stamped, persist=False) new_manifest[fp_key] = fingerprint chunks_total += len(stamped) except Exception as exc: # noqa: BLE001 logger.warning("KB ingest failed for %s: %s", fp, exc) skipped += 1 flush = getattr(vs, "flush", None) if callable(flush): flush() if new_manifest != manifest: _save_kb_manifest(new_manifest) if chunks_total: import asyncio from app.retrieval.semantic_cache import invalidate_semantic_cache_for_tenant try: asyncio.get_running_loop() pending = True except RuntimeError: pending = False if pending: asyncio.create_task(invalidate_semantic_cache_for_tenant(tenant_id)) else: asyncio.run(invalidate_semantic_cache_for_tenant(tenant_id)) logger.info( "KB ingest complete: files=%d chunks=%d skipped=%d unchanged=%d", len(files), chunks_total, skipped, skipped_unchanged, ) return { "files": len(files), "chunks": chunks_total, "skipped": skipped, "unchanged": skipped_unchanged, }