Download src/knowledge_domain/compose.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 16 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_domain/compose.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/knowledge_domain/compose.py
-
curl -L -o compose.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_domain/compose.py
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 | |
| 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)}, | |
| ) | |