test: gate real pgvector corpus lifecycle
This commit is contained in:
@@ -0,0 +1,42 @@
|
||||
# Evidence Task 5D — Real pgvector lifecycle gate
|
||||
|
||||
## Status
|
||||
|
||||
Complete. The Docker-backed L0 gate uses one persistent `pgvector/pgvector:pg16`
|
||||
database and the production migrations, direct reader/writer `PgVectorStore`,
|
||||
`CorpusStore`, `CorpusPipeline.run_as_job`/JobRunner, ACTIVE Evidence retrieval,
|
||||
search-pack fusion, owned session artifact copy, retention, and explicit GC.
|
||||
|
||||
## Lifecycle covered
|
||||
|
||||
- Four real corpus publications with retention set to two generations.
|
||||
- A higher-similarity stale vector proves ACTIVE metadata filtering happens before LIMIT
|
||||
for normal Evidence retrieval and the search-pack fusion path.
|
||||
- A removed source is absent from ACTIVE retrieval and cannot be copied to a session.
|
||||
- An injected process death occurs after one real committed vector upsert. Resume uses the
|
||||
real run ID, preserves that record, fills the missing records, and produces no duplicate keys.
|
||||
- Database engines and direct store objects are disposed/recreated before persisted ACTIVE
|
||||
retrieval is checked again.
|
||||
- An exact canonical vector-only orphan generation is discovered and removed by explicit GC.
|
||||
- Filesystem and vector inventories converge exactly to ACTIVE plus one rollback; a second GC
|
||||
is a no-op.
|
||||
- Owned session artifact bytes and SHA-256 match the ACTIVE canonical document.
|
||||
|
||||
## Production bug found and fixed
|
||||
|
||||
Production migration `003_roles.sql` intentionally restricted `vector_writer`, but omitted
|
||||
the privileges used by the production generation lifecycle: `SELECT(metadata)` for inventory
|
||||
and `DELETE` for cleanup on `vectors.evidence`. Consequently a real job published successfully
|
||||
and then failed in `retention_cleanup` on its first run.
|
||||
|
||||
Added versioned migration `004_evidence_generation_gc.sql` granting only those two Evidence
|
||||
generation-management privileges. Runtime application code was not redesigned.
|
||||
|
||||
## Verification
|
||||
|
||||
- Target lifecycle: `1 passed` (Docker-backed).
|
||||
- Full harness: `681 passed, 5 deselected`.
|
||||
- Scoped Ruff: passed.
|
||||
- `git diff --check`: passed.
|
||||
|
||||
The existing Pydantic serialization and legacy-workspace deprecation warnings remain unchanged.
|
||||
@@ -0,0 +1,236 @@
|
||||
"""L0 gate for the complete durable Evidence/pgvector lifecycle."""
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
|
||||
import pytest
|
||||
from sqlalchemy import create_engine, text
|
||||
from testcontainers.postgres import PostgresContainer
|
||||
|
||||
from tht.adapters.evidence import FilesystemEvidenceSource
|
||||
from tht.adapters.vector.pgvector import PgVectorStore
|
||||
from tht.cli.vector_migrate_cmd import migrate
|
||||
from tht.config import DatabaseConfig
|
||||
from tht.corpus.chunk import ChunkPolicy
|
||||
from tht.corpus.pipeline import CorpusPipeline
|
||||
from tht.corpus.store import CorpusStore
|
||||
from tht.ports.vector import VectorRecord, VectorWriteRecord
|
||||
from tht.search import combined_search
|
||||
from tht.search.evidence import ActiveEvidenceSearcher, resolve_evidence_file
|
||||
|
||||
|
||||
DIMENSIONS = 768
|
||||
|
||||
|
||||
class DeterministicEmbedder:
|
||||
def embed_documents(self, texts):
|
||||
return [self.embed_query(text) for text in texts]
|
||||
|
||||
def embed_query(self, text):
|
||||
vector = [0.0] * DIMENSIONS
|
||||
vector[0] = 0.8
|
||||
vector[1] = 0.6
|
||||
return vector
|
||||
|
||||
|
||||
class EvidenceDelegate:
|
||||
"""Adapt the real multi-collection port to the runtime search protocol."""
|
||||
|
||||
def __init__(self, store):
|
||||
self.store = store
|
||||
|
||||
def search(self, embedding, top_n=10, kinds=None, metadata_filter=None):
|
||||
return self.store.search(
|
||||
["evidence"], embedding, limit=top_n, kinds=kinds,
|
||||
metadata_filter=metadata_filter,
|
||||
)
|
||||
|
||||
|
||||
class InterruptAfterRealPartialUpsert:
|
||||
"""Crash after a committed real row, as a process death would."""
|
||||
|
||||
def __init__(self, store):
|
||||
self.store = store
|
||||
self.interrupt = True
|
||||
|
||||
def __getattr__(self, name):
|
||||
return getattr(self.store, name)
|
||||
|
||||
def upsert(self, collection, records):
|
||||
if self.interrupt and len(records) > 1:
|
||||
self.interrupt = False
|
||||
self.store.upsert(collection, records[:1])
|
||||
raise KeyboardInterrupt("injected process death after committed vector row")
|
||||
return self.store.upsert(collection, records)
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def persistent_pgvector():
|
||||
with PostgresContainer("pgvector/pgvector:pg16") as postgres:
|
||||
migrate(postgres.get_connection_url())
|
||||
admin = create_engine(postgres.get_connection_url())
|
||||
with admin.begin() as connection:
|
||||
connection.exec_driver_sql(
|
||||
"ALTER ROLE vector_reader LOGIN PASSWORD 'reader-lifecycle'"
|
||||
)
|
||||
connection.exec_driver_sql(
|
||||
"ALTER ROLE vector_writer LOGIN PASSWORD 'writer-lifecycle'"
|
||||
)
|
||||
url = admin.url
|
||||
common = dict(
|
||||
host=url.host, port=url.port, database=url.database, schema="vectors"
|
||||
)
|
||||
reader = DatabaseConfig(
|
||||
**common, user="vector_reader", password="reader-lifecycle"
|
||||
)
|
||||
writer = DatabaseConfig(
|
||||
**common, user="vector_writer", password="writer-lifecycle"
|
||||
)
|
||||
yield postgres, admin, reader, writer
|
||||
admin.dispose()
|
||||
|
||||
|
||||
def _pipeline(root, source_root, vectors):
|
||||
return CorpusPipeline(
|
||||
store=CorpusStore(root / "corpus"),
|
||||
sources=[FilesystemEvidenceSource(source_root)],
|
||||
embedder=DeterministicEmbedder(),
|
||||
vector_store=vectors,
|
||||
embedding_model="deterministic-l0",
|
||||
embedding_dimensions=DIMENSIONS,
|
||||
chunk_policy=ChunkPolicy(version="lifecycle-v1", max_chars=48),
|
||||
pipeline_version="evidence-v1",
|
||||
retain_published_generations=2,
|
||||
)
|
||||
|
||||
|
||||
def _publish(pipeline, root, serial):
|
||||
return pipeline.run_as_job(
|
||||
workspace_id="pgvector-lifecycle",
|
||||
workspace_root=root,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + f"{serial:x}" * 64,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.l0
|
||||
def test_real_pgvector_corpus_job_lifecycle(tmp_path, persistent_pgvector):
|
||||
postgres, admin, reader_config, writer_config = persistent_pgvector
|
||||
source_root = tmp_path / "sources"
|
||||
source_root.mkdir()
|
||||
kept = source_root / "kept.md"
|
||||
removed = source_root / "removed.md"
|
||||
removed.write_text("removed evidence generation zero", encoding="utf-8")
|
||||
|
||||
vectors = PgVectorStore(reader_config, writer_config, expected_dimension=DIMENSIONS)
|
||||
generations = []
|
||||
for serial in range(3):
|
||||
kept.write_text(f"active evidence generation {serial}", encoding="utf-8")
|
||||
result = _publish(_pipeline(tmp_path, source_root, vectors), tmp_path, serial + 1)
|
||||
assert result.status == "succeeded"
|
||||
generations.append(result.generation)
|
||||
|
||||
# A stale, closer row must not consume LIMIT before ACTIVE filtering.
|
||||
stale_generation = generations[-2]
|
||||
stale = VectorWriteRecord(
|
||||
record=VectorRecord(
|
||||
id="chunk:stale-perfect-match", kind="evidence", ref="doc:stale",
|
||||
title="stale forbidden", content="stale forbidden",
|
||||
metadata={
|
||||
"document_id": "doc:stale",
|
||||
"vector_generation": stale_generation,
|
||||
},
|
||||
),
|
||||
embedding=[1.0] + [0.0] * (DIMENSIONS - 1),
|
||||
content_hash="sha256:" + "a" * 64,
|
||||
)
|
||||
vectors.upsert("evidence", [stale])
|
||||
runtime = ActiveEvidenceSearcher(CorpusStore(tmp_path / "corpus"), EvidenceDelegate(vectors))
|
||||
query = DeterministicEmbedder().embed_query("active")
|
||||
hits = runtime.search(query, top_n=1, kinds=["evidence"])
|
||||
assert len(hits) == 1 and hits[0].title != "stale forbidden"
|
||||
packed = combined_search(
|
||||
"active", lsh_hits=None, store=runtime, embedder=DeterministicEmbedder(),
|
||||
top=1, rrf_k=60, kinds=["evidence"], query_vec=query,
|
||||
)
|
||||
assert len(packed) == 1 and packed[0].label != "stale forbidden"
|
||||
|
||||
# Fourth publication removes a document and creates multiple chunks for crash recovery.
|
||||
removed.unlink()
|
||||
kept.write_text("active fourth generation " * 8, encoding="utf-8")
|
||||
crashing = InterruptAfterRealPartialUpsert(vectors)
|
||||
candidate = _pipeline(tmp_path, source_root, crashing)
|
||||
with pytest.raises(KeyboardInterrupt, match="injected process death"):
|
||||
_publish(candidate, tmp_path, 4)
|
||||
runs = tmp_path / ".tht-jobs" / "evidence" / "runs"
|
||||
crashed_run = max(runs.iterdir(), key=lambda path: path.stat().st_mtime_ns).name
|
||||
before = vectors.existing_hashes("evidence", ["evidence"])
|
||||
intent = json.loads(
|
||||
(runs / crashed_run / "artifacts" / "vector-intent.json").read_text()
|
||||
)["records"]
|
||||
already_present = set(intent) & set(before)
|
||||
assert len(already_present) == 1
|
||||
resumed = candidate.run_as_job(
|
||||
workspace_id="pgvector-lifecycle", workspace_root=tmp_path,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + "4" * 64,
|
||||
resume_run_id=crashed_run,
|
||||
)
|
||||
assert resumed.status == "succeeded" and resumed.resumed_from == crashed_run
|
||||
generations.append(resumed.generation)
|
||||
after = vectors.existing_hashes("evidence", ["evidence"])
|
||||
assert {key: after[key] for key in already_present} == {
|
||||
key: before[key] for key in already_present
|
||||
}
|
||||
assert set(intent).issubset(after)
|
||||
with admin.connect() as connection:
|
||||
duplicate_count = connection.execute(text(
|
||||
"SELECT count(*) - count(DISTINCT record_key) FROM vectors.evidence"
|
||||
)).scalar_one()
|
||||
assert duplicate_count == 0
|
||||
|
||||
runtime = ActiveEvidenceSearcher(CorpusStore(tmp_path / "corpus"), EvidenceDelegate(vectors))
|
||||
assert all(hit.ref != "doc:stale" for hit in runtime.search(query, top_n=20, kinds=["evidence"]))
|
||||
manifest = CorpusStore(tmp_path / "corpus").active_manifest()
|
||||
removed_document = next(
|
||||
doc for doc in _pipeline(tmp_path, source_root, vectors).store.manifest(generations[-2]).documents
|
||||
if "removed.md" in doc.source_uri
|
||||
)
|
||||
assert resolve_evidence_file(
|
||||
CorpusStore(tmp_path / "corpus"), removed_document.document_id,
|
||||
materialized_root=tmp_path / "session",
|
||||
) == ""
|
||||
|
||||
active_document = manifest.documents[0]
|
||||
owned = CorpusStore(tmp_path / "corpus").materialize_document(
|
||||
active_document.document_id, tmp_path / "session" / "active-evidence.md"
|
||||
)
|
||||
assert owned.read_bytes() == active_document.content.encode()
|
||||
assert hashlib.sha256(owned.read_bytes()).hexdigest() == active_document.content_hash[7:]
|
||||
|
||||
# Recreate engines and stores against the same persisted database.
|
||||
vectors._reader.dispose()
|
||||
vectors._writer.dispose()
|
||||
recreated = PgVectorStore(reader_config, writer_config, expected_dimension=DIMENSIONS)
|
||||
assert recreated.health().ok is True
|
||||
recreated_hits = ActiveEvidenceSearcher(
|
||||
CorpusStore(tmp_path / "corpus"), EvidenceDelegate(recreated)
|
||||
).search(query, top_n=2, kinds=["evidence"])
|
||||
assert recreated_hits and all(hit.metadata["vector_generation"] == resumed.generation for hit in recreated_hits)
|
||||
|
||||
orphan = "gen:" + "f" * 32
|
||||
recreated.upsert("evidence", [VectorWriteRecord(
|
||||
record=VectorRecord(
|
||||
id="chunk:exact-vector-orphan", kind="evidence", ref="doc:orphan",
|
||||
title="orphan", content="orphan",
|
||||
metadata={"document_id": "doc:orphan", "vector_generation": orphan},
|
||||
),
|
||||
embedding=query, content_hash="sha256:" + "f" * 64,
|
||||
)])
|
||||
final_pipeline = _pipeline(tmp_path, source_root, recreated)
|
||||
report = final_pipeline.gc(workspace_root=tmp_path)
|
||||
assert report["evicted"] == [orphan]
|
||||
expected = set(generations[-2:])
|
||||
assert set(CorpusStore(tmp_path / "corpus").list_generations()) == expected
|
||||
assert set(recreated.list_evidence_generations("evidence")) == expected
|
||||
assert final_pipeline.gc(workspace_root=tmp_path)["evicted"] == []
|
||||
@@ -22,7 +22,7 @@ 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"]
|
||||
assert [item.version for item in before.pending] == ["001", "002", "003", "004"]
|
||||
|
||||
migrate(database_url)
|
||||
migrate(database_url)
|
||||
@@ -30,7 +30,7 @@ def test_migrations_are_clean_and_idempotent(database_url):
|
||||
status = migration_status(database_url)
|
||||
assert status.pending == ()
|
||||
assert status.drifted == ()
|
||||
assert [item.version for item in status.applied] == ["001", "002", "003"]
|
||||
assert [item.version for item in status.applied] == ["001", "002", "003", "004"]
|
||||
|
||||
|
||||
def test_schema_matches_direct_adapter_contract(database_url):
|
||||
@@ -156,7 +156,7 @@ def test_status_json_is_pristine(database_url, monkeypatch):
|
||||
|
||||
assert result.exit_code == 0, result.output
|
||||
assert json.loads(result.stdout) == {
|
||||
"applied": ["001", "002", "003"],
|
||||
"applied": ["001", "002", "003", "004"],
|
||||
"drifted": [],
|
||||
"pending": [],
|
||||
}
|
||||
@@ -249,7 +249,7 @@ def test_hostile_admin_search_path_cannot_shadow_migration_objects(database_url)
|
||||
with verification.connect() as connection:
|
||||
assert connection.execute(
|
||||
text("SELECT count(*) FROM public.tht_vector_migrations")
|
||||
).scalar_one() == 3
|
||||
).scalar_one() == 4
|
||||
assert connection.execute(
|
||||
text("SELECT count(*) FROM shadow.tht_vector_migrations")
|
||||
).scalar_one() == 0
|
||||
@@ -279,10 +279,10 @@ 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 = _copy_migrations(tmp_path)
|
||||
(migrations / "004_first.sql").write_text("CREATE TABLE public.must_rollback (id int);\n")
|
||||
(migrations / "005_broken.sql").write_text("THIS IS NOT SQL;\n")
|
||||
(migrations / "005_first.sql").write_text("CREATE TABLE public.must_rollback (id int);\n")
|
||||
(migrations / "006_broken.sql").write_text("THIS IS NOT SQL;\n")
|
||||
|
||||
with pytest.raises(MigrationError, match="005_broken.sql"):
|
||||
with pytest.raises(MigrationError, match="006_broken.sql"):
|
||||
migrate(database_url, migrations)
|
||||
|
||||
engine = create_engine(database_url)
|
||||
@@ -290,8 +290,8 @@ def test_failed_batch_rolls_back_schema_and_ledger(database_url, tmp_path):
|
||||
assert connection.execute(text("SELECT to_regclass('public.must_rollback')")).scalar() is None
|
||||
engine.dispose()
|
||||
status = migration_status(database_url, migrations)
|
||||
assert [item.version for item in status.applied] == ["001", "002", "003"]
|
||||
assert [item.version for item in status.pending] == ["004", "005"]
|
||||
assert [item.version for item in status.applied] == ["001", "002", "003", "004"]
|
||||
assert [item.version for item in status.pending] == ["005", "006"]
|
||||
|
||||
|
||||
def _copy_migrations(tmp_path: Path) -> Path:
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
-- The writer owns derived-generation reconciliation but not runtime similarity reads.
|
||||
GRANT SELECT (metadata) ON vectors.evidence TO vector_writer;
|
||||
GRANT DELETE ON vectors.evidence TO vector_writer;
|
||||
Reference in New Issue
Block a user