From c0d50e9b08612ee11ee39e2124a0db799e543938 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 01:21:29 +0200 Subject: [PATCH] feat(vector): version pgvector schema --- .superpowers/sdd/pgvector-task-2-report.md | 52 +++++ harness/migrations/vector/001_extensions.sql | 1 + .../migrations/vector/002_schema_tables.sql | 36 +++ harness/migrations/vector/003_roles.sql | 23 ++ harness/tests/l0/test_vector_migrations.py | 205 ++++++++++++++++++ harness/tht/cli/__init__.py | 1 + harness/tht/cli/vector_migrate_cmd.py | 191 ++++++++++++++++ 7 files changed, 509 insertions(+) create mode 100644 .superpowers/sdd/pgvector-task-2-report.md create mode 100644 harness/migrations/vector/001_extensions.sql create mode 100644 harness/migrations/vector/002_schema_tables.sql create mode 100644 harness/migrations/vector/003_roles.sql create mode 100644 harness/tests/l0/test_vector_migrations.py create mode 100644 harness/tht/cli/vector_migrate_cmd.py diff --git a/.superpowers/sdd/pgvector-task-2-report.md b/.superpowers/sdd/pgvector-task-2-report.md new file mode 100644 index 00000000..7fdb5ea2 --- /dev/null +++ b/.superpowers/sdd/pgvector-task-2-report.md @@ -0,0 +1,52 @@ +# Local pgvector Task 2 report + +## Outcome + +Implemented ordered, idempotent production migrations and the `tht vector migrate` +interface, including `tht vector migrate --status --json` with pristine JSON output. + +## Implementation + +- `001_extensions.sql` installs pgvector. +- `002_schema_tables.sql` creates `vectors.schema_records`, `vectors.evidence`, and + `vectors.memory` with the `VectorWriteRecord` columns and `vector(768)` embeddings. +- `003_roles.sql` creates passwordless `NOLOGIN` reader/writer roles. Deployments inject + credentials (or grant these roles to separately-created login roles); no production secret + is stored in the repository. +- Reader authority is schema usage plus table `SELECT`. +- Writer authority is schema usage, table `INSERT`/`UPDATE`, narrow hash-probe column `SELECT`, + and sequence `USAGE`. It has no `DELETE`, broad row `SELECT`, DDL, or ownership authority. +- The migration runner discovers ordered SQL files, records SHA-256 checksums in + `public.tht_vector_migrations`, serializes runners with a transaction-scoped advisory lock, + and applies the full pending batch in one transaction. +- Status distinguishes applied, pending, and checksum-drifted migrations. Apply refuses drift. + A failed migration rolls back both prior migrations in that batch and ledger writes. + +## TDD evidence + +RED was observed with a real `pgvector/pgvector:pg16` testcontainer: 6 failures for the missing +module, missing command, and missing schema. + +GREEN verification: + +- Focused migration + direct adapter integration: `23 passed`. +- Full harness from the documented `harness/` cwd: `473 passed, 5 deselected`. +- Targeted Ruff (`tht` plus the new L0 test): clean. +- `git diff --check`: clean. + +The new L0 coverage exercises clean install, idempotent rerun, pristine JSON status, checksum +drift, transaction rollback, exact tables/columns/dimensions, role isolation, sequence authority, +and the real `PgVectorStore.health()` plus `VectorWriteRecord` upsert path. + +## Existing repository lint baseline + +The requested full `ruff check .` was run. It reports 34 pre-existing violations in unrelated +test files (unused imports and one-line semicolon statements). None are in Task 2 files; changing +them would exceed this task's scope. The complete harness test gate is green. + +## Self-review + +No unresolved Task 2 correctness concern found. One deliberate contract choice is worth noting: +writer `INSERT` and `UPDATE` are table-level because the approved direct adapter health probe uses +`has_table_privilege` for those authorities. Least privilege is retained by withholding broad +`SELECT`, `DELETE`, DDL, ownership, and credentials. diff --git a/harness/migrations/vector/001_extensions.sql b/harness/migrations/vector/001_extensions.sql new file mode 100644 index 00000000..0aa0fc22 --- /dev/null +++ b/harness/migrations/vector/001_extensions.sql @@ -0,0 +1 @@ +CREATE EXTENSION IF NOT EXISTS vector; diff --git a/harness/migrations/vector/002_schema_tables.sql b/harness/migrations/vector/002_schema_tables.sql new file mode 100644 index 00000000..640d593c --- /dev/null +++ b/harness/migrations/vector/002_schema_tables.sql @@ -0,0 +1,36 @@ +CREATE SCHEMA IF NOT EXISTS vectors; + +REVOKE ALL ON SCHEMA vectors FROM PUBLIC; + +CREATE TABLE IF NOT EXISTS vectors.schema_records ( + id bigserial PRIMARY KEY, + record_key text UNIQUE NOT NULL, + kind text NOT NULL, + content_hash text NOT NULL, + metadata jsonb NOT NULL, + embedding vector(768) NOT NULL, + indexed_at timestamptz NOT NULL DEFAULT now() +); + +CREATE TABLE IF NOT EXISTS vectors.evidence ( + id bigserial PRIMARY KEY, + record_key text UNIQUE NOT NULL, + kind text NOT NULL, + content_hash text NOT NULL, + metadata jsonb NOT NULL, + embedding vector(768) NOT NULL, + indexed_at timestamptz NOT NULL DEFAULT now() +); + +CREATE TABLE IF NOT EXISTS vectors.memory ( + id bigserial PRIMARY KEY, + record_key text UNIQUE NOT NULL, + kind text NOT NULL, + content_hash text NOT NULL, + metadata jsonb NOT NULL, + embedding vector(768) NOT NULL, + indexed_at timestamptz NOT NULL DEFAULT now() +); + +REVOKE ALL ON ALL TABLES IN SCHEMA vectors FROM PUBLIC; +REVOKE ALL ON ALL SEQUENCES IN SCHEMA vectors FROM PUBLIC; diff --git a/harness/migrations/vector/003_roles.sql b/harness/migrations/vector/003_roles.sql new file mode 100644 index 00000000..34de7cf5 --- /dev/null +++ b/harness/migrations/vector/003_roles.sql @@ -0,0 +1,23 @@ +DO $roles$ +BEGIN + IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'vector_reader') THEN + CREATE ROLE vector_reader NOLOGIN; + END IF; + IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'vector_writer') THEN + CREATE ROLE vector_writer NOLOGIN; + END IF; +END +$roles$; + +REVOKE ALL ON SCHEMA vectors FROM vector_reader, vector_writer; +REVOKE ALL ON ALL TABLES IN SCHEMA vectors FROM vector_reader, vector_writer; +REVOKE ALL ON ALL SEQUENCES IN SCHEMA vectors FROM vector_reader, vector_writer; + +GRANT USAGE ON SCHEMA vectors TO vector_reader, vector_writer; +GRANT SELECT ON ALL TABLES IN SCHEMA vectors TO vector_reader; + +GRANT INSERT, UPDATE +ON vectors.schema_records, vectors.evidence, vectors.memory TO vector_writer; +GRANT SELECT (record_key, kind, content_hash) +ON vectors.schema_records, vectors.evidence, vectors.memory TO vector_writer; +GRANT USAGE ON ALL SEQUENCES IN SCHEMA vectors TO vector_writer; diff --git a/harness/tests/l0/test_vector_migrations.py b/harness/tests/l0/test_vector_migrations.py new file mode 100644 index 00000000..15a208bf --- /dev/null +++ b/harness/tests/l0/test_vector_migrations.py @@ -0,0 +1,205 @@ +import json +from pathlib import Path + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.exc import ProgrammingError +from testcontainers.postgres import PostgresContainer +from typer.testing import CliRunner + +from tht.cli import app +from tht.config import DatabaseConfig +from tht.ports.vector import VectorRecord, VectorWriteRecord + + +@pytest.fixture(scope="module") +def database_url(): + with PostgresContainer("pgvector/pgvector:pg16") as postgres: + yield postgres.get_connection_url() + + +def test_migrations_are_clean_and_idempotent(database_url): + from tht.cli.vector_migrate_cmd import migrate, migration_status + + before = migration_status(database_url) + assert [item.version for item in before.pending] == ["001", "002", "003"] + + migrate(database_url) + migrate(database_url) + + status = migration_status(database_url) + assert status.pending == () + assert status.drifted == () + assert [item.version for item in status.applied] == ["001", "002", "003"] + + +def test_schema_matches_direct_adapter_contract(database_url): + engine = create_engine(database_url) + with engine.connect() as connection: + rows = connection.execute( + text( + "SELECT table_name, column_name, data_type, udt_name " + "FROM information_schema.columns WHERE table_schema = 'vectors' " + "ORDER BY table_name, ordinal_position" + ) + ).all() + vector_types = connection.execute( + text( + "SELECT c.relname, format_type(a.atttypid, a.atttypmod) " + "FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace " + "JOIN pg_attribute a ON a.attrelid = c.oid AND a.attname = 'embedding' " + "WHERE n.nspname = 'vectors' ORDER BY c.relname" + ) + ).all() + engine.dispose() + + tables = {row.table_name for row in rows} + assert tables == {"evidence", "memory", "schema_records"} + required = {"id", "record_key", "kind", "content_hash", "metadata", "embedding", "indexed_at"} + for table in tables: + assert {row.column_name for row in rows if row.table_name == table} == required + assert vector_types == [ + ("evidence", "vector(768)"), + ("memory", "vector(768)"), + ("schema_records", "vector(768)"), + ] + + +def test_roles_have_runtime_privileges_only(database_url): + from tht.cli.vector_migrate_cmd import migrate + + migrate(database_url) + admin = create_engine(database_url) + with admin.begin() as connection: + connection.exec_driver_sql("ALTER ROLE vector_reader LOGIN PASSWORD 'reader-test-only'") + connection.exec_driver_sql("ALTER ROLE vector_writer LOGIN PASSWORD 'writer-test-only'") + url = admin.url + reader = create_engine(url.set(username="vector_reader", password="reader-test-only")) + writer = create_engine(url.set(username="vector_writer", password="writer-test-only")) + + from tht.adapters.vector.pgvector import PgVectorStore + + common = { + "host": url.host, + "port": url.port, + "database": url.database, + "schema": "vectors", + } + reader_config = DatabaseConfig( + **common, user="vector_reader", password="reader-test-only" + ) + writer_config = DatabaseConfig( + **common, user="vector_writer", password="writer-test-only" + ) + store = PgVectorStore(reader_config, writer_config, expected_dimension=768) + assert store.health().ok is True + assert store.upsert( + "memory", + [ + VectorWriteRecord( + record=VectorRecord( + id="adapter-write", + kind="memory", + ref="session:test", + title="test", + content="test", + ), + embedding=[0.0] * 768, + content_hash="adapter-hash", + ) + ], + ) == 1 + + with reader.connect() as connection: + connection.execute(text("SELECT metadata, embedding FROM vectors.memory")).all() + with pytest.raises(ProgrammingError): + with reader.begin() as connection: + connection.execute( + text( + "INSERT INTO vectors.memory " + "(record_key, kind, content_hash, metadata, embedding) " + "VALUES ('reader-write', 'memory', 'x', '{}', " + "array_fill(0, ARRAY[768])::vector)" + ) + ) + + with writer.begin() as connection: + connection.execute( + text( + "INSERT INTO vectors.memory " + "(record_key, kind, content_hash, metadata, embedding) " + "VALUES ('writer-ok', 'memory', 'x', '{}', array_fill(0, ARRAY[768])::vector)" + ) + ) + assert connection.execute( + text("SELECT content_hash FROM vectors.memory WHERE record_key = 'writer-ok'") + ).scalar_one() == "x" + connection.execute( + text("UPDATE vectors.memory SET content_hash = 'y' WHERE record_key = 'writer-ok'") + ) + with pytest.raises(ProgrammingError): + with writer.connect() as connection: + connection.execute(text("SELECT metadata FROM vectors.memory")).all() + with pytest.raises(ProgrammingError): + with writer.begin() as connection: + connection.execute(text("DELETE FROM vectors.memory WHERE record_key = 'writer-ok'")) + + reader.dispose() + writer.dispose() + admin.dispose() + + +def test_status_json_is_pristine(database_url, monkeypatch): + monkeypatch.setenv("THT_VECTOR_ADMIN_URL", database_url) + result = CliRunner().invoke(app, ["vector", "migrate", "--status", "--json"]) + + assert result.exit_code == 0, result.output + assert json.loads(result.stdout) == { + "applied": ["001", "002", "003"], + "drifted": [], + "pending": [], + } + assert result.stderr == "" + + +def test_checksum_drift_is_reported_and_refused(database_url, tmp_path): + from tht.cli.vector_migrate_cmd import MigrationError, migrate, migration_status + + migrations = _copy_migrations(tmp_path) + migrate(database_url, migrations) + (migrations / "002_schema_tables.sql").write_text("SELECT 2;\n") + + assert [item.version for item in migration_status(database_url, migrations).drifted] == [ + "002" + ] + with pytest.raises(MigrationError, match="checksum drift"): + migrate(database_url, migrations) + + +def test_failed_batch_rolls_back_schema_and_ledger(database_url, tmp_path): + from tht.cli.vector_migrate_cmd import MigrationError, migrate, migration_status + + migrations = tmp_path / "failed" + migrations.mkdir() + (migrations / "101_first.sql").write_text("CREATE TABLE public.must_rollback (id int);\n") + (migrations / "102_broken.sql").write_text("THIS IS NOT SQL;\n") + + with pytest.raises(MigrationError, match="102_broken.sql"): + migrate(database_url, migrations) + + engine = create_engine(database_url) + with engine.connect() as connection: + assert connection.execute(text("SELECT to_regclass('public.must_rollback')")).scalar() is None + engine.dispose() + status = migration_status(database_url, migrations) + assert status.applied == () + assert [item.version for item in status.pending] == ["101", "102"] + + +def _copy_migrations(tmp_path: Path) -> Path: + source = Path(__file__).parents[2] / "migrations" / "vector" + target = tmp_path / "migrations" + target.mkdir() + for migration in source.glob("*.sql"): + (target / migration.name).write_bytes(migration.read_bytes()) + return target diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index 3ed06d88..afa18f6f 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -53,6 +53,7 @@ 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 app.add_typer(phase_app, name="phase") app.add_typer(config_app, name="config") diff --git a/harness/tht/cli/vector_migrate_cmd.py b/harness/tht/cli/vector_migrate_cmd.py new file mode 100644 index 00000000..51e9c1f8 --- /dev/null +++ b/harness/tht/cli/vector_migrate_cmd.py @@ -0,0 +1,191 @@ +"""Versioned, transactional migrations for the direct pgvector schema.""" + +from __future__ import annotations + +import hashlib +import json +import re +from dataclasses import dataclass +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 = Path(__file__).parents[2] / "migrations" / "vector" +_MIGRATION_NAME = re.compile(r"^(?P\d+)_(?P[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: Path + checksum: str + + +@dataclass(frozen=True) +class MigrationStatus: + applied: tuple[Migration, ...] + pending: tuple[Migration, ...] + drifted: tuple[Migration, ...] + + +def _discover(directory: Path) -> tuple[Migration, ...]: + migrations = [] + seen_versions: set[str] = set() + for path in sorted(directory.glob("*.sql")): + match = _MIGRATION_NAME.fullmatch(path.name) + if match is None: + raise MigrationError(f"Invalid migration filename: {path.name}") + version = match.group("version") + if version in seen_versions: + raise MigrationError(f"Duplicate migration version: {version}") + seen_versions.add(version) + migrations.append( + Migration( + version=version, + name=match.group("name"), + path=path, + checksum=hashlib.sha256(path.read_bytes()).hexdigest(), + ) + ) + if not migrations: + raise MigrationError(f"No migrations found in {directory}") + return tuple(migrations) + + +def _applied(connection) -> dict[str, str]: + exists = connection.execute(text("SELECT 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 migration_status( + database_url: str, migrations_dir: Path | str = MIGRATIONS_DIR +) -> MigrationStatus: + migrations = _discover(Path(migrations_dir)) + engine = create_engine(database_url) + try: + with engine.connect() as connection: + applied_checksums = _applied(connection) + finally: + engine.dispose() + 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: Path | str = MIGRATIONS_DIR) -> MigrationStatus: + migrations = _discover(Path(migrations_dir)) + engine = create_engine(database_url) + current: Migration | None = None + try: + with engine.begin() as connection: + connection.execute(text("SELECT 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 now() + )""" + ) + connection.exec_driver_sql( + "REVOKE ALL ON public.tht_vector_migrations FROM PUBLIC" + ) + applied_checksums = _applied(connection) + 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"]