"""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=`, 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()