Files
ThothII/harness/tests/l0/test_pgvector_store.py

408 lines
17 KiB
Python

import pytest
from psycopg2.errors import InsufficientPrivilege
from sqlalchemy import create_engine, text
from sqlalchemy.exc import ProgrammingError
from testcontainers.postgres import PostgresContainer
from tht.adapters.vector.thoth_http import ThothHttpVectorStore
from tht.config import DatabaseConfig
from tht.ports.vector import (
VectorReadUnavailable,
VectorRecord,
VectorStoreError,
VectorWriteRecord,
VectorWriteUnavailable,
)
def _record(content_hash: str, embedding: list[float], *, kind: str = "memory"):
return VectorWriteRecord(
record=VectorRecord(
id=f"record:{content_hash}",
kind=kind,
ref="session:test",
title=content_hash,
content=f"content {content_hash}",
metadata={"content_hash": content_hash},
),
embedding=embedding,
content_hash=content_hash,
)
@pytest.fixture(scope="module")
def vector_configs():
with PostgresContainer("pgvector/pgvector:pg16") as pg:
host = pg.get_container_host_ip()
port = int(pg.get_exposed_port(5432))
admin_config = DatabaseConfig(
host=host,
port=port,
database=pg.dbname,
schema="vectors",
user=pg.username,
password=pg.password,
)
engine = create_engine(pg.get_connection_url())
with engine.begin() as connection:
connection.exec_driver_sql("CREATE SCHEMA vectors")
# Match the co-located Supabase deployment: tables are in vectors, extension in public.
connection.exec_driver_sql("CREATE EXTENSION vector WITH SCHEMA public")
for table in ("schema_records", "evidence", "memory"):
connection.exec_driver_sql(f"""
CREATE TABLE vectors.{table} (
id bigserial PRIMARY KEY,
record_key text UNIQUE NOT NULL,
kind text NOT NULL,
content_hash text NOT NULL,
metadata jsonb NOT NULL,
embedding public.vector(2) NOT NULL,
indexed_at timestamptz NOT NULL DEFAULT now()
)
""")
connection.exec_driver_sql("CREATE ROLE vector_l0_reader LOGIN PASSWORD 'reader'")
connection.exec_driver_sql("CREATE ROLE vector_l0_writer LOGIN PASSWORD 'writer'")
connection.exec_driver_sql(
"CREATE ROLE vector_l0_no_sequence LOGIN PASSWORD 'no_sequence'"
)
connection.exec_driver_sql(
"GRANT USAGE ON SCHEMA vectors TO vector_l0_reader, vector_l0_writer, "
"vector_l0_no_sequence"
)
connection.exec_driver_sql(
"GRANT SELECT ON ALL TABLES IN SCHEMA vectors TO vector_l0_reader"
)
connection.exec_driver_sql(
"GRANT USAGE, SELECT ON ALL SEQUENCES IN SCHEMA vectors TO vector_l0_writer"
)
for table in ("schema_records", "evidence", "memory"):
connection.exec_driver_sql(
f"GRANT INSERT, UPDATE ON vectors.{table} "
"TO vector_l0_writer, vector_l0_no_sequence"
)
if table == "evidence":
connection.exec_driver_sql(
"GRANT DELETE ON vectors.evidence TO vector_l0_writer"
)
connection.exec_driver_sql(
"GRANT SELECT (kind, metadata) ON vectors.evidence TO vector_l0_writer"
)
connection.exec_driver_sql(
f"GRANT SELECT (record_key, kind, content_hash) "
f"ON vectors.{table} TO vector_l0_writer, vector_l0_no_sequence"
)
engine.dispose()
reader_config = admin_config.model_copy(
update={"user": "vector_l0_reader", "password": "reader"}
)
writer_config = admin_config.model_copy(
update={"user": "vector_l0_writer", "password": "writer"}
)
no_sequence_config = admin_config.model_copy(
update={"user": "vector_l0_no_sequence", "password": "no_sequence"}
)
yield admin_config, reader_config, writer_config, no_sequence_config
@pytest.fixture
def store(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
_, reader_config, writer_config, _ = vector_configs
store = PgVectorStore(reader_config, writer_config, expected_dimension=2)
store.upsert("memory", [_record("reset", [0.0, 1.0])])
yield store
def test_pgvector_round_trip_hash_and_upsert(store):
assert store.upsert("memory", [_record("a", [1.0, 0.0])]) == 1
assert store.existing_hashes("memory", ["memory"])["record:a"] == "a"
hits = store.search(["memory"], [1.0, 0.0], limit=5, kinds=["memory"])
assert hits[0].metadata["content_hash"] == "a"
assert hits[0].id == "record:a"
assert store.upsert("memory", [_record("a", [0.8, 0.2])]) == 1
assert store.search(["memory"], [0.8, 0.2], limit=1)[0].id == "record:a"
def test_pgvector_lists_and_deletes_exact_evidence_generation(store):
generation = "gen:" + "a" * 32
value = VectorWriteRecord(
record=VectorRecord(
id="evidence-generation-a", kind="evidence", ref="doc:a", title="a",
content="content", metadata={"vector_generation": generation, "workspace_id": "default"},
),
embedding=[1.0, 0.0], content_hash="sha256:" + "a" * 64,
)
store.upsert("evidence", [value])
assert generation in store.list_evidence_generations("evidence", "default")
assert store.delete_generation("evidence", generation, "default") == 1
assert generation not in store.list_evidence_generations("evidence", "default")
def test_pgvector_generation_cleanup_isolated_between_workspaces(store):
generation = "gen:" + "b" * 32
records = [VectorWriteRecord(
record=VectorRecord(
id=f"evidence-{workspace}", kind="evidence", ref=f"doc:{workspace}",
title=workspace, content=workspace,
metadata={"vector_generation": generation, "workspace_id": workspace},
), embedding=[1.0, 0.0], content_hash="sha256:" + key * 64,
) for workspace, key in (("workspace-a", "b"), ("workspace-b", "c"))]
store.upsert("evidence", records)
assert store.delete_generation("evidence", generation, "workspace-a") == 1
assert generation not in store.list_evidence_generations("evidence", "workspace-a")
assert generation in store.list_evidence_generations("evidence", "workspace-b")
def test_pgvector_search_filters_kinds_before_limit(store):
store.upsert("memory", [_record("solved", [1.0, 0.0], kind="solved_question")])
hits = store.search("memory".split(), [1.0, 0.0], limit=1, kinds=["memory"])
assert len(hits) == 1
assert hits[0].kind == "memory"
def test_pgvector_multi_collection_search_skips_collections_unrelated_to_kinds(store):
store.upsert("evidence", [_record("evidence", [1.0, 0.0], kind="evidence")])
hits = store.search(["evidence", "memory"], [1.0, 0.0], limit=3, kinds=["memory"])
assert hits
assert {hit.kind for hit in hits} == {"memory"}
def test_pgvector_multi_collection_kind_filter_matches_http_adapter(store):
class Reader:
def search_similar(self, collection, embedding, limit, kinds=None):
if collection != "memory" or "memory" not in (kinds or []):
return []
return [
{
"similarity": 1.0,
"metadata": {
"record_key": "record:a",
"kind": "memory",
"ref": "session:test",
"title": "a",
"content": "content a",
"content_hash": "a",
},
}
]
direct = store.search(["evidence", "memory"], [1.0, 0.0], limit=1, kinds=["memory"])
http = ThothHttpVectorStore(Reader(), None).search(
["evidence", "memory"], [1.0, 0.0], limit=1, kinds=["memory"]
)
assert [(hit.id, hit.kind) for hit in direct] == [(hit.id, hit.kind) for hit in http]
def test_pgvector_search_rejects_unknown_kind_globally(store):
with pytest.raises(VectorStoreError, match="Kind not allowed"):
store.search(["memory"], [1.0, 0.0], limit=1, kinds=["unknown"])
@pytest.mark.parametrize("limit", [True, False, 1.0, 0, -1])
def test_pgvector_search_requires_strict_positive_limit(store, limit):
with pytest.raises(ValueError, match="positive integer"):
store.search(["memory"], [1.0, 0.0], limit=limit)
def test_pgvector_allowlists_collections(store):
with pytest.raises(VectorStoreError, match="Collection not allowed"):
store.search(["memory; DROP SCHEMA vectors"], [1.0, 0.0], limit=1)
with pytest.raises(VectorStoreError, match="Collection not allowed"):
store.upsert("unknown", [])
def test_pgvector_rejects_kinds_not_belonging_to_collection(store):
with pytest.raises(VectorStoreError, match="Kind not allowed"):
store.existing_hashes("evidence", ["memory"])
with pytest.raises(VectorStoreError, match="Kind not allowed"):
store.upsert("evidence", [_record("wrong", [1.0, 0.0])])
def test_pgvector_separates_read_and_write_credentials(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
_, reader_config, writer_config, _ = vector_configs
reader = PgVectorStore(reader_config, expected_dimension=2)
assert reader.capabilities.search is True
assert reader.capabilities.upsert is False
with pytest.raises(VectorWriteUnavailable):
reader.upsert("memory", [])
writer = PgVectorStore(None, writer_config, expected_dimension=2)
assert writer.capabilities.search is False
assert writer.capabilities.upsert is True
with pytest.raises(VectorReadUnavailable):
writer.search(["memory"], [1.0, 0.0], limit=1)
def test_pgvector_database_roles_are_least_privilege(vector_configs):
_, reader_config, writer_config, _ = vector_configs
reader_engine = create_engine(
f"postgresql+psycopg2://{reader_config.user}:{reader_config.password}"
f"@{reader_config.host}:{reader_config.port}/{reader_config.database}"
)
writer_engine = create_engine(
f"postgresql+psycopg2://{writer_config.user}:{writer_config.password}"
f"@{writer_config.host}:{writer_config.port}/{writer_config.database}"
)
with pytest.raises(ProgrammingError):
with reader_engine.begin() as connection:
connection.execute(
text(
"INSERT INTO vectors.memory "
"(record_key, kind, content_hash, metadata, embedding) "
"VALUES ('forbidden', 'memory', 'x', '{}', '[1,0]')"
)
)
with pytest.raises(ProgrammingError):
with writer_engine.connect() as connection:
connection.execute(
text(
"SELECT metadata, 1 - (embedding <=> '[1,0]'::vector) AS similarity "
"FROM vectors.memory ORDER BY embedding <=> '[1,0]'::vector LIMIT 1"
)
)
reader_engine.dispose()
writer_engine.dispose()
def test_pgvector_writer_health_requires_sequence_usage(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
admin_config, _, _, no_sequence_config = vector_configs
store = PgVectorStore(None, no_sequence_config, expected_dimension=2)
health = store.health()
assert health.ok is False
assert health.write_reachable is False
assert health.write_detail == (
"vector schema incomplete: missing sequence privileges evidence, memory, schema_records"
)
with pytest.raises(VectorWriteUnavailable) as error:
store.upsert("memory", [_record("needs-sequence", [1.0, 0.0])])
assert isinstance(error.value.__cause__, InsufficientPrivilege)
admin_engine = create_engine(
f"postgresql+psycopg2://{admin_config.user}:{admin_config.password}"
f"@{admin_config.host}:{admin_config.port}/{admin_config.database}"
)
with admin_engine.begin() as connection:
connection.exec_driver_sql(
"GRANT USAGE ON ALL SEQUENCES IN SCHEMA vectors TO vector_l0_no_sequence"
)
admin_engine.dispose()
assert store.health().ok is True
assert store.upsert("memory", [_record("has-sequence", [1.0, 0.0])]) == 1
def test_pgvector_health_requires_schema_usage_for_reader_and_writer(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
admin_config, reader_config, writer_config, _ = vector_configs
admin_engine = create_engine(
f"postgresql+psycopg2://{admin_config.user}:{admin_config.password}"
f"@{admin_config.host}:{admin_config.port}/{admin_config.database}"
)
store = PgVectorStore(reader_config, writer_config, expected_dimension=2)
with admin_engine.begin() as connection:
connection.exec_driver_sql(
f"REVOKE USAGE ON SCHEMA vectors FROM {reader_config.user}, {writer_config.user}"
)
health = store.health()
assert health.read_reachable is False and health.write_reachable is False
assert "missing schema usage" in health.read_detail
assert "missing schema usage" in health.write_detail
with pytest.raises(VectorReadUnavailable, match="Vector read operation unavailable"):
store.search(["memory"], [1.0, 0.0], limit=1)
with pytest.raises(VectorWriteUnavailable, match="Vector write operation unavailable"):
store.upsert("memory", [_record("blocked", [1.0, 0.0])])
with admin_engine.begin() as connection:
connection.exec_driver_sql(
f"GRANT USAGE ON SCHEMA vectors TO {reader_config.user}, {writer_config.user}"
)
admin_engine.dispose()
assert store.health().ok is True
def test_pgvector_maps_unavailable_connections_without_leaking_password(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
_, reader_config, writer_config, _ = vector_configs
password = "never-leak-this"
reader = reader_config.model_copy(update={"port": 1, "password": password})
writer = writer_config.model_copy(update={"port": 1, "password": password})
with pytest.raises(VectorReadUnavailable) as read_error:
PgVectorStore(reader, None).search(["memory"], [1.0, 0.0], limit=1)
with pytest.raises(VectorWriteUnavailable) as hash_error:
PgVectorStore(None, writer).existing_hashes("memory", ["memory"])
with pytest.raises(VectorWriteUnavailable) as write_error:
PgVectorStore(None, writer).upsert("memory", [_record("x", [1.0, 0.0])])
assert password not in str(read_error.value)
assert password not in str(hash_error.value)
assert password not in str(write_error.value)
def test_pgvector_health_reports_dimension_and_each_connection(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
_, reader_config, writer_config, _ = vector_configs
health = PgVectorStore(reader_config, writer_config, expected_dimension=2).health()
assert health.ok is True
assert health.read_reachable is True
assert health.write_reachable is True
assert health.observed_dimensions == (2,)
assert health.dimension_compatible is True
mismatch = PgVectorStore(reader_config, None, expected_dimension=3).health()
assert mismatch.ok is False
assert mismatch.read_reachable is False
assert mismatch.read_detail == (
"embedding dimension mismatch: evidence=2, memory=2, schema_records=2"
)
assert mismatch.dimension_compatible is False
def test_pgvector_health_rejects_clean_and_partial_schemas(vector_configs):
from tht.adapters.vector.pgvector import PgVectorStore
admin_config, _, _, _ = vector_configs
engine = create_engine(
f"postgresql+psycopg2://{admin_config.user}:{admin_config.password}"
f"@{admin_config.host}:{admin_config.port}/{admin_config.database}"
)
with engine.begin() as connection:
connection.exec_driver_sql("CREATE SCHEMA clean_vectors")
connection.exec_driver_sql("CREATE SCHEMA partial_vectors")
connection.exec_driver_sql(
"CREATE TABLE partial_vectors.memory "
"(record_key text, kind text, content_hash text, metadata jsonb)"
)
engine.dispose()
clean = PgVectorStore(
admin_config.model_copy(update={"db_schema": "clean_vectors"}),
expected_dimension=2,
).health()
assert clean.ok is False
assert clean.read_reachable is False
assert clean.read_detail == (
"vector schema incomplete: missing tables evidence, memory, schema_records"
)
partial = PgVectorStore(
admin_config.model_copy(update={"db_schema": "partial_vectors"}),
expected_dimension=2,
).health()
assert partial.ok is False
assert partial.read_reachable is False
assert partial.read_detail == (
"vector schema incomplete: missing tables evidence, schema_records; "
"missing embedding columns memory"
)