Files
Codex cffa60772e
Publish documentation / publish (push) Successful in 2m12s
feat: complete catalog-driven preprocessing
2026-09-06 17:49:35 +02:00

201 lines
7.4 KiB
Python

from __future__ import annotations
import re
from typing import TYPE_CHECKING
from pydantic import BaseModel
from tht.mschema.models import Annotations, PhysicalSchema
if TYPE_CHECKING:
from tht.evidence.model import EvidenceDoc
MAX_EXAMPLES_IN_RECORD = 5
class VectorRecord(BaseModel):
id: str
kind: str # evidence | schema_table | schema_column | schema_relationship
ref: str # file/chiave canonica di provenienza
title: str
content: str
metadata: dict = {}
def qdrant_semantic_kind(kind: str) -> str:
if kind in {"schema_table", "schema_column", "schema_relationship"}:
return "schema"
if kind in {"memory", "solved_question"}:
return "memory"
if kind == "evidence":
return "evidence"
raise ValueError(f"Unsupported vector kind: {kind}")
def qdrant_payload(
record: VectorRecord,
*,
content_hash: str,
workspace_id: str,
workspace_revision: str | None = None,
) -> dict:
semantic_kind = qdrant_semantic_kind(record.kind)
return {
**record.metadata,
"workspace_id": workspace_id,
**({"workspace_revision": workspace_revision} if workspace_revision else {}),
"kind": semantic_kind,
"record_kind": record.kind,
"record_key": record.id,
"ref": record.ref,
"title": record.title,
"content": record.content,
"content_hash": content_hash,
}
def split_markdown(text: str, max_chars: int) -> list[str]:
"""Spezza un markdown: intero se sta nel limite, altrimenti per heading '##',
e in ultima istanza per accumulo greedy di righe."""
if len(text) <= max_chars:
return [text]
parts = re.split(r"(?=^## )", text, flags=re.MULTILINE)
chunks: list[str] = []
for part in parts:
part = part.strip("\n")
if not part:
continue
if len(part) <= max_chars:
chunks.append(part)
continue
current: list[str] = []
size = 0
for line in part.splitlines():
if size + len(line) > max_chars and current:
chunks.append("\n".join(current))
current, size = [], 0
current.append(line)
size += len(line) + 1
if current:
chunks.append("\n".join(current))
return chunks
def evidence_records(docs: list[EvidenceDoc], max_chunk_chars: int) -> list[VectorRecord]:
"""Record per tutte le evidence presenti: la sola presenza basta a indicizzarle."""
records: list[VectorRecord] = []
for doc in docs:
content = f"{doc.title}\n\n{doc.body}"
for i, chunk in enumerate(split_markdown(content, max_chunk_chars)):
records.append(
VectorRecord(
id=f"evidence:{doc.id}:{i}",
kind="evidence",
ref=str(doc.path) if doc.path else doc.id,
title=doc.title,
content=chunk,
metadata={
"status": doc.status, "tier": doc.tier,
"tables": doc.tables, "concepts": doc.concepts,
},
)
)
return records
def schema_records(physical: PhysicalSchema, annotations: Annotations) -> list[VectorRecord]:
"""Un record per tabella e uno per colonna, da mschema (physical + annotations)."""
records: list[VectorRecord] = []
for table_name, table in physical.tables.items():
ann_t = annotations.tables.get(table_name)
t_desc = (ann_t.description if ann_t and ann_t.description else table.comment)
t_concepts = ann_t.concepts if ann_t else []
lines = [f"Tabella {table_name}", t_desc]
if t_concepts:
lines.append("Concetti: " + ", ".join(t_concepts))
lines.append("Colonne: " + ", ".join(table.columns))
records.append(
VectorRecord(
id=f"schema_table:{table_name}", kind="schema_table", ref=table_name,
title=table_name, content="\n".join(filter(None, lines)),
)
)
for column_name, column in table.columns.items():
ann_c = ann_t.columns.get(column_name) if ann_t else None
c_desc = (ann_c.description if ann_c and ann_c.description else column.comment)
lines = [f"Colonna {table_name}.{column_name} ({column.type})", c_desc]
if ann_c and ann_c.synonyms:
lines.append("Sinonimi: " + ", ".join(ann_c.synonyms))
if column.examples:
lines.append("Esempi: " + ", ".join(column.examples[:MAX_EXAMPLES_IN_RECORD]))
records.append(
VectorRecord(
id=f"schema_column:{table_name}.{column_name}", kind="schema_column",
ref=f"{table_name}.{column_name}", title=f"{table_name}.{column_name}",
content="\n".join(filter(None, lines)),
)
)
return records
def catalog_schema_records(snapshot) -> list[VectorRecord]:
"""Build the complete schema slice from a Catalog Metadata Snapshot."""
records: list[VectorRecord] = []
for table in snapshot.tables:
column_names = [column.name for column in table.columns]
lines = [f"Tabella {table.name}"]
if table.description:
lines.append(table.description)
lines.append("Colonne: " + ", ".join(column_names))
records.append(
VectorRecord(
id=f"schema_table:{table.id}",
kind="schema_table",
ref=table.name,
title=table.name,
content="\n".join(lines),
metadata={"tables": [table.name]},
)
)
for column in table.columns:
lines = [f"Colonna {table.name}.{column.name}", f"Tipo: {column.data_type}"]
if column.description:
lines.append(column.description)
if column.sensitive:
lines.append("Dato sensibile")
records.append(
VectorRecord(
id=f"schema_column:{column.id}",
kind="schema_column",
ref=f"{table.name}.{column.name}",
title=f"{table.name}.{column.name}",
content="\n".join(lines),
metadata={"tables": [table.name], "sensitive": column.sensitive},
)
)
for relationship in snapshot.relationships:
pairs = ", ".join(
f"{relationship.source_table}.{source} -> "
f"{relationship.target_table}.{target}"
for source, target in zip(
relationship.source_columns, relationship.target_columns, strict=True
)
)
ref = f"{relationship.source_table}->{relationship.target_table}"
records.append(
VectorRecord(
id=f"schema_relationship:{relationship.id}",
kind="schema_relationship",
ref=ref,
title=ref,
content=f"Relazione {relationship.origin}: {pairs}",
metadata={
"tables": [relationship.source_table, relationship.target_table],
"origin": relationship.origin,
"source_table": relationship.source_table,
"target_table": relationship.target_table,
},
)
)
return records