ishaq101's picture
/feat knowledge management (#20)
b68816f
Raw History Blame Contribute Delete
16 kB
"""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)},
)