fix: remove dead vector migrate command

This commit is contained in:
2026-08-08 21:45:16 +02:00
parent 8826f8ac6b
commit 569bb7f1db
10 changed files with 100 additions and 262 deletions
+18 -19
View File
@@ -36,25 +36,24 @@ def main(
pass
from tht.cli.config_cmd import config_app # noqa: E402
from tht.cli.cte_cmd import cte_app # noqa: E402
from tht.cli.datamart_cmd import datamart_app # noqa: E402
from tht.cli.db_cmd import db_app # noqa: E402
from tht.cli.decision_cmd import decision_app # noqa: E402
from tht.cli.doctor_cmd import doctor # noqa: E402
from tht.cli.evidence_cmd import evidence_app # noqa: E402
from tht.cli.formula_cmd import formula_app # noqa: E402
from tht.cli.lsh_cmd import lsh_app # noqa: E402
from tht.cli.memory_cmd import memory_app # noqa: E402
from tht.cli.ollama_cmd import ollama_app # noqa: E402
from tht.cli.phase_cmd import phase_app # noqa: E402
from tht.cli.preprocess_cmd import preprocess_app # noqa: E402
from tht.cli.schema_cmd import schema_app # noqa: E402
from tht.cli.search_cmd import search_app # noqa: E402
from tht.cli.session_cmd import session_app # noqa: E402
from tht.cli.sql_cmd import sql_app # noqa: E402
from tht.cli.vector_cmd import vector_app # noqa: E402
import tht.cli.vector_migrate_cmd # noqa: E402, F401
from tht.cli.config_cmd import config_app
from tht.cli.cte_cmd import cte_app
from tht.cli.datamart_cmd import datamart_app
from tht.cli.db_cmd import db_app
from tht.cli.decision_cmd import decision_app
from tht.cli.doctor_cmd import doctor
from tht.cli.evidence_cmd import evidence_app
from tht.cli.formula_cmd import formula_app
from tht.cli.lsh_cmd import lsh_app
from tht.cli.memory_cmd import memory_app
from tht.cli.ollama_cmd import ollama_app
from tht.cli.phase_cmd import phase_app
from tht.cli.preprocess_cmd import preprocess_app
from tht.cli.schema_cmd import schema_app
from tht.cli.search_cmd import search_app
from tht.cli.session_cmd import session_app
from tht.cli.sql_cmd import sql_app
from tht.cli.vector_cmd import vector_app
app.add_typer(phase_app, name="phase")
app.add_typer(preprocess_app, name="preprocess")
+10 -10
View File
@@ -19,7 +19,7 @@ from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli.session_cmd import load_snapshot_or_exit
from tht.cli.vector_cmd import require_vector_cfg
memory_app = typer.Typer(help="Review memory (registro canonico + indice pgvector)")
memory_app = typer.Typer(help="Review memory (registro canonico + indice semantico)")
DECISION_OPT = typer.Option(None, "--decision", help="Seq da promuovere (ripetibile).")
@@ -28,7 +28,7 @@ def registry_path(cfg) -> Path:
def _resync_memory(cfg):
"""Risincronizza l'indice pgvector col registro corrente (incrementale)."""
"""Risincronizza l'indice semantico col registro corrente (incrementale)."""
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder, sync_canonical_records
from tht.memory import load_registry, memory_vector_records
@@ -111,8 +111,8 @@ def promote_cmd(
typer.secho(msg, fg=typer.colors.YELLOW)
return
# Promozione nel registro: riuscita. L'indicizzazione su pgvector puo' fallire
# (tabella mancante o vectordb irraggiungibile da questa postazione): in quel
# Promozione nel registro: riuscita. L'indicizzazione semantica puo' fallire
# (runtime non pronto o vectordb irraggiungibile da questa postazione): in quel
# caso le memorie restano nel registro ma NON sono trovate da `tht memory
# search` finche' non si reindicizza sul server. `indexed` rende lo stato
# leggibile da Pi, cosi' il reviewer lo vede invece di perderlo nello stderr.
@@ -125,7 +125,7 @@ def promote_cmd(
indexed = False
warning = (
f"{len(promoted)} memorie promosse nel registro, ma l'indice vettoriale "
"NON e' stato sincronizzato (tabella pgvector mancante o irraggiungibile): "
"NON e' stato sincronizzato (runtime vettoriale mancante o irraggiungibile): "
"NON saranno trovate da `tht memory search` finche' non reindicizzi sul "
"server (`tht vector init`, poi `tht memory index`)."
)
@@ -154,11 +154,11 @@ def save_one_cmd(
json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."),
config: Path = CONFIG_OPT,
) -> None:
"""Upsert mirato (una riga) della memoria di una decisione su pgvector (D11).
"""Upsert mirato (una riga) della memoria di una decisione nel semantic store (D11).
Promuove la decisione nel registro locale (idempotente) e fa un singolo upsert
con dedup hash client-side -- niente full-resync. Il factory seleziona il writer
REST su workstation oppure il writer pgvector diretto sul profilo server.
del runtime vettoriale attivo.
"""
import json as _json
@@ -180,7 +180,7 @@ def save_one_cmd(
count = save_one_memory(records, decision, store=store, embedder=embedder)
msg = (
f"{count} memoria salvata su pgvector (decision_seq {decision})."
f"{count} memoria salvata nell'indice semantico (decision_seq {decision})."
if count
else f"Nessun upsert (decisione {decision} assente/stale o memoria gia' aggiornata)."
)
@@ -196,7 +196,7 @@ def clear_cmd(
yes: bool = typer.Option(False, "--yes", "-y", help="Salta la richiesta di conferma."),
config: Path = CONFIG_OPT,
) -> None:
"""Cancella TUTTA la review memory: registro canonico + indice pgvector (kind=memory)."""
"""Cancella TUTTA la review memory: registro canonico + indice semantico (kind=memory)."""
from tht.memory import load_registry
cfg = _load_config_or_exit(config)
@@ -222,7 +222,7 @@ def clear_cmd(
@memory_app.command("index")
def index_cmd(config: Path = CONFIG_OPT) -> None:
"""Sincronizza il registro memory su pgvector (full-resync)."""
"""Sincronizza il registro memory nell'indice semantico (full-resync)."""
from tht.cli.vector_cmd import _print_stats
cfg = _load_config_or_exit(config)
+1 -1
View File
@@ -49,7 +49,7 @@ def search_cmd(
False, "--json", help="Output JSON machine-readable per Pi (sopprime le tabelle a video)."
),
) -> None:
"""Ricerca combinata LSH + pgvector con ranking RRF spiegabile."""
"""Ricerca combinata LSH + semantic search con ranking RRF spiegabile."""
from rich.console import Console
from rich.table import Table
+3 -4
View File
@@ -8,7 +8,7 @@ from tht.cli.schema_cmd import _load_config_or_exit, annotations_path, physical_
from tht.ports.vector import VectorWriteRecord
from tht.vectorstore.store import SyncStats, content_hash
vector_app = typer.Typer(help="Indice semantico pgvector (derivato, rigenerabile)")
vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)")
def make_embedder(embeddings_cfg):
@@ -33,8 +33,7 @@ def require_vector_cfg(cfg):
def open_searcher(cfg):
"""Searcher per la LETTURA (similarity search): via REST se `vector_rest` è configurato,
altrimenti connessione diretta (dev/test)."""
"""Searcher per la lettura semantic search sul runtime vettoriale attivo."""
from tht.adapters.factory import build_vector_store
from tht.vectorstore.reader import tables_for_kinds
@@ -119,7 +118,7 @@ def init_cmd(
@vector_app.command("index-schema")
def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
"""Embedda e sincronizza i record schema (tabelle e colonne) da mschema."""
"""Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store."""
from tht.mschema.models import Annotations, PhysicalSchema
from tht.vectorstore.records import schema_records
-227
View File
@@ -1,227 +0,0 @@
"""Versioned, transactional migrations for the direct pgvector schema."""
from __future__ import annotations
import hashlib
import json
import re
from dataclasses import dataclass
from importlib.resources import files
from importlib.resources.abc import Traversable
from pathlib import Path
import typer
from sqlalchemy import create_engine, text
from sqlalchemy.exc import SQLAlchemyError
from tht.cli.vector_cmd import vector_app
MIGRATIONS_DIR = files("tht").joinpath("migrations", "vector")
_MIGRATION_NAME = re.compile(r"^(?P<version>\d+)_(?P<name>[a-z0-9_]+)\.sql$")
_LOCK_KEY = 7_304_708_654_221_909_028
class MigrationError(RuntimeError):
"""Raised when migration discovery or application is unsafe."""
@dataclass(frozen=True)
class Migration:
version: str
name: str
path: Traversable
checksum: str
@dataclass(frozen=True)
class MigrationStatus:
applied: tuple[Migration, ...]
pending: tuple[Migration, ...]
drifted: tuple[Migration, ...]
def _migration_source(directory: Traversable | Path | str) -> Traversable:
return Path(directory) if isinstance(directory, (str, Path)) else directory
def _discover(directory: Traversable | Path | str) -> tuple[Migration, ...]:
source = _migration_source(directory)
migrations = []
seen_versions: set[int] = set()
paths = [path for path in source.iterdir() if path.name.endswith(".sql")]
parsed = []
for path in paths:
match = _MIGRATION_NAME.fullmatch(path.name)
if match is None:
raise MigrationError(f"Invalid migration filename: {path.name}")
version = match.group("version")
numeric_version = int(version)
if numeric_version in seen_versions:
raise MigrationError(f"Duplicate migration version: {numeric_version}")
seen_versions.add(numeric_version)
parsed.append((numeric_version, version, match.group("name"), path))
for _, version, name, path in sorted(parsed, key=lambda item: item[0]):
migrations.append(
Migration(
version=version,
name=name,
path=path,
checksum=hashlib.sha256(path.read_bytes()).hexdigest(),
)
)
if not migrations:
raise MigrationError(f"No migrations found in {source}")
return tuple(migrations)
def _applied(connection) -> dict[str, str]:
exists = connection.execute(
text("SELECT pg_catalog.to_regclass('public.tht_vector_migrations')")
).scalar()
if exists is None:
return {}
return dict(
connection.execute(
text("SELECT version, checksum FROM public.tht_vector_migrations")
).all()
)
def _reject_unknown_versions(
migrations: tuple[Migration, ...], applied_checksums: dict[str, str]
) -> None:
local_versions = {migration.version for migration in migrations}
unknown = sorted(
set(applied_checksums) - local_versions,
key=lambda version: (0, int(version)) if version.isdigit() else (1, version),
)
if unknown:
raise MigrationError(
"Database migration versions absent from local manifest: " + ", ".join(unknown)
)
def migration_status(
database_url: str, migrations_dir: Traversable | Path | str = MIGRATIONS_DIR
) -> MigrationStatus:
migrations = _discover(migrations_dir)
engine = create_engine(database_url)
try:
with engine.connect() as connection:
connection.exec_driver_sql("SET LOCAL search_path = pg_catalog, pg_temp")
applied_checksums = _applied(connection)
finally:
engine.dispose()
_reject_unknown_versions(migrations, applied_checksums)
applied = tuple(
migration
for migration in migrations
if applied_checksums.get(migration.version) == migration.checksum
)
drifted = tuple(
migration
for migration in migrations
if migration.version in applied_checksums
and applied_checksums[migration.version] != migration.checksum
)
pending = tuple(
migration for migration in migrations if migration.version not in applied_checksums
)
return MigrationStatus(applied=applied, pending=pending, drifted=drifted)
def migrate(
database_url: str, migrations_dir: Traversable | Path | str = MIGRATIONS_DIR
) -> MigrationStatus:
migrations = _discover(migrations_dir)
engine = create_engine(database_url)
current: Migration | None = None
try:
with engine.begin() as connection:
connection.exec_driver_sql("SET LOCAL search_path = pg_catalog, pg_temp")
connection.execute(
text("SELECT pg_catalog.pg_advisory_xact_lock(:key)"), {"key": _LOCK_KEY}
)
connection.exec_driver_sql(
"""CREATE TABLE IF NOT EXISTS public.tht_vector_migrations (
version text PRIMARY KEY,
name text NOT NULL,
checksum text NOT NULL,
applied_at timestamptz NOT NULL DEFAULT pg_catalog.now()
)"""
)
connection.exec_driver_sql(
"REVOKE ALL ON public.tht_vector_migrations FROM PUBLIC"
)
applied_checksums = _applied(connection)
_reject_unknown_versions(migrations, applied_checksums)
drifted = [
item
for item in migrations
if item.version in applied_checksums
and applied_checksums[item.version] != item.checksum
]
if drifted:
versions = ", ".join(item.version for item in drifted)
raise MigrationError(f"Migration checksum drift: {versions}")
for current in migrations:
if current.version in applied_checksums:
continue
connection.exec_driver_sql(current.path.read_text())
connection.execute(
text(
"INSERT INTO public.tht_vector_migrations (version, name, checksum) "
"VALUES (:version, :name, :checksum)"
),
{
"version": current.version,
"name": current.name,
"checksum": current.checksum,
},
)
except MigrationError:
raise
except SQLAlchemyError as exc:
filename = current.path.name if current is not None else "migration setup"
raise MigrationError(f"Failed to apply {filename}: {type(exc).__name__}") from exc
finally:
engine.dispose()
return migration_status(database_url, migrations_dir)
def _payload(status: MigrationStatus) -> dict[str, list[str]]:
return {
"applied": [item.version for item in status.applied],
"drifted": [item.version for item in status.drifted],
"pending": [item.version for item in status.pending],
}
@vector_app.command("migrate")
def migrate_cmd(
database_url: str = typer.Option(
..., "--database-url", envvar="THT_VECTOR_ADMIN_URL", help="Admin PostgreSQL URL."
),
status_only: bool = typer.Option(False, "--status", help="Inspect without applying."),
json_output: bool = typer.Option(False, "--json", help="Emit pristine JSON."),
) -> None:
"""Apply or inspect the local pgvector schema migrations."""
try:
status = migration_status(database_url) if status_only else migrate(database_url)
except (MigrationError, SQLAlchemyError) as exc:
if json_output:
typer.echo(json.dumps({"error": str(exc)}, sort_keys=True))
else:
typer.echo(f"ERROR: {exc}", err=True)
raise typer.Exit(code=1) from None
payload = _payload(status)
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.echo(
f"Applied: {len(status.applied)}; pending: {len(status.pending)}; "
f"drifted: {len(status.drifted)}"
)
__all__ = ["MigrationError", "MigrationStatus", "migrate", "migration_status"]