Download src/knowledge_ingest/knowledge_store.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 45.7 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_ingest/knowledge_store.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/knowledge_ingest/knowledge_store.py
-
curl -L -o knowledge_store.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_ingest/knowledge_store.py
45.7 kB
| """KnowledgeStore - persisting artifacts, runs and candidate entries (N4). | |
| Writes the three tables the pipeline produces: | |
| knowledge_documents one parsed artifact version (content-addressed) | |
| knowledge_extraction_runs one row per extraction of one document | |
| knowledge_entries one row per candidate, all four kinds | |
| Schema of record: `docs/knowledge/KNOWLEDGE_PERSISTENCE_CONTRACT.md` section 2. | |
| Go owns the dedorch migration (CLAUDE.md 2.2); this module only INSERTs and | |
| SELECTs and executes no DDL, ever. | |
| **Deliberately NOT a never-throw seam.** The seams listed in CLAUDE.md 5.4 exist | |
| so a *chat turn* degrades to soft output rather than 500ing at the user. An | |
| ingestion job is the opposite case: it is asynchronous, the caller polls for its | |
| outcome, and a job that silently lost its rows while reporting success is | |
| strictly worse than one that reports `failed`. Failures here propagate and the | |
| job settles `failed`. Do not "fix" this into the house pattern. | |
| Two invariants worth stating because they are easy to break later: | |
| - **`payload` is the source of truth.** The promoted columns (`label`, | |
| `mention_count`, `extraction_status`, `diff_status`, `definition_conflict`) | |
| are a read optimisation so the review queue can be ordered in SQL without | |
| opening jsonb. A writer sets both from the same object in the same insert. | |
| - **Entries are written per run, never upserted across runs.** Two extractions | |
| of one document produce two full sets of rows sharing `entity_id`s. That is | |
| the point: it is what makes a quality change attributable. | |
| """ | |
| from __future__ import annotations | |
| import uuid | |
| from datetime import UTC, datetime | |
| from typing import Any | |
| from sqlalchemy import select, update | |
| from src.db.postgres.connection import AsyncSessionLocal | |
| from src.db.postgres.models import ( | |
| KnowledgeDocumentRow, | |
| KnowledgeEntryRow, | |
| KnowledgeExtractionRunRow, | |
| KnowledgeReviewRow, | |
| ) | |
| from src.knowledge_extraction import keys, mutation | |
| from src.middlewares.logging import get_logger | |
| logger = get_logger("knowledge_store") | |
| # How each branch's output maps onto a `knowledge_entries` row. The id field | |
| # differs per kind because each is derived from different content (`ids.py`), | |
| # and the label is whatever a reviewer would recognise the row by in a list. | |
| _ENTRY_KINDS: dict[str, tuple[str, str]] = { | |
| "glossary": ("term_id", "term"), | |
| "rule": ("rule_id", "statement"), | |
| "formula": ("formula_id", "name"), | |
| # The per-document brief. `domain` is reserved for the scope-level | |
| # semantic model composed in `src/knowledge_domain/` (v4). | |
| "document": ("brief_id", "title"), | |
| # Scope-level, not produced by a run - see `save_declaration`. | |
| "domain": ("domain_id", "domain_name"), | |
| } | |
| # m4: the body a reviewer or a consumer actually wants, per kind. First | |
| # non-empty wins, so a kind whose primary field abstained still promotes | |
| # something rather than silently leaving the column null. | |
| _CONTENT_FIELDS: dict[str, tuple[str, ...]] = { | |
| "glossary": ("definition",), | |
| "rule": ("statement", "consequence"), | |
| "formula": ("formula_latex", "name"), | |
| "document": ("purpose_verbatim", "title"), | |
| "domain": ("purpose_verbatim", "domain_name"), | |
| } | |
| # m5: the provenance keys worth lifting out. `page` is 0-based as the parser | |
| # reports it and `page_no` is the 1-based number a human reads; both travel | |
| # because they are derived from one value and cannot disagree. | |
| _PROVENANCE_KEYS = ("doc_id", "page", "page_no", "section_no", "chunk_id") | |
| def _content(kind: str, payload: dict[str, Any]) -> str | None: | |
| """m4: the entry's body, promoted out of `payload`.""" | |
| for field in _CONTENT_FIELDS.get(kind, ()): | |
| value = payload.get(field) | |
| if isinstance(value, str) and value.strip(): | |
| return value.strip() | |
| return None | |
| def _provenance(payload: dict[str, Any]) -> list[dict]: | |
| """m5: provenance as a LIST, so a fused entry can carry several sources. | |
| The extraction contract emits a single object per entry today; it is | |
| wrapped rather than reshaped, because the list is what an entry drawn from | |
| more than one document needs and normalising the shape here means every | |
| consumer reads one thing. | |
| """ | |
| raw = payload.get("provenance") | |
| if isinstance(raw, dict): | |
| raw = [raw] | |
| if not isinstance(raw, list): | |
| return [] | |
| out: list[dict] = [] | |
| for item in raw: | |
| if not isinstance(item, dict): | |
| continue | |
| kept = {k: item[k] for k in _PROVENANCE_KEYS if item.get(k) is not None} | |
| if kept: | |
| out.append(kept) | |
| return out | |
| def _doc_ids(doc_id: str, provenance: list[dict]) -> list[str]: | |
| """m3: every document behind this entry, this one first. | |
| Derived rather than declared, so it cannot drift from the provenance it | |
| summarises. For an ordinary entry it is `[doc_id]`; it grows only when a | |
| definition was carried forward from an earlier document that a later one | |
| re-mentions (i2), which is the fusion case m3 was asked for. | |
| """ | |
| ids = [doc_id] | |
| for item in provenance: | |
| source = item.get("doc_id") | |
| if isinstance(source, str) and source and source not in ids: | |
| ids.append(source) | |
| return ids | |
| def _has_live_source(row: Any, withdrawn: set[str]) -> bool: | |
| """Does this entry still rest on a document that has not been withdrawn? (i7) | |
| Reads `doc_ids` when it is populated and falls back to `doc_id` when it is | |
| not, so rows written before m3 landed keep answering correctly instead of | |
| all looking orphaned at once. | |
| """ | |
| sources = [d for d in (getattr(row, "doc_ids", None) or []) if isinstance(d, str)] | |
| if not sources: | |
| sources = [row.doc_id] | |
| return any(source not in withdrawn for source in sources) | |
| def _label(payload: dict[str, Any], key: str) -> str | None: | |
| """The listing label, trimmed. A rule's `statement` can be a paragraph, and | |
| a column meant for scanning a queue should not carry one.""" | |
| value = payload.get(key) | |
| if not isinstance(value, str): | |
| return None | |
| value = value.strip() | |
| return (value[:197] + "...") if len(value) > 200 else value or None | |
| class KnowledgeStore: | |
| """Insert and read the pipeline's persisted output.""" | |
| async def find_document( | |
| self, doc_id: str, content_hash: str, version: int = 1 | |
| ) -> KnowledgeDocumentRow | None: | |
| """The content-addressing lookup: has this exact artifact been stored? | |
| This is what makes a second click on unchanged documents free. The key | |
| is `(doc_id, content_hash, version)` and it is UNIQUE in the schema, so | |
| a race between two jobs surfaces as an integrity error rather than a | |
| duplicate row. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow).where( | |
| KnowledgeDocumentRow.doc_id == doc_id, | |
| KnowledgeDocumentRow.content_hash == content_hash, | |
| KnowledgeDocumentRow.version == version, | |
| ) | |
| ) | |
| return result.scalars().first() | |
| async def find_document_by_hash( | |
| self, doc_id: str, content_hash: str | |
| ) -> KnowledgeDocumentRow | None: | |
| """The same content-addressing lookup, across every version (i5). | |
| Version is assigned by us now, not read off the artifact, so the reuse | |
| check cannot key on it: computing the next version and *then* looking | |
| it up would miss the row that already holds these exact bytes and | |
| re-parse a document we already have — turning the idempotency guarantee | |
| into a paid no-op. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow) | |
| .where( | |
| KnowledgeDocumentRow.doc_id == doc_id, | |
| KnowledgeDocumentRow.content_hash == content_hash, | |
| ) | |
| .order_by(KnowledgeDocumentRow.version.desc()) | |
| ) | |
| return result.scalars().first() | |
| async def latest_document(self, doc_id: str) -> KnowledgeDocumentRow | None: | |
| """The current version of a document — the highest `version` (i5). | |
| What "current" means for a document that has been amended: the newest | |
| version supersedes, and the older rows stay for the audit trail rather | |
| than being rewritten. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow) | |
| .where(KnowledgeDocumentRow.doc_id == doc_id) | |
| .order_by(KnowledgeDocumentRow.version.desc()) | |
| ) | |
| return result.scalars().first() | |
| async def save_document( | |
| self, | |
| artifact: Any, | |
| *, | |
| scope_id: str, | |
| user_id: str, | |
| document_id: str | None = None, | |
| ) -> KnowledgeDocumentRow: | |
| """Persist one `ParsedDocument`, or return the row that already holds it. | |
| `artifact` is the parsing half's model; it is stored whole in `artifact` | |
| (jsonb) with its provenance lifted into columns. The four fingerprints | |
| are kept separate on purpose - `parser_config` covers the parser's | |
| settings, `normaliser_version` our post-processing, `schema_version` the | |
| artifact's shape, `parser_version` the parser build - because none of | |
| them moves when the others do, and a stale artifact once went unnoticed | |
| for a day precisely because they were conflated. | |
| """ | |
| payload = ( | |
| artifact.model_dump(mode="json") | |
| if hasattr(artifact, "model_dump") | |
| else dict(artifact) | |
| ) | |
| doc_id = payload["doc_id"] | |
| content_hash = payload["content_hash"] | |
| existing = await self.find_document_by_hash(doc_id, content_hash) | |
| if existing is not None: | |
| logger.info( | |
| "knowledge_document_reused", | |
| extra={"doc_id": doc_id, "knowledge_document_id": existing.id}, | |
| ) | |
| return existing | |
| # i5 (ruling 2026-09-14): identity is `doc_id`, so the same document | |
| # arriving with different bytes is a NEW VERSION of it, never a new | |
| # document. That is what keeps an amended file's approvals attached — | |
| # `entity_id` derives from `doc_id`, which does not move. The artifact's | |
| # own `version` is ignored: the parsing half stamps 1 on every parse and | |
| # has no way to know what this tenant already holds. | |
| previous = await self.latest_document(doc_id) | |
| version = (previous.version + 1) if previous else 1 | |
| if previous is not None: | |
| logger.info( | |
| "knowledge_document_versioned", | |
| extra={"doc_id": doc_id, "version": version, "supersedes": previous.id}, | |
| ) | |
| row = KnowledgeDocumentRow( | |
| id=str(uuid.uuid4()), | |
| doc_id=doc_id, | |
| document_id=document_id, | |
| scope_id=scope_id, | |
| user_id=user_id, | |
| source_path=payload.get("source_path") or "", | |
| content_hash=content_hash, | |
| n_pages=int(payload.get("n_pages") or 0), | |
| version=version, | |
| schema_version=payload.get("schema_version") or "", | |
| parser_name=payload.get("parser_name") or "mistral", | |
| parser_version=payload.get("parser_version"), | |
| parser_backend=payload.get("parser_backend"), | |
| parser_config=payload.get("parser_config"), | |
| normaliser_version=payload.get("normaliser_version"), | |
| raw_output_dir=payload.get("raw_output_dir"), | |
| artifact=payload, | |
| ) | |
| async with AsyncSessionLocal() as session: | |
| session.add(row) | |
| await session.commit() | |
| await session.refresh(row) | |
| logger.info( | |
| "knowledge_document_saved", | |
| extra={ | |
| "doc_id": doc_id, | |
| "knowledge_document_id": row.id, | |
| "chunks": len(payload.get("chunks") or []), | |
| }, | |
| ) | |
| return row | |
| async def start_run( | |
| self, | |
| *, | |
| knowledge_document_id: str, | |
| doc_id: str, | |
| scope_id: str, | |
| user_id: str, | |
| branches: list[str], | |
| model_deployment: str | None = None, | |
| ) -> KnowledgeExtractionRunRow: | |
| """Open a run row before extraction begins. | |
| Opened first, not written at the end: a run that crashes mid-way must | |
| still be visible, otherwise the only trace of an expensive failed | |
| extraction is its absence. | |
| """ | |
| row = KnowledgeExtractionRunRow( | |
| id=str(uuid.uuid4()), | |
| knowledge_document_id=knowledge_document_id, | |
| doc_id=doc_id, | |
| scope_id=scope_id, | |
| user_id=user_id, | |
| status="running", | |
| branches=branches, | |
| model_deployment=model_deployment, | |
| ) | |
| async with AsyncSessionLocal() as session: | |
| session.add(row) | |
| await session.commit() | |
| await session.refresh(row) | |
| return row | |
| async def finish_run( | |
| self, | |
| run_id: str, | |
| *, | |
| status: str, | |
| usages: list[dict] | None = None, | |
| rejected: list[dict] | None = None, | |
| counts: dict | None = None, | |
| error_message: str | None = None, | |
| ) -> None: | |
| """Close a run with its usage, rejections and counts. | |
| `model_version` and `system_fingerprint` come from the API response, not | |
| from config: the deployment name stays stable while the model behind it | |
| moves, and a run without them recorded is a measurement of an unknown. | |
| """ | |
| usages = usages or [] | |
| async with AsyncSessionLocal() as session: | |
| row = await session.get(KnowledgeExtractionRunRow, run_id) | |
| if row is None: | |
| raise LookupError(f"no extraction run {run_id}") | |
| row.status = status | |
| row.usage = usages | |
| row.rejected = rejected or [] | |
| row.counts = counts or {} | |
| row.error_message = error_message | |
| row.n_calls = len(usages) | |
| row.prompt_tokens = sum(u.get("prompt_tokens", 0) or 0 for u in usages) | |
| row.cached_tokens = sum(u.get("cached_tokens", 0) or 0 for u in usages) | |
| row.completion_tokens = sum( | |
| u.get("completion_tokens", 0) or 0 for u in usages | |
| ) | |
| # Whatever actually served the calls. Distinct values would mean one | |
| # run spanned a model change, which is worth seeing rather than | |
| # averaging away, so the first non-null is recorded and the set is | |
| # visible in `usage`. | |
| for key in ("model_version", "system_fingerprint"): | |
| value = next((u.get(key) for u in usages if u.get(key)), None) | |
| if value: | |
| setattr(row, key, value) | |
| row.completed_at = datetime.now(UTC) | |
| await session.commit() | |
| async def save_entries( | |
| self, | |
| *, | |
| run_id: str, | |
| knowledge_document_id: str, | |
| doc_id: str, | |
| scope_id: str, | |
| user_id: str, | |
| entries_by_kind: dict[str, list[dict]], | |
| ) -> int: | |
| """Write every candidate entry for one run. Returns the row count. | |
| One flush for the whole run: a half-written entry set is not a useful | |
| thing to leave behind, and the caller settles the job on the exception. | |
| """ | |
| rows: list[KnowledgeEntryRow] = [] | |
| # `(run_id, kind, entity_id)` is UNIQUE, and ids are content-derived, so | |
| # two entries CAN collide inside one run: `formula_id` falls back | |
| # name -> latex -> chunk_id, and two formulas sharing a name share an id. | |
| # Writing both is impossible; dropping one silently is worse than saying | |
| # so, hence the warning and the count on the run row. | |
| seen: set[tuple[str, str]] = set() | |
| collisions: list[str] = [] | |
| # m1/m6: `label` carries the readable `TYPE_SUBJECT` key. Rules key off | |
| # the term they govern, and a rule payload holds `term_ids` rather than | |
| # words, so the surfaces are collected from this run's glossary first. | |
| surfaces = keys.term_surfaces(entries_by_kind.get("glossary")) | |
| for kind, payloads in entries_by_kind.items(): | |
| if kind not in _ENTRY_KINDS: | |
| raise ValueError(f"unknown entry kind {kind!r}") | |
| id_field, label_field = _ENTRY_KINDS[kind] | |
| for payload in payloads: | |
| entity_id = payload.get(id_field) | |
| if not entity_id: | |
| # Not recoverable by guessing: the id IS the thing an | |
| # expert's approval hangs on, and inventing one would make | |
| # the row unreviewable in a way nothing downstream detects. | |
| raise ValueError(f"{kind} entry has no {id_field}") | |
| if (kind, entity_id) in seen: | |
| collisions.append(f"{kind}:{entity_id}") | |
| continue | |
| seen.add((kind, entity_id)) | |
| prov = _provenance(payload) | |
| rows.append( | |
| KnowledgeEntryRow( | |
| id=str(uuid.uuid4()), | |
| run_id=run_id, | |
| knowledge_document_id=knowledge_document_id, | |
| doc_id=doc_id, | |
| scope_id=scope_id, | |
| user_id=user_id, | |
| kind=kind, | |
| entity_id=entity_id, | |
| # m6 (2026-09-08): the key, not the term. The readable | |
| # term is still in `payload`, so nothing is lost to a | |
| # reader — but the column an agent searches now holds | |
| # something it can type from memory. `_label` remains | |
| # the fallback for a kind with no registered prefix. | |
| label=keys.concept_key( | |
| kind, | |
| payload, | |
| entity_id=entity_id, | |
| doc_id=doc_id, | |
| surfaces=surfaces, | |
| ) | |
| or _label(payload, label_field), | |
| mention_count=int(payload.get("mention_count") or 0), | |
| extraction_status=payload.get("extraction_status"), | |
| diff_status=payload.get("diff_status"), | |
| definition_conflict=bool( | |
| payload.get("definition_conflict", False) | |
| ), | |
| payload=payload, | |
| # m3/m4/m5: promoted from the SAME object in the SAME | |
| # insert, so the column and the jsonb can never | |
| # disagree about what this entry says. | |
| content=_content(kind, payload), | |
| provenance=prov, | |
| doc_ids=_doc_ids(doc_id, prov), | |
| ) | |
| ) | |
| if not rows: | |
| return 0 | |
| async with AsyncSessionLocal() as session: | |
| session.add_all(rows) | |
| await session.commit() | |
| if collisions: | |
| logger.warning( | |
| "knowledge_entry_id_collision", | |
| extra={"run_id": run_id, "doc_id": doc_id, | |
| "dropped": len(collisions), "ids": collisions[:10]}, | |
| ) | |
| logger.info( | |
| "knowledge_entries_saved", | |
| extra={"run_id": run_id, "doc_id": doc_id, "count": len(rows), | |
| "collisions": len(collisions)}, | |
| ) | |
| self.last_collisions = collisions | |
| return len(rows) | |
| async def list_entries( | |
| self, | |
| *, | |
| scope_id: str, | |
| doc_id: str | None = None, | |
| kind: str | None = None, | |
| run_id: str | None = None, | |
| limit: int = 500, | |
| ) -> list[KnowledgeEntryRow]: | |
| """Read candidates back. Always filtered by tenant. | |
| `scope_id` is not optional: every read filters on it, so a missing | |
| filter is a cross-tenant read rather than a broader result set. | |
| """ | |
| stmt = select(KnowledgeEntryRow).where(KnowledgeEntryRow.scope_id == scope_id) | |
| if doc_id: | |
| stmt = stmt.where(KnowledgeEntryRow.doc_id == doc_id) | |
| if kind: | |
| stmt = stmt.where(KnowledgeEntryRow.kind == kind) | |
| if run_id: | |
| stmt = stmt.where(KnowledgeEntryRow.run_id == run_id) | |
| stmt = stmt.order_by( | |
| KnowledgeEntryRow.definition_conflict.desc(), | |
| KnowledgeEntryRow.mention_count.desc(), | |
| ).limit(limit) | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute(stmt) | |
| return list(result.scalars().all()) | |
| async def get_by_key( | |
| self, scope_id: str, key: str | |
| ) -> KnowledgeEntryRow | None: | |
| """One entry by its readable key (m8) — `GLOSSARY_PA`, not a search. | |
| This is the lookup the 2026-09-08 discussion asked for: an agent that | |
| meets a term it does not know writes the key and hits the row. **No | |
| model, no embedding, no scan of payloads** — the whole reason the key | |
| exists is that finding a glossary term should not cost a round trip to | |
| a language model. | |
| Keys are uppercased on the way in so a caller that types one from | |
| memory in lower case still lands, matching how they are minted. | |
| Note for whoever adds the index: this filters on `label`, which has no | |
| index today, so Postgres scans within the tenant. Correct, and fine at | |
| the current corpus size — but it is an index away from being the O(1) | |
| the design intends, and that index rides along with the m3/m4/m5 DDL. | |
| """ | |
| stmt = ( | |
| select(KnowledgeEntryRow) | |
| .where(KnowledgeEntryRow.scope_id == scope_id) | |
| .where(KnowledgeEntryRow.label == (key or "").strip().upper()) | |
| .order_by(KnowledgeEntryRow.created_at.desc()) | |
| .limit(1) | |
| ) | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute(stmt) | |
| return result.scalars().first() | |
| async def glossary_index( | |
| self, scope_id: str, *, limit: int = 1000 | |
| ) -> dict[str, dict]: | |
| """The company's glossary as a dictionary keyed by concept key (m8). | |
| `{"GLOSSARY_PA": {...payload...}, "GLOSSARY_QTY": {...}}` — the shape | |
| Mas asked for when he said the glossary should be key-value rather than | |
| a list. Searching a list is O(N) and searching this is O(1); it is also | |
| the shape the JSON export writes, so the file on disk and the lookup in | |
| memory agree with each other rather than drifting. | |
| A row whose key could not be built is skipped rather than given a | |
| made-up one: a glossary you cannot trust the keys of is worse than a | |
| smaller glossary. | |
| """ | |
| rows = await self.list_entries( | |
| scope_id=scope_id, kind="glossary", limit=limit | |
| ) | |
| index: dict[str, dict] = {} | |
| for row in rows: | |
| if row.label and isinstance(row.payload, dict): | |
| index.setdefault(row.label, row.payload) | |
| return index | |
| async def review_diff( | |
| self, scope_id: str, entity_id: str, kind: str = "glossary" | |
| ) -> dict | None: | |
| """Before / after for one entry, for the review screen (r4). | |
| Returns None when there is nothing to compare — no decision on record, | |
| or the entry is gone. `changed` is empty when the approval still | |
| stands, which is the common case and the one worth making cheap: the | |
| screen can say "nothing to re-read" instead of showing two identical | |
| panes. | |
| `before` is what the expert had in front of them; for an `edited` | |
| decision that is their own corrected text, not the model's claim. | |
| """ | |
| reviews = await self.latest_reviews(scope_id, [entity_id]) | |
| review = reviews.get(entity_id) | |
| if review is None: | |
| return None | |
| rows = await self.list_entries(scope_id=scope_id, kind=kind, limit=1000) | |
| current = next( | |
| (r for r in sorted(rows, key=lambda r: r.created_at or 0, reverse=True) | |
| if r.entity_id == entity_id), | |
| None, | |
| ) | |
| if current is None: | |
| return None | |
| if review.decision == "edited" and review.edited_payload: | |
| before = review.edited_payload | |
| else: | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeEntryRow).where( | |
| KnowledgeEntryRow.id == review.entry_id | |
| ) | |
| ) | |
| reviewed_row = result.scalars().first() | |
| before = reviewed_row.payload if reviewed_row else None | |
| changed = mutation.changed_fields(kind, before, current.payload) | |
| return { | |
| "entity_id": entity_id, | |
| "kind": kind, | |
| "decision": review.decision, | |
| "decided_at": review.created_at, | |
| "changed": changed, | |
| "status": "waiting_approval" if changed else review.decision, | |
| "before": mutation.ruled_content(kind, before), | |
| "after": mutation.ruled_content(kind, current.payload), | |
| } | |
| async def withdraw_document(self, scope_id: str, doc_id: str) -> int: | |
| """Mark every version of a document withdrawn (i7). Returns how many. | |
| **Nothing is deleted**, here or downstream. The entries the document | |
| produced stay exactly where they are, and so do the reviews attached to | |
| them — an approval is an act someone took, and destroying the record of | |
| what it referred to would destroy the audit trail that makes the | |
| approval mean anything. | |
| What withdrawal *does* change is what the knowledge base will serve: a | |
| withdrawn document stops feeding the approved baseline, so the pipeline | |
| cannot keep answering from a standard the client has taken back. That | |
| was the concern behind the whole discussion — never serve knowledge the | |
| source no longer supports. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| update(KnowledgeDocumentRow) | |
| .where( | |
| KnowledgeDocumentRow.scope_id == scope_id, | |
| KnowledgeDocumentRow.doc_id == doc_id, | |
| KnowledgeDocumentRow.withdrawn_at.is_(None), | |
| ) | |
| .values(withdrawn_at=datetime.now(UTC)) | |
| ) | |
| await session.commit() | |
| count = result.rowcount or 0 | |
| logger.info( | |
| "knowledge_document_withdrawn", | |
| extra={"scope_id": scope_id, "doc_id": doc_id, "versions": count}, | |
| ) | |
| return count | |
| async def withdrawn_doc_ids(self, scope_id: str) -> set[str]: | |
| """Documents this company has taken back.""" | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow.doc_id).where( | |
| KnowledgeDocumentRow.scope_id == scope_id, | |
| KnowledgeDocumentRow.withdrawn_at.is_not(None), | |
| ) | |
| ) | |
| return {r for (r,) in result.all() if r} | |
| async def orphaned_entities(self, scope_id: str, kind: str = "glossary") -> set[str]: | |
| """Entries whose every source document has been withdrawn (i7). | |
| Surfaced for the expert to rule on rather than deleted: whether the | |
| knowledge leaves with the document is their judgement, not the | |
| pipeline's. A term defined in a withdrawn standard may well still be | |
| the company's term. | |
| """ | |
| withdrawn = await self.withdrawn_doc_ids(scope_id) | |
| if not withdrawn: | |
| return set() | |
| rows = await self.list_entries(scope_id=scope_id, kind=kind, limit=1000) | |
| return {r.entity_id for r in rows if not _has_live_source(r, withdrawn)} | |
| async def impacted_by_document( | |
| self, scope_id: str, doc_id: str, limit: int = 1000 | |
| ) -> list[KnowledgeEntryRow]: | |
| """Which entries a document contributed to — the impact map (i4). | |
| This is what makes a re-ingest proportionate: a document is amended, | |
| these are the entries that can possibly have moved, and **only these** | |
| go back to the reviewer. Everything else the company has approved is | |
| untouched and stays untouched. | |
| No separate lineage structure is needed for it — the contribution is | |
| already recorded on the entry. Once `doc_ids` lands (m3) this also has | |
| to match membership in that list, for entries fused from several | |
| documents; today every entry has exactly one source, so the single | |
| column is the whole answer. | |
| """ | |
| return await self.list_entries(scope_id=scope_id, doc_id=doc_id, limit=limit) | |
| async def latest_run_for( | |
| self, scope_id: str, doc_id: str | |
| ) -> KnowledgeExtractionRunRow | None: | |
| """The most recent extraction of one document. | |
| Reads default to this rather than unioning every historical run's | |
| entries, which is never what a caller wants. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeExtractionRunRow) | |
| .where( | |
| KnowledgeExtractionRunRow.scope_id == scope_id, | |
| KnowledgeExtractionRunRow.doc_id == doc_id, | |
| ) | |
| .order_by(KnowledgeExtractionRunRow.created_at.desc()) | |
| .limit(1) | |
| ) | |
| return result.scalars().first() | |
| async def list_documents( | |
| self, scope_id: str, limit: int = 50 | |
| ) -> list[KnowledgeDocumentRow]: | |
| """Stored artifacts for one tenant, newest first. | |
| What the re-extract route is picked from: prompt retuning runs against a | |
| stored artifact and must never re-pay the parser for a byte-identical | |
| result. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow) | |
| .where(KnowledgeDocumentRow.scope_id == scope_id) | |
| .order_by(KnowledgeDocumentRow.created_at.desc()) | |
| .limit(limit) | |
| ) | |
| return list(result.scalars().all()) | |
| async def get_document( | |
| self, knowledge_document_id: str, scope_id: str | |
| ) -> KnowledgeDocumentRow | None: | |
| """One artifact, tenant-checked. | |
| `scope_id` is a predicate rather than a post-hoc assertion: filtering in | |
| SQL means a wrong tenant gets "not found" instead of a row it then has | |
| to be trusted not to use. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeDocumentRow).where( | |
| KnowledgeDocumentRow.id == knowledge_document_id, | |
| KnowledgeDocumentRow.scope_id == scope_id, | |
| ) | |
| ) | |
| return result.scalars().first() | |
| async def find_entry( | |
| self, entity_id: str, scope_id: str | |
| ) -> KnowledgeEntryRow | None: | |
| """The most recent entry row for one durable entity id. | |
| A decision is recorded against `entity_id`, which survives re-extraction; | |
| `entry_id` does not. So a review posted after a re-run must attach to the | |
| NEWEST row for that entity, which is what this returns. | |
| """ | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeEntryRow) | |
| .where( | |
| KnowledgeEntryRow.entity_id == entity_id, | |
| KnowledgeEntryRow.scope_id == scope_id, | |
| ) | |
| .order_by(KnowledgeEntryRow.created_at.desc()) | |
| .limit(1) | |
| ) | |
| return result.scalars().first() | |
| async def save_review( | |
| self, | |
| *, | |
| entity_id: str, | |
| scope_id: str, | |
| decision: str, | |
| reviewer_id: str, | |
| edited_payload: dict | None = None, | |
| note: str | None = None, | |
| decided_by: str = "human", | |
| authorised_by: str | None = None, | |
| ) -> KnowledgeReviewRow: | |
| """Record one expert decision. APPEND-ONLY. | |
| A changed mind is a NEW row, never an update: the point of the table is | |
| that it records what was ruled and when, and an update destroys the | |
| history the whole span-verification apparatus exists to produce. | |
| `content_hash` is derived server-side from the artifact the entry came | |
| from, never accepted from the caller — a stale client would otherwise | |
| attribute a decision to the wrong artifact version, which is exactly the | |
| confusion the field exists to prevent. | |
| """ | |
| entry = await self.find_entry(entity_id, scope_id) | |
| if entry is None: | |
| raise LookupError(f"no entry {entity_id} in this scope") | |
| async with AsyncSessionLocal() as session: | |
| document = await session.get( | |
| KnowledgeDocumentRow, entry.knowledge_document_id | |
| ) | |
| content_hash = document.content_hash if document else "" | |
| row = KnowledgeReviewRow( | |
| id=str(uuid.uuid4()), | |
| entry_id=entry.id, | |
| entity_id=entity_id, | |
| doc_id=entry.doc_id, | |
| scope_id=scope_id, | |
| content_hash=content_hash, | |
| decision=decision, | |
| edited_payload=edited_payload, | |
| note=note, | |
| reviewer_id=reviewer_id, | |
| decided_by=decided_by, | |
| authorised_by=authorised_by, | |
| ) | |
| session.add(row) | |
| await session.commit() | |
| await session.refresh(row) | |
| logger.info( | |
| "knowledge_review_saved", | |
| extra={ | |
| "entity_id": entity_id, | |
| "decision": decision, | |
| "decided_by": decided_by, | |
| "content_hash": content_hash[:12], | |
| }, | |
| ) | |
| return row | |
| async def auto_approve_run( | |
| self, *, run_id: str, scope_id: str, authorised_by: str | |
| ) -> int: | |
| """Approve every entry a run produced, on a named person's authority (r1). | |
| This is what `review_mode="llm"` does, and it is worth being plain about | |
| what it is: **no model re-examines anything.** The entries were written | |
| by a model already; choosing this mode means accepting them without a | |
| human looking. There is no second opinion in here. | |
| What makes it offerable rather than reckless is the record. Every row | |
| carries `decided_by="llm"` and `authorised_by=<the person who chose | |
| it>`, so an entry approved this way can always answer "who allowed | |
| this?" — which was the entire point when it was agreed: if something | |
| goes wrong it was their choice, made knowingly, and the system can show | |
| whose. | |
| Written through the ordinary append-only review path, so an expert can | |
| later disagree with any of it exactly as they would with each other. | |
| """ | |
| rows = await self.list_entries(scope_id=scope_id, run_id=run_id, limit=1000) | |
| approved = 0 | |
| for row in rows: | |
| try: | |
| await self.save_review( | |
| entity_id=row.entity_id, | |
| scope_id=scope_id, | |
| decision="approved", | |
| reviewer_id="llm", | |
| decided_by="llm", | |
| authorised_by=authorised_by, | |
| note=f"auto-approved (review_mode=llm) on behalf of {authorised_by}", | |
| ) | |
| approved += 1 | |
| except LookupError: | |
| # The entry vanished between listing and reviewing. Not fatal: | |
| # an un-approved entry is a reviewable one, which is the safe | |
| # direction to fail in. | |
| logger.warning( | |
| "knowledge_auto_approve_skipped", | |
| extra={"entity_id": row.entity_id, "run_id": run_id}, | |
| ) | |
| logger.info( | |
| "knowledge_run_auto_approved", | |
| extra={ | |
| "run_id": run_id, | |
| "scope_id": scope_id, | |
| "approved": approved, | |
| "authorised_by": authorised_by, | |
| }, | |
| ) | |
| return approved | |
| async def latest_reviews( | |
| self, scope_id: str, entity_ids: list[str] | None = None | |
| ) -> dict[str, KnowledgeReviewRow]: | |
| """Current decision per entity — the latest row wins. | |
| Returned as a map so a queue read can annotate rows in one pass rather | |
| than a query per row. | |
| """ | |
| stmt = select(KnowledgeReviewRow).where( | |
| KnowledgeReviewRow.scope_id == scope_id | |
| ) | |
| if entity_ids: | |
| stmt = stmt.where(KnowledgeReviewRow.entity_id.in_(entity_ids)) | |
| stmt = stmt.order_by(KnowledgeReviewRow.created_at.asc()) | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute(stmt) | |
| # Ascending, so a later row overwrites an earlier one and the map | |
| # ends up holding the current decision. | |
| return {r.entity_id: r for r in result.scalars().all()} | |
| async def approved_glossary(self, scope_id: str) -> list[dict]: | |
| """Approved glossary entries — the baseline the next run diffs against. | |
| THIS IS WHAT CLOSES THE LOOP: an expert's approvals become the active | |
| glossary a later extraction classifies new candidates against, which is | |
| what turns a batch extractor into a knowledge base that improves as it | |
| is worked. An `edited` decision contributes the expert's corrected | |
| payload, not the model's original claim. | |
| """ | |
| fresh, _stale = await self._approval_state(scope_id, "glossary") | |
| return fresh | |
| async def stale_approvals( | |
| self, scope_id: str, kind: str = "glossary" | |
| ) -> set[str]: | |
| """Entity ids whose approval no longer stands — 'termutasi' (i3). | |
| A later document changed the content the expert ruled on, so the entry | |
| is back to waiting approval and must not feed the baseline in the | |
| meantime. What counts as changed is `mutation.RULED_FIELDS`, and it is | |
| deliberately narrow — see that module for why a re-mention must not | |
| re-open anything. | |
| """ | |
| _fresh, stale = await self._approval_state(scope_id, kind) | |
| return stale | |
| async def _approval_state( | |
| self, scope_id: str, kind: str | |
| ) -> tuple[list[dict], set[str]]: | |
| """Split approved entries into those that still stand and those that do not. | |
| One pass, because the baseline and the status display must never | |
| disagree about whether an approval is current — the same defect as | |
| having two notions of readiness, which is why `readiness.py` shares one | |
| predicate between the report API and Help. | |
| """ | |
| reviews = await self.latest_reviews(scope_id) | |
| approved = { | |
| entity_id: review | |
| for entity_id, review in reviews.items() | |
| if review.decision in ("approved", "edited") | |
| } | |
| if not approved: | |
| return [], set() | |
| reviewed_ids = [r.entry_id for r in approved.values() if r.entry_id] | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeEntryRow).where( | |
| KnowledgeEntryRow.scope_id == scope_id, | |
| KnowledgeEntryRow.kind == kind, | |
| KnowledgeEntryRow.entity_id.in_(list(approved)), | |
| ) | |
| ) | |
| rows = list(result.scalars().all()) | |
| # The exact rows the decisions were made against. Without these the | |
| # only available comparison is document-level `content_hash`, which | |
| # flips every entry in an edited document including the untouched | |
| # ones — the outcome the i3 ruling rejected. | |
| reviewed_rows: dict[str, KnowledgeEntryRow] = {} | |
| if reviewed_ids: | |
| reviewed = await session.execute( | |
| select(KnowledgeEntryRow).where( | |
| KnowledgeEntryRow.id.in_(reviewed_ids) | |
| ) | |
| ) | |
| reviewed_rows = {r.id: r for r in reviewed.scalars().all()} | |
| withdrawn = await self.withdrawn_doc_ids(scope_id) | |
| out: list[dict] = [] | |
| stale: set[str] = set() | |
| seen: set[str] = set() | |
| for row in sorted(rows, key=lambda r: r.created_at or 0, reverse=True): | |
| if row.entity_id in seen: | |
| continue | |
| seen.add(row.entity_id) | |
| # i7: a withdrawn source stops feeding the baseline, approval or | |
| # not. The approval itself is untouched and the rows stay — this | |
| # only stops the pipeline answering from a document the client has | |
| # taken back. | |
| if withdrawn and not _has_live_source(row, withdrawn): | |
| continue | |
| review = approved[row.entity_id] | |
| if review.decision == "edited" and review.edited_payload: | |
| # An edit is a stronger statement than an approval: the expert | |
| # wrote the text themselves, so a later run producing different | |
| # model output does not unsettle it. Their payload IS the entry. | |
| out.append(review.edited_payload) | |
| continue | |
| ruled = reviewed_rows.get(review.entry_id) | |
| if ruled is not None and mutation.has_mutated( | |
| kind, ruled.payload, row.payload | |
| ): | |
| stale.add(row.entity_id) | |
| continue | |
| out.append(row.payload) | |
| if stale: | |
| logger.info( | |
| "knowledge_approvals_stale", | |
| extra={"scope_id": scope_id, "kind": kind, "count": len(stale)}, | |
| ) | |
| return out, stale | |
| async def save_declaration( | |
| self, *, scope_id: str, user_id: str, payload: dict | |
| ) -> KnowledgeEntryRow: | |
| """Persist a scope-level domain declaration as a `kind='domain'` entry. | |
| Deliberately an ENTRY rather than its own table: it has to be reviewable, | |
| and `knowledge_reviews` keys on `entity_id`. Making it an entry means the | |
| expert approves, edits or rejects a domain declaration through exactly | |
| the same route as a glossary term - no second review mechanism to build, | |
| and no second one to keep consistent. | |
| Parentless by nature: it is composed across every document in the scope, | |
| so it belongs to no run and no document. A DB CHECK allows that only for | |
| this kind. | |
| Re-proposing REPLACES the standing draft rather than accumulating | |
| versions, because a draft is not a finding. The expert's decisions live | |
| in `knowledge_reviews` and are keyed on `entity_id`, so they survive the | |
| replacement - which is the whole reason the id is deterministic. | |
| """ | |
| entity_id = payload["domain_id"] | |
| async with AsyncSessionLocal() as session: | |
| existing = await session.execute( | |
| select(KnowledgeEntryRow).where( | |
| KnowledgeEntryRow.scope_id == scope_id, | |
| KnowledgeEntryRow.kind == "domain", | |
| KnowledgeEntryRow.entity_id == entity_id, | |
| ) | |
| ) | |
| row = existing.scalars().first() | |
| if row is None: | |
| row = KnowledgeEntryRow( | |
| id=str(uuid.uuid4()), | |
| run_id=None, | |
| knowledge_document_id=None, | |
| doc_id=scope_id, # scope-level; there is no document | |
| scope_id=scope_id, | |
| user_id=user_id, | |
| kind="domain", | |
| entity_id=entity_id, | |
| ) | |
| session.add(row) | |
| # m6: the declaration is keyed like every other entry, so one | |
| # lookup path covers all five kinds. | |
| row.label = keys.concept_key( | |
| "domain", payload, entity_id=entity_id | |
| ) or _label(payload, "domain_name") | |
| row.payload = payload | |
| row.extraction_status = "ok" | |
| await session.commit() | |
| await session.refresh(row) | |
| logger.info( | |
| "knowledge_declaration_saved", | |
| extra={"scope_id": scope_id, "entity_id": entity_id}, | |
| ) | |
| return row | |
| async def get_declaration(self, scope_id: str) -> KnowledgeEntryRow | None: | |
| """The standing domain declaration for a scope, approved or not.""" | |
| async with AsyncSessionLocal() as session: | |
| result = await session.execute( | |
| select(KnowledgeEntryRow).where( | |
| KnowledgeEntryRow.scope_id == scope_id, | |
| KnowledgeEntryRow.kind == "domain", | |
| ) | |
| ) | |
| return result.scalars().first() | |