feat: consolidate database management work
Add catalog-owned logical relationships and runtime snapshots, extend the database-management UI and validation coverage, and document the updated operational workflow. Keep active sensitive-generation status in a tooltip and indicator, and update the layout E2E to follow the history action in its new database-scoped location.
This commit is contained in:
@@ -0,0 +1,403 @@
|
||||
import json
|
||||
from datetime import UTC, datetime
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
from typer.testing import CliRunner
|
||||
|
||||
from tht.cli import app
|
||||
from tht.config import load_config
|
||||
from tht.mschema.context import SchemaContextError, load_schema_context
|
||||
from tht.mschema.models import (
|
||||
Annotations,
|
||||
ColumnPhysical,
|
||||
ForeignKey,
|
||||
PhysicalSchema,
|
||||
TableAnnotation,
|
||||
TablePhysical,
|
||||
)
|
||||
from tht.mschema.render import to_mschema_text
|
||||
|
||||
|
||||
REVISION = "a" * 40
|
||||
RUNNER = CliRunner()
|
||||
|
||||
|
||||
def _workspace(tmp_path, snapshot: dict) -> object:
|
||||
artifacts = tmp_path / "artifacts"
|
||||
relationships = tmp_path / "effective-relationships.json"
|
||||
relationships.write_text(json.dumps(snapshot))
|
||||
physical = PhysicalSchema(
|
||||
database="warehouse",
|
||||
schema="analytics",
|
||||
introspected_at=datetime(2026, 1, 1, tzinfo=UTC),
|
||||
tables={
|
||||
"users": TablePhysical(
|
||||
columns={"id": ColumnPhysical(type="bigint", pk=True)},
|
||||
),
|
||||
"orders": TablePhysical(
|
||||
columns={
|
||||
"physical_user_id": ColumnPhysical(type="bigint"),
|
||||
"annotated_user_id": ColumnPhysical(type="bigint"),
|
||||
"generated_user_id": ColumnPhysical(type="bigint"),
|
||||
},
|
||||
foreign_keys=[
|
||||
ForeignKey(
|
||||
columns=["physical_user_id"],
|
||||
ref_table="users",
|
||||
ref_columns=["id"],
|
||||
)
|
||||
],
|
||||
),
|
||||
},
|
||||
)
|
||||
physical.to_yaml(artifacts / "mschema" / "physical.yaml")
|
||||
Annotations(
|
||||
tables={
|
||||
"orders": TableAnnotation(
|
||||
description="Customer orders",
|
||||
foreign_keys=[
|
||||
ForeignKey(
|
||||
columns=["annotated_user_id"],
|
||||
ref_table="users",
|
||||
ref_columns=["id"],
|
||||
)
|
||||
],
|
||||
)
|
||||
}
|
||||
).to_yaml(artifacts / "mschema" / "annotations.yaml")
|
||||
config = tmp_path / "runtime.yaml"
|
||||
config.write_text(
|
||||
f"""
|
||||
runtime_identity:
|
||||
workspace_id: demo
|
||||
workspace_revision: {REVISION}
|
||||
source_identity: workspace://demo
|
||||
database: {{database: warehouse, schema: analytics, user: reader, password: secret, transport: direct}}
|
||||
vector_db: {{database: vectors, schema: public, user: reader, password: secret}}
|
||||
embeddings: {{base_url: 'http://localhost:11434', model: 'qwen3-embedding:0.6b', dim: 1024}}
|
||||
paths:
|
||||
artifacts: {artifacts}
|
||||
indexes: {tmp_path / 'indexes'}
|
||||
sessions: {tmp_path / 'sessions'}
|
||||
effective_relationships: {relationships}
|
||||
"""
|
||||
)
|
||||
return load_config(config)
|
||||
|
||||
|
||||
def test_effective_snapshot_is_the_exclusive_relationship_source(tmp_path):
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "orders",
|
||||
"sourceColumns": ["generated_user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "generated",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
context = load_schema_context(cfg)
|
||||
rendered = to_mschema_text(
|
||||
context.physical,
|
||||
context.annotations,
|
||||
effective_relationships=context.effective_relationships,
|
||||
)
|
||||
|
||||
assert "-- Customer orders" in rendered
|
||||
assert "orders.generated_user_id=users.id" in rendered
|
||||
assert "orders.physical_user_id=users.id" not in rendered
|
||||
assert "orders.annotated_user_id=users.id" not in rendered
|
||||
|
||||
|
||||
def test_declared_effective_snapshot_must_exist(tmp_path):
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{"schemaVersion": 1, "workspaceId": "demo", "relationships": []},
|
||||
)
|
||||
cfg.paths.effective_relationships.unlink()
|
||||
|
||||
with pytest.raises(SchemaContextError, match="snapshot is missing"):
|
||||
load_schema_context(cfg)
|
||||
|
||||
|
||||
def test_data_root_override_preserves_the_runtime_relationship_snapshot_path(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
monkeypatch.setenv("THT_DATA_ROOT", str(tmp_path / "mounted-data"))
|
||||
monkeypatch.delenv("THT_HOME", raising=False)
|
||||
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{"schemaVersion": 1, "workspaceId": "demo", "relationships": []},
|
||||
)
|
||||
|
||||
assert cfg.paths.effective_relationships == tmp_path / "effective-relationships.json"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("snapshot", "message"),
|
||||
[
|
||||
(
|
||||
{"schemaVersion": 1, "workspaceId": "another", "relationships": []},
|
||||
"workspace does not match",
|
||||
),
|
||||
(
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "missing_orders",
|
||||
"sourceColumns": ["user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "generated",
|
||||
}
|
||||
],
|
||||
},
|
||||
"endpoint table is absent",
|
||||
),
|
||||
(
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "orders",
|
||||
"sourceColumns": ["missing_user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "manual",
|
||||
}
|
||||
],
|
||||
},
|
||||
"endpoint column is absent",
|
||||
),
|
||||
(
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "orders",
|
||||
"sourceColumns": ["generated_user_id", "annotated_user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "physical",
|
||||
}
|
||||
],
|
||||
},
|
||||
"snapshot is invalid",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_effective_snapshot_rejects_wrong_workspace_and_orphan_endpoints(
|
||||
tmp_path, snapshot, message
|
||||
):
|
||||
cfg = _workspace(tmp_path, snapshot)
|
||||
|
||||
with pytest.raises(SchemaContextError, match=message):
|
||||
load_schema_context(cfg)
|
||||
|
||||
|
||||
def test_legacy_relationship_merge_remains_available_and_deduplicated(tmp_path):
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{"schemaVersion": 1, "workspaceId": "demo", "relationships": []},
|
||||
)
|
||||
cfg.paths.effective_relationships = None
|
||||
annotations = Annotations.from_yaml(
|
||||
cfg.paths.artifacts / "mschema" / "annotations.yaml"
|
||||
)
|
||||
annotations.tables["orders"].foreign_keys.append(
|
||||
ForeignKey(
|
||||
columns=["annotated_user_id"],
|
||||
ref_table="users",
|
||||
ref_columns=["id"],
|
||||
)
|
||||
)
|
||||
annotations.to_yaml(cfg.paths.artifacts / "mschema" / "annotations.yaml")
|
||||
|
||||
context = load_schema_context(cfg)
|
||||
rendered = to_mschema_text(context.physical, context.annotations)
|
||||
|
||||
assert context.effective_relationships is None
|
||||
assert rendered.count("orders.physical_user_id=users.id") == 1
|
||||
assert rendered.count("orders.annotated_user_id=users.id") == 1
|
||||
|
||||
|
||||
def test_schema_render_uses_the_effective_snapshot(tmp_path):
|
||||
_workspace(
|
||||
tmp_path,
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "orders",
|
||||
"sourceColumns": ["generated_user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "manual",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
result = RUNNER.invoke(
|
||||
app,
|
||||
["schema", "render", "--format", "mschema-text", "-c", str(tmp_path / "runtime.yaml")],
|
||||
)
|
||||
|
||||
assert result.exit_code == 0, result.output
|
||||
assert "orders.generated_user_id=users.id" in result.stdout
|
||||
assert "orders.physical_user_id=users.id" not in result.stdout
|
||||
assert "orders.annotated_user_id=users.id" not in result.stdout
|
||||
|
||||
|
||||
def test_schema_search_uses_the_effective_snapshot(tmp_path, monkeypatch):
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{
|
||||
"schemaVersion": 1,
|
||||
"workspaceId": "demo",
|
||||
"relationships": [
|
||||
{
|
||||
"sourceTable": "orders",
|
||||
"sourceColumns": ["generated_user_id"],
|
||||
"targetTable": "users",
|
||||
"targetColumns": ["id"],
|
||||
"origin": "generated",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
physical = cfg.paths.artifacts / "mschema" / "physical.yaml"
|
||||
|
||||
class Searcher:
|
||||
def search(self, _vector, **_kwargs):
|
||||
return [
|
||||
SimpleNamespace(
|
||||
id="orders",
|
||||
kind="schema_table",
|
||||
ref="orders",
|
||||
title="orders",
|
||||
similarity=1.0,
|
||||
content="orders",
|
||||
metadata={},
|
||||
)
|
||||
]
|
||||
|
||||
class Embedder:
|
||||
def embed_query(self, _text):
|
||||
return [1.0, 0.0, 0.0]
|
||||
|
||||
import tht.cli.search_cmd as search_module
|
||||
import tht.cli.vector_cmd as vector_module
|
||||
import tht.evidence as evidence_module
|
||||
|
||||
monkeypatch.setattr(
|
||||
search_module,
|
||||
"_leased_dwh_snapshot",
|
||||
lambda _cfg, _ctx: SimpleNamespace(physical=physical, lsh_dir=tmp_path / "no-lsh"),
|
||||
)
|
||||
monkeypatch.setattr(vector_module, "require_vector_cfg", lambda _cfg: None)
|
||||
monkeypatch.setattr(vector_module, "open_searcher", lambda _cfg: Searcher())
|
||||
monkeypatch.setattr(vector_module, "make_embedder", lambda _cfg: Embedder())
|
||||
monkeypatch.setattr(evidence_module, "validate_corpus_workspace", lambda *_args: None)
|
||||
monkeypatch.setattr(evidence_module, "active_searcher", lambda _cfg, searcher, **_kwargs: searcher)
|
||||
|
||||
result = RUNNER.invoke(
|
||||
app,
|
||||
[
|
||||
"search",
|
||||
"find",
|
||||
"orders",
|
||||
"--kind",
|
||||
"schema",
|
||||
"--json",
|
||||
"-c",
|
||||
str(tmp_path / "runtime.yaml"),
|
||||
],
|
||||
)
|
||||
|
||||
assert result.exit_code == 0, result.output
|
||||
payload = json.loads(result.stdout)
|
||||
assert "orders.generated_user_id=users.id" in payload["mschema"]
|
||||
assert "orders.physical_user_id=users.id" not in payload["mschema"]
|
||||
assert "orders.annotated_user_id=users.id" not in payload["mschema"]
|
||||
|
||||
|
||||
def test_schema_search_fails_closed_before_returning_no_candidates(tmp_path, monkeypatch):
|
||||
cfg = _workspace(
|
||||
tmp_path,
|
||||
{"schemaVersion": 1, "workspaceId": "demo", "relationships": []},
|
||||
)
|
||||
cfg.paths.effective_relationships.unlink()
|
||||
physical = cfg.paths.artifacts / "mschema" / "physical.yaml"
|
||||
|
||||
import tht.cli.search_cmd as search_module
|
||||
import tht.cli.vector_cmd as vector_module
|
||||
import tht.evidence as evidence_module
|
||||
import tht.search as search_core
|
||||
|
||||
monkeypatch.setattr(
|
||||
search_module,
|
||||
"_leased_dwh_snapshot",
|
||||
lambda _cfg, _ctx: SimpleNamespace(physical=physical, lsh_dir=tmp_path / "no-lsh"),
|
||||
)
|
||||
monkeypatch.setattr(vector_module, "require_vector_cfg", lambda _cfg: None)
|
||||
monkeypatch.setattr(vector_module, "open_searcher", lambda _cfg: object())
|
||||
monkeypatch.setattr(vector_module, "make_embedder", lambda _cfg: object())
|
||||
monkeypatch.setattr(evidence_module, "validate_corpus_workspace", lambda *_args: None)
|
||||
monkeypatch.setattr(evidence_module, "active_searcher", lambda _cfg, searcher, **_kwargs: searcher)
|
||||
monkeypatch.setattr(search_core, "combined_search", lambda **_kwargs: [])
|
||||
|
||||
result = RUNNER.invoke(
|
||||
app,
|
||||
[
|
||||
"search",
|
||||
"find",
|
||||
"orders",
|
||||
"--kind",
|
||||
"schema",
|
||||
"--json",
|
||||
"-c",
|
||||
str(tmp_path / "runtime.yaml"),
|
||||
],
|
||||
)
|
||||
|
||||
assert result.exit_code == 1
|
||||
assert "effective relationship snapshot is missing" in result.stderr
|
||||
|
||||
|
||||
def test_suggest_fks_cannot_write_when_relationships_are_catalog_managed(tmp_path):
|
||||
_workspace(
|
||||
tmp_path,
|
||||
{"schemaVersion": 1, "workspaceId": "demo", "relationships": []},
|
||||
)
|
||||
annotations = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
|
||||
before = annotations.read_text()
|
||||
|
||||
result = RUNNER.invoke(
|
||||
app,
|
||||
[
|
||||
"schema",
|
||||
"suggest-fks",
|
||||
"--write",
|
||||
"-c",
|
||||
str(tmp_path / "runtime.yaml"),
|
||||
],
|
||||
)
|
||||
|
||||
assert result.exit_code == 1
|
||||
assert "managed by the catalog" in result.stderr
|
||||
assert annotations.read_text() == before
|
||||
@@ -10,6 +10,12 @@ from tht.adapters.factory import build_dwh
|
||||
from tht.cli.config_cmd import CONFIG_OPT
|
||||
from tht.config import ConfigError, load_config
|
||||
from tht.db.sampling import is_text_type
|
||||
from tht.mschema.context import (
|
||||
SchemaContextError,
|
||||
annotations_path,
|
||||
load_schema_context,
|
||||
physical_path,
|
||||
)
|
||||
from tht.mschema.eligibility import classify_all
|
||||
|
||||
schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)")
|
||||
@@ -40,20 +46,6 @@ def _load_config_or_exit(config: Path):
|
||||
raise typer.Exit(code=1)
|
||||
|
||||
|
||||
def physical_path(cfg) -> Path:
|
||||
from tht.jobs.dwh_pipeline import resolve_dwh_snapshot
|
||||
|
||||
if not (cfg.paths.artifacts.parent / ".tht-dwh").exists():
|
||||
return cfg.paths.artifacts / "mschema" / "physical.yaml"
|
||||
return resolve_dwh_snapshot(cfg).physical
|
||||
|
||||
|
||||
def annotations_path(cfg) -> Path:
|
||||
if cfg.paths.annotations_root is not None:
|
||||
return cfg.paths.annotations_root / "mschema" / "annotations.yaml"
|
||||
return cfg.paths.artifacts / "mschema" / "annotations.yaml"
|
||||
|
||||
|
||||
def refresh_catalog(cfg, *, dwh=None, output_path: Path | None = None):
|
||||
"""Run the existing catalog algorithm and persist its canonical output."""
|
||||
target = dwh if dwh is not None else build_dwh(cfg)
|
||||
@@ -432,6 +424,21 @@ def suggest_fks_cmd(
|
||||
from tht.mschema.models import Annotations, PhysicalSchema, TableAnnotation
|
||||
|
||||
cfg = _load_config_or_exit(config)
|
||||
if write and cfg.paths.effective_relationships is not None:
|
||||
error = "effective relationships are managed by the catalog; --write is disabled"
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": "catalog_relationships_managed",
|
||||
"error": error,
|
||||
"operation": "schema_suggest_fks",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
typer.secho(f"ERRORE: {error}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1)
|
||||
phys_file = physical_path(cfg)
|
||||
if not phys_file.exists():
|
||||
if json_output:
|
||||
@@ -556,7 +563,6 @@ def render_cmd(
|
||||
"""Serializza mschema (physical + annotations) nel formato richiesto."""
|
||||
import json
|
||||
|
||||
from tht.mschema.models import Annotations, PhysicalSchema
|
||||
from tht.mschema.render import to_markdown, to_mschema_text, to_schema_dict
|
||||
|
||||
cfg = _load_config_or_exit(config)
|
||||
@@ -564,19 +570,40 @@ def render_cmd(
|
||||
if not phys_file.exists():
|
||||
typer.secho(
|
||||
f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.",
|
||||
fg=typer.colors.RED, err=True,
|
||||
fg=typer.colors.RED,
|
||||
err=True,
|
||||
)
|
||||
raise typer.Exit(code=1)
|
||||
physical = PhysicalSchema.from_yaml(phys_file)
|
||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||
try:
|
||||
context = load_schema_context(cfg, physical_file=phys_file)
|
||||
except SchemaContextError as exc:
|
||||
typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
table_filter = list(tables) if tables else None
|
||||
|
||||
if format == "markdown":
|
||||
out = to_markdown(physical, annotations)
|
||||
out = to_markdown(
|
||||
context.physical,
|
||||
context.annotations,
|
||||
effective_relationships=context.effective_relationships,
|
||||
)
|
||||
elif format == "mschema-text":
|
||||
out = to_mschema_text(physical, annotations, tables=table_filter)
|
||||
out = to_mschema_text(
|
||||
context.physical,
|
||||
context.annotations,
|
||||
tables=table_filter,
|
||||
effective_relationships=context.effective_relationships,
|
||||
)
|
||||
elif format == "schema-dict":
|
||||
out = json.dumps(to_schema_dict(physical, annotations), ensure_ascii=False, indent=2)
|
||||
out = json.dumps(
|
||||
to_schema_dict(
|
||||
context.physical,
|
||||
context.annotations,
|
||||
effective_relationships=context.effective_relationships,
|
||||
),
|
||||
ensure_ascii=False,
|
||||
indent=2,
|
||||
)
|
||||
else:
|
||||
typer.secho(f"ERRORE: formato sconosciuto: {format}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1)
|
||||
|
||||
@@ -231,8 +231,7 @@ def search_cmd(
|
||||
)
|
||||
|
||||
if kind == "schema":
|
||||
from tht.cli.schema_cmd import annotations_path
|
||||
from tht.mschema.models import Annotations, PhysicalSchema
|
||||
from tht.mschema.context import SchemaContextError, load_schema_context
|
||||
from tht.mschema.render import to_mschema_text
|
||||
from tht.search import schema_tables
|
||||
|
||||
@@ -243,6 +242,11 @@ def search_cmd(
|
||||
fg=typer.colors.RED, err=True,
|
||||
)
|
||||
raise typer.Exit(code=1)
|
||||
try:
|
||||
schema_context = load_schema_context(cfg, physical_file=phys_file)
|
||||
except SchemaContextError as exc:
|
||||
typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
|
||||
candidates = combined_search(
|
||||
keyword=keyword, lsh_hits=lsh_hits,
|
||||
@@ -258,10 +262,13 @@ def search_cmd(
|
||||
typer.secho(f"Nessuna tabella candidata per '{keyword}'.", fg=typer.colors.YELLOW)
|
||||
return
|
||||
|
||||
physical = PhysicalSchema.from_yaml(phys_file)
|
||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||
selected = [t for t, _ in ranked]
|
||||
mschema = to_mschema_text(physical, annotations, tables=selected)
|
||||
mschema = to_mschema_text(
|
||||
schema_context.physical,
|
||||
schema_context.annotations,
|
||||
tables=selected,
|
||||
effective_relationships=schema_context.effective_relationships,
|
||||
)
|
||||
|
||||
if json_out:
|
||||
typer.echo(json.dumps(
|
||||
|
||||
@@ -327,6 +327,9 @@ class PathsConfig(BaseModel):
|
||||
# Revision-qualified curated FK annotations root (P5). When absent, legacy
|
||||
# `artifacts/mschema/annotations.yaml` remains the annotations source.
|
||||
annotations_root: Path | None = None
|
||||
# Runtime-only, backend-derived effective relationship snapshot. When present,
|
||||
# it is the exclusive FK source; it is not an authored workspace artifact.
|
||||
effective_relationships: Path | None = None
|
||||
|
||||
|
||||
class RuntimeIdentityConfig(BaseModel):
|
||||
@@ -676,6 +679,7 @@ def load_config(path: Path) -> Config:
|
||||
indexes=resolved.indexes,
|
||||
memory=cfg.paths.memory,
|
||||
annotations_root=cfg.paths.annotations_root,
|
||||
effective_relationships=cfg.paths.effective_relationships,
|
||||
)
|
||||
}
|
||||
)
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
import json
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Annotated, Literal
|
||||
|
||||
from pydantic import BaseModel, Field, ValidationError, model_validator
|
||||
|
||||
from tht.mschema.models import Annotations, ForeignKey, PhysicalSchema
|
||||
|
||||
|
||||
class SchemaContextError(ValueError):
|
||||
"""The effective schema inputs cannot be used safely."""
|
||||
|
||||
|
||||
RelationshipName = Annotated[str, Field(min_length=1)]
|
||||
|
||||
|
||||
class EffectiveRelationship(BaseModel):
|
||||
source_table: RelationshipName = Field(alias="sourceTable")
|
||||
source_columns: list[RelationshipName] = Field(alias="sourceColumns", min_length=1)
|
||||
target_table: RelationshipName = Field(alias="targetTable")
|
||||
target_columns: list[RelationshipName] = Field(alias="targetColumns", min_length=1)
|
||||
origin: Literal["physical", "generated", "manual"]
|
||||
|
||||
model_config = {"extra": "forbid", "populate_by_name": True}
|
||||
|
||||
@model_validator(mode="after")
|
||||
def columns_are_paired(self):
|
||||
if len(self.source_columns) != len(self.target_columns):
|
||||
raise ValueError("sourceColumns and targetColumns must have the same length")
|
||||
return self
|
||||
|
||||
|
||||
class EffectiveRelationshipSnapshot(BaseModel):
|
||||
schema_version: Literal[1] = Field(alias="schemaVersion")
|
||||
workspace_id: RelationshipName = Field(alias="workspaceId")
|
||||
relationships: list[EffectiveRelationship]
|
||||
|
||||
model_config = {"extra": "forbid", "populate_by_name": True}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SchemaContext:
|
||||
physical: PhysicalSchema
|
||||
annotations: Annotations
|
||||
# None means legacy physical + annotation FK merging. A dict, including an
|
||||
# empty one, means the catalog snapshot is the exclusive relationship source.
|
||||
effective_relationships: dict[str, list[ForeignKey]] | None
|
||||
|
||||
|
||||
def physical_path(cfg) -> Path:
|
||||
from tht.jobs.dwh_pipeline import resolve_dwh_snapshot
|
||||
|
||||
if not (cfg.paths.artifacts.parent / ".tht-dwh").exists():
|
||||
return cfg.paths.artifacts / "mschema" / "physical.yaml"
|
||||
return resolve_dwh_snapshot(cfg).physical
|
||||
|
||||
|
||||
def annotations_path(cfg) -> Path:
|
||||
if cfg.paths.annotations_root is not None:
|
||||
return cfg.paths.annotations_root / "mschema" / "annotations.yaml"
|
||||
return cfg.paths.artifacts / "mschema" / "annotations.yaml"
|
||||
|
||||
|
||||
def _load_effective_relationships(cfg, physical: PhysicalSchema) -> dict[str, list[ForeignKey]] | None:
|
||||
path = cfg.paths.effective_relationships
|
||||
if path is None:
|
||||
return None
|
||||
if not path.is_file():
|
||||
raise SchemaContextError(f"effective relationship snapshot is missing: {path}")
|
||||
try:
|
||||
raw = json.loads(path.read_text())
|
||||
snapshot = EffectiveRelationshipSnapshot.model_validate(raw)
|
||||
except (OSError, json.JSONDecodeError, ValidationError) as exc:
|
||||
raise SchemaContextError(f"effective relationship snapshot is invalid: {path}") from exc
|
||||
if snapshot.workspace_id != cfg._workspace_id:
|
||||
raise SchemaContextError(
|
||||
"effective relationship snapshot workspace does not match runtime workspace"
|
||||
)
|
||||
|
||||
by_table: dict[str, list[ForeignKey]] = {}
|
||||
seen: set[tuple[str, tuple[str, ...], str, tuple[str, ...]]] = set()
|
||||
for relationship in snapshot.relationships:
|
||||
source = physical.tables.get(relationship.source_table)
|
||||
target = physical.tables.get(relationship.target_table)
|
||||
if source is None or target is None:
|
||||
raise SchemaContextError(
|
||||
"effective relationship endpoint table is absent from physical schema: "
|
||||
f"{relationship.source_table}->{relationship.target_table}"
|
||||
)
|
||||
missing_source = [name for name in relationship.source_columns if name not in source.columns]
|
||||
missing_target = [name for name in relationship.target_columns if name not in target.columns]
|
||||
if missing_source or missing_target:
|
||||
missing = ", ".join(
|
||||
[f"{relationship.source_table}.{name}" for name in missing_source]
|
||||
+ [f"{relationship.target_table}.{name}" for name in missing_target]
|
||||
)
|
||||
raise SchemaContextError(
|
||||
f"effective relationship endpoint column is absent from physical schema: {missing}"
|
||||
)
|
||||
key = (
|
||||
relationship.source_table,
|
||||
tuple(relationship.source_columns),
|
||||
relationship.target_table,
|
||||
tuple(relationship.target_columns),
|
||||
)
|
||||
if key in seen:
|
||||
continue
|
||||
seen.add(key)
|
||||
by_table.setdefault(relationship.source_table, []).append(
|
||||
ForeignKey(
|
||||
columns=relationship.source_columns,
|
||||
ref_table=relationship.target_table,
|
||||
ref_columns=relationship.target_columns,
|
||||
)
|
||||
)
|
||||
return by_table
|
||||
|
||||
|
||||
def load_schema_context(cfg, *, physical_file: Path | None = None) -> SchemaContext:
|
||||
physical_file = physical_file or physical_path(cfg)
|
||||
if not physical_file.is_file():
|
||||
raise SchemaContextError(f"physical schema is missing: {physical_file}")
|
||||
try:
|
||||
physical = PhysicalSchema.from_yaml(physical_file)
|
||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||
except (OSError, ValidationError, ValueError) as exc:
|
||||
raise SchemaContextError("physical schema or annotations are invalid") from exc
|
||||
return SchemaContext(
|
||||
physical=physical,
|
||||
annotations=annotations,
|
||||
effective_relationships=_load_effective_relationships(cfg, physical),
|
||||
)
|
||||
@@ -7,16 +7,23 @@ MAX_EXAMPLES_IN_PROMPT = 5
|
||||
|
||||
|
||||
def table_foreign_keys(
|
||||
physical: PhysicalSchema, annotations: Annotations, table: str
|
||||
physical: PhysicalSchema,
|
||||
annotations: Annotations,
|
||||
table: str,
|
||||
effective_relationships: dict[str, list[ForeignKey]] | None = None,
|
||||
) -> list[ForeignKey]:
|
||||
"""FK fisiche + FK logiche dalle annotations (dedup su columns/ref)."""
|
||||
if effective_relationships is not None:
|
||||
return list(effective_relationships.get(table, []))
|
||||
fks = list(physical.tables[table].foreign_keys)
|
||||
ann = annotations.tables.get(table)
|
||||
if ann:
|
||||
seen = {(tuple(f.columns), f.ref_table, tuple(f.ref_columns)) for f in fks}
|
||||
for fk in ann.foreign_keys:
|
||||
if (tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns)) not in seen:
|
||||
key = (tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))
|
||||
if key not in seen:
|
||||
fks.append(fk)
|
||||
seen.add(key)
|
||||
return fks
|
||||
|
||||
|
||||
@@ -47,6 +54,8 @@ def to_mschema_text(
|
||||
physical: PhysicalSchema,
|
||||
annotations: Annotations | None = None,
|
||||
tables: list[str] | None = None,
|
||||
*,
|
||||
effective_relationships: dict[str, list[ForeignKey]] | None = None,
|
||||
) -> str:
|
||||
"""Serializzazione testuale in stile ThothAI (【Schema】/【Foreign keys】)."""
|
||||
annotations = annotations or Annotations()
|
||||
@@ -73,7 +82,9 @@ def to_mschema_text(
|
||||
shown = ", ".join(column.examples[:MAX_EXAMPLES_IN_PROMPT])
|
||||
lines.append(f" -- Examples: {shown}")
|
||||
lines.append(");")
|
||||
for fk in table_foreign_keys(physical, annotations, table_name):
|
||||
for fk in table_foreign_keys(
|
||||
physical, annotations, table_name, effective_relationships
|
||||
):
|
||||
for src, dst in zip(fk.columns, fk.ref_columns):
|
||||
fk_lines.append(f"{table_name}.{src}={fk.ref_table}.{dst}")
|
||||
lines.extend(["", "【Foreign keys】", *fk_lines])
|
||||
@@ -81,7 +92,10 @@ def to_mschema_text(
|
||||
|
||||
|
||||
def to_schema_dict(
|
||||
physical: PhysicalSchema, annotations: Annotations | None = None
|
||||
physical: PhysicalSchema,
|
||||
annotations: Annotations | None = None,
|
||||
*,
|
||||
effective_relationships: dict[str, list[ForeignKey]] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Vista compatibile con le logiche AV-SQL (schema_dict)."""
|
||||
annotations = annotations or Annotations()
|
||||
@@ -103,13 +117,20 @@ def to_schema_dict(
|
||||
"primary_keys": [c for c in cols if table.columns[c].pk],
|
||||
"foreign_keys": [
|
||||
{"columns": fk.columns, "ref_table": fk.ref_table, "ref_columns": fk.ref_columns}
|
||||
for fk in table_foreign_keys(physical, annotations, table_name)
|
||||
for fk in table_foreign_keys(
|
||||
physical, annotations, table_name, effective_relationships
|
||||
)
|
||||
],
|
||||
}
|
||||
return out
|
||||
|
||||
|
||||
def to_markdown(physical: PhysicalSchema, annotations: Annotations | None = None) -> str:
|
||||
def to_markdown(
|
||||
physical: PhysicalSchema,
|
||||
annotations: Annotations | None = None,
|
||||
*,
|
||||
effective_relationships: dict[str, list[ForeignKey]] | None = None,
|
||||
) -> str:
|
||||
"""Report leggibile per il reviewer."""
|
||||
annotations = annotations or Annotations()
|
||||
lines = [
|
||||
@@ -146,7 +167,9 @@ def to_markdown(physical: PhysicalSchema, annotations: Annotations | None = None
|
||||
f"| {column_name} | {column.type} | {'sì' if column.nullable else 'no'} "
|
||||
f"| {'sì' if column.pk else ''} | {cdesc} | {examples} |"
|
||||
)
|
||||
fks = table_foreign_keys(physical, annotations, table_name)
|
||||
fks = table_foreign_keys(
|
||||
physical, annotations, table_name, effective_relationships
|
||||
)
|
||||
if fks:
|
||||
lines += ["", "Foreign keys:"]
|
||||
for fk in fks:
|
||||
|
||||
Reference in New Issue
Block a user