"""Compose a `DomainContext` from a company's persisted entries (V4-5..V4-7). Reads `knowledge_entries` and joins glossary -> formula -> rules into measures. **Nothing here calls a model.** Every field this module produces already exists in the database; composition is a join, which is what makes it verifiable and free to run. The declared half of the context (identity, boundary, conventions, entities, dimensions) is NOT produced here - it is authored through the review loop (decided 2026-09-08: the model proposes, the expert approves). This module leaves those fields empty and reports the gap in `coverage.declared`. """ from __future__ import annotations from datetime import UTC, datetime from typing import Any from src.knowledge_domain.models import ( Authority, Conventions, Coverage, Definition, DomainContext, Identity, Measure, SourceRef, Unit, ) from src.knowledge_extraction.cluster.normalize import AbbrevIndex, normalize from src.knowledge_ingest.knowledge_store import KnowledgeStore from src.middlewares.logging import get_logger logger = get_logger("knowledge_domain") # A measure with no definition AND no formula is a bare surface - the filter saw # a term and nothing was learned about it. Those belong in the review queue, not # in a planner's context, where they add tokens and no information. _MIN_USEFUL = ("definition", "formula_id") def _surface_forms(payload: dict[str, Any]) -> list[str]: """Every way the documents write this term, deduped, order preserved. `source_wording` is included precisely when it DISAGREES with the canonical name - that disagreement is a locked product decision and a planner meeting either form in a user's question should recognise it. """ forms: list[str] = [] for key in ("term", "full_name", "source_wording"): value = payload.get(key) if isinstance(value, str) and value.strip() and value.strip() not in forms: forms.append(value.strip()) return forms class DomainComposer: """Builds one company's domain context out of its approved entries.""" def __init__(self, store: KnowledgeStore | None = None) -> None: self._store = store or KnowledgeStore() async def compose( self, scope_id: str, *, include_unapproved: bool = False, catalog: Any | None = None, ) -> DomainContext: """Compose the context for one company. ⚠️ **`include_unapproved` is a DEVELOPMENT flag and must never be True on the request path.** Candidates are the model's claims; letting them into a planner's context means the pipeline's guesses steer the analysis, which is exactly what the review step exists to prevent. It exists because `knowledge_reviews` is empty today and a feature nobody can see working is hard to build against (decided 2026-09-08). Every entry carries `approved`, so a consumer can always tell them apart, and turning it on logs loudly. """ if include_unapproved: logger.warning( "domain_context_including_unapproved", extra={"scope_id": scope_id, "reason": "development flag"}, ) glossary = await self._store.list_entries( scope_id=scope_id, kind="glossary", limit=1000 ) formulas = await self._store.list_entries( scope_id=scope_id, kind="formula", limit=1000 ) rules = await self._store.list_entries( scope_id=scope_id, kind="rule", limit=1000 ) # The per-document brief. Since the V4-4 split this is kind="document"; # kind="domain" now means something else entirely and must NOT be read # here. Reading both was a transitional line from before the split, and # leaving it would have made the first declared domain entry ever # written get folded into `authority.sources[]` as if it were a file. briefs = await self._store.list_entries( scope_id=scope_id, kind="document", limit=200 ) # The DECLARED half - identity, boundary, conventions, entities, # dimensions. Expert-authored (the model proposes, the expert approves), # so it must be an entry to be reviewable: `knowledge_reviews` keys on # `entity_id`. Nothing writes one yet (V4-9/V4-10), so this is normally # empty - but the seam is here so that when one lands it is applied # instead of silently ignored. declared_rows = await self._store.list_entries( scope_id=scope_id, kind="domain", limit=5 ) reviews = await self._store.latest_reviews(scope_id) approved_ids = { entity_id for entity_id, review in reviews.items() if review.decision in ("approved", "edited") } def _payload(row): """The expert's correction wins over the model's original claim.""" review = reviews.get(row.entity_id) if review is not None and review.decision == "edited" and review.edited_payload: return review.edited_payload return row.payload def _usable(rows): return [ r for r in rows if include_unapproved or r.entity_id in approved_ids ] formula_by_id = {r.entity_id: _payload(r) for r in _usable(formulas)} rule_rows = _usable(rules) # Rules that name a term govern that measure; rules that name none are # domain-wide. Splitting here means nothing is listed twice. governs: dict[str, list[str]] = {} policies: list[str] = [] for row in rule_rows: payload = _payload(row) term_ids = [t for t in (payload.get("term_ids") or []) if t] if term_ids: for term_id in term_ids: governs.setdefault(term_id, []).append(row.entity_id) else: policies.append(row.entity_id) # `term_id` is document-scoped by design - right for review, wrong for a # domain. Two documents defining PA would otherwise be two measures, the # same duplication fixed inside one document on 2026-09-08 reappearing # across documents. Merge on the canonical surface, using the SAME # normalisation that clusters mentions: one rule at both scales. index = AbbrevIndex([]) index.learn_inline( form for row in _usable(glossary) for form in _surface_forms(_payload(row)) ) merged: dict[str, Measure] = {} for row in _usable(glossary): payload = _payload(row) formula_id = payload.get("defining_formula_id") formula = formula_by_id.get(formula_id) if formula_id else None if not any(payload.get(k) for k in _MIN_USEFUL) and not formula: continue name = payload.get("term") or row.label or row.entity_id key = index.canonical_key(name) or normalize(name) or row.entity_id approved = row.entity_id in approved_ids definition = payload.get("definition") existing = merged.get(key) if existing is None: merged[key] = Measure( key=key, term_ids=[row.entity_id], name=name, surface_forms=_surface_forms(payload), definition=definition, definitions=( [Definition(text=definition, doc_id=row.doc_id, term_id=row.entity_id, approved=approved)] if definition else [] ), source_docs=[row.doc_id], unit=(formula or {}).get("unit"), formula_id=formula_id, formula_latex=(formula or {}).get("formula_latex"), inputs=[ v.get("symbol") for v in ((formula or {}).get("variables") or []) if v.get("symbol") ], governed_by=list(governs.get(row.entity_id, [])), catalog_binding=None, mention_count=row.mention_count or 0, approved=approved, ) continue # Second document, same measure. Union everything; never overwrite. existing.term_ids.append(row.entity_id) if row.doc_id not in existing.source_docs: existing.source_docs.append(row.doc_id) for form in _surface_forms(payload): if form not in existing.surface_forms: existing.surface_forms.append(form) existing.governed_by.extend( r for r in governs.get(row.entity_id, []) if r not in existing.governed_by ) existing.mention_count += row.mention_count or 0 existing.approved = existing.approved or approved # Prefer the shortest name - "PA" reads better in a card than # "Physical of Availability (PA)". if len(name) < len(existing.name): existing.name = name existing.formula_id = existing.formula_id or formula_id existing.formula_latex = existing.formula_latex or (formula or {}).get( "formula_latex" ) existing.unit = existing.unit or (formula or {}).get("unit") if definition: existing.definitions.append( Definition(text=definition, doc_id=row.doc_id, term_id=row.entity_id, approved=approved) ) # A DISAGREEMENT is the finding, not a problem to resolve. Two # documents saying different things about one measure is exactly # what an expert must rule on, so both are kept and the measure # is marked rather than one silently winning. distinct = {d.text.strip() for d in existing.definitions} existing.contested = len(distinct) > 1 if existing.definition is None: existing.definition = definition measures = list(merged.values()) # Most-established first: a term the corpus keeps returning to is the one # a reader should meet first, and it is the sensible truncation order # when the card has to be capped. measures.sort(key=lambda m: (-m.mention_count, m.name.casefold())) if catalog is not None: self._bind_to_catalog(measures, catalog) units: list[Unit] = [] seen_units: set[str] = set() for measure in measures: if measure.unit and measure.unit not in seen_units: seen_units.add(measure.unit) units.append(Unit(symbol=measure.unit)) subdomains: list[str] = [] sources: list[SourceRef] = [] for row in briefs: payload = _payload(row) for tag in payload.get("subdomains") or []: if tag not in subdomains: subdomains.append(tag) sources.append( SourceRef( doc_id=row.doc_id, title=payload.get("title"), n_pages=0, approved_at=getattr(reviews.get(row.entity_id), "created_at", None), ) ) # Apply the declared half when an expert has ruled on one. declared = next( (_payload(r) for r in declared_rows if include_unapproved or r.entity_id in approved_ids), None, ) identity = Identity(subdomains=subdomains) conventions = Conventions(units=units) if declared: identity = Identity( domain_name=declared.get("domain_name"), # Declared subdomains win over the tags rolled up from briefs: # an expert saying what the domain covers beats an aggregate of # per-term guesses. subdomains=list(declared.get("subdomains") or subdomains), boundary=declared.get("boundary"), boundary_note=declared.get("boundary_note"), ) conv = declared.get("conventions") or {} conventions = Conventions( time_grain=conv.get("time_grain"), calendar=conv.get("calendar"), units=units or [Unit(**u) for u in (conv.get("units") or [])], ) total = len(glossary) + len(formulas) + len(rules) + len(briefs) + len( declared_rows ) context = DomainContext( scope_id=scope_id, # Derived, overlaid with the declared half when an expert has # approved one. `domain_name` and `boundary` stay empty until then. identity=identity, measures=measures, conventions=conventions, policies=policies, authority=Authority( sources=sources, as_of=datetime.now(UTC), coverage=Coverage( n_documents=len(sources), n_measures=len(measures), n_policies=len(policies), n_entries_total=total, n_entries_approved=len(approved_ids), approved_ratio=( (len(approved_ids) / total) if (total and reviews) else None ), declared=declared is not None, ), ), ) logger.info( "domain_context_composed", extra={ "scope_id": scope_id, "measures": len(measures), "policies": len(policies), "approved": len(approved_ids), "unapproved_included": include_unapproved, }, ) return context @staticmethod def _bind_to_catalog(measures: list[Measure], catalog: Any) -> None: """Resolve each measure to a real `column_id` where one matches. The field that turns a glossary the planner reads into a semantic layer it computes against: a measure naming no column is documentation. Deliberately conservative and **self-disabling** - it matches a measure's surface forms against column names and binds only on an unambiguous hit. A wrong binding is worse than none: it would point the planner at real data under the wrong name, which is the force-mapping failure this whole model exists to prevent. Mirrors `fk_inference.py`'s posture. """ columns: dict[str, str] = {} for source in getattr(catalog, "sources", []) or []: for table in getattr(source, "tables", []) or []: for column in getattr(table, "columns", []) or []: name = (getattr(column, "name", "") or "").strip().casefold() column_id = getattr(column, "column_id", None) if not name or not column_id: continue # A name occurring in two tables is ambiguous; drop it rather # than pick one. columns[name] = "" if name in columns else column_id bound = 0 for measure in measures: for form in measure.surface_forms: column_id = columns.get(form.strip().casefold()) if column_id: measure.catalog_binding = column_id bound += 1 break logger.info( "domain_measures_bound", extra={"bound": bound, "measures": len(measures)}, )