diff --git a/.superpowers/sdd/evidence-task-5c-report.md b/.superpowers/sdd/evidence-task-5c-report.md index 45f3c021..9bcd4282 100644 --- a/.superpowers/sdd/evidence-task-5c-report.md +++ b/.superpowers/sdd/evidence-task-5c-report.md @@ -102,3 +102,23 @@ JobRunner tests. Focused default/mixed/no-ACTIVE/search-pack tests pass, the real Docker pgvector lifecycle passes, and the final full harness plus scoped Ruff/diff invocation completed with exit code 0. + +## Workspace-scoped Evidence isolation + +- Evidence manifests, vector metadata, and record keys now carry the stable JobRunner workspace id + derived from the configured workspace identity (config stem), never credentials or absolute paths. +- Every ACTIVE server-side predicate includes `workspace_id`. Legacy unscoped rows therefore fail + closed and cannot appear in Evidence results. +- Vector generation inventory and deletion require the workspace namespace across the port, direct + pgvector adapter, HTTP client/adapter, and allowlisted RPC SQL. Legacy unscoped RPC overloads are + explicitly dropped during migration; destructive SQL matches collection, kind, generation, and + workspace together. +- GC recovers the persisted namespace from ACTIVE for explicit/restarted cleanup and can only list + or delete that workspace's generations. Real shared-pgvector coverage proves deleting a generation + for workspace A preserves the same generation in workspace B. +- `PipelineResult.model_dump` now serializes fields explicitly instead of `dataclasses.asdict`, + avoiding deepcopy of immutable `FrozenDict` metadata while preserving pristine JSON CLI output. + +Final focused verification: `89 passed` across corpus/CLI JSON, direct/HTTP parity, migrations, and +real Docker pgvector lifecycle; scoped Ruff and `git diff --check` clean. A contemporaneous full-suite +run reached unrelated Task 6 immutable-file tamper tests; those files were deliberately not changed. diff --git a/harness/scripts/create_vector_reader_rpc.sql b/harness/scripts/create_vector_reader_rpc.sql index 43c4ba31..08897921 100644 --- a/harness/scripts/create_vector_reader_rpc.sql +++ b/harness/scripts/create_vector_reader_rpc.sql @@ -43,7 +43,8 @@ begin table_name <> 'evidence' or not (metadata_filter ? 'vector_generation') or not (metadata_filter ? 'document_ids') - or jsonb_object_length(metadata_filter) <> 2 + or not (metadata_filter ? 'workspace_id') + or jsonb_object_length(metadata_filter) <> 3 or jsonb_typeof(metadata_filter->'document_ids') <> 'array' ) then raise exception 'invalid Evidence metadata filter'; @@ -56,6 +57,7 @@ begin where ($3 is null or t.kind = any ($3)) and ($4 is null or ( t.metadata->>''vector_generation'' = $4->>''vector_generation'' + and t.metadata->>''workspace_id'' = $4->>''workspace_id'' and t.metadata->>''document_id'' in ( select jsonb_array_elements_text($4->''document_ids'') ) diff --git a/harness/scripts/create_vector_writer_rpc.sql b/harness/scripts/create_vector_writer_rpc.sql index d006ac5e..538e215b 100644 --- a/harness/scripts/create_vector_writer_rpc.sql +++ b/harness/scripts/create_vector_writer_rpc.sql @@ -40,14 +40,18 @@ begin end; $$; -create or replace function public.list_evidence_generations(table_name text, kind text) +drop function if exists public.list_evidence_generations(text, text); + +create or replace function public.list_evidence_generations( + table_name text, kind text, workspace_id text +) returns table(generation text) language plpgsql security definer set search_path = public, vectors, extensions as $$ begin - if table_name <> 'evidence' or kind <> 'evidence' then + if table_name <> 'evidence' or kind <> 'evidence' or workspace_id !~ '^[a-z][a-z0-9_-]{0,63}$' then raise exception 'only exact Evidence generations may be listed'; end if; return query @@ -55,6 +59,7 @@ begin from vectors.evidence e where e.kind = 'evidence' and e.metadata->>'vector_generation' ~ '^gen:[0-9a-f]{32}$' + and e.metadata->>'workspace_id' = workspace_id order by 1; end; $$; @@ -120,8 +125,10 @@ begin end; $$; +drop function if exists public.delete_vector_generation(text, text, text); + create or replace function public.delete_vector_generation( - table_name text, kind text, generation text + table_name text, kind text, generation text, workspace_id text ) returns jsonb language plpgsql @@ -130,11 +137,13 @@ set search_path = public, vectors, extensions as $$ declare affected integer; begin - if table_name <> 'evidence' or kind <> 'evidence' or generation !~ '^gen:[0-9a-f]{32}$' then + if table_name <> 'evidence' or kind <> 'evidence' or generation !~ '^gen:[0-9a-f]{32}$' + or workspace_id !~ '^[a-z][a-z0-9_-]{0,63}$' then raise exception 'only an exact Evidence generation may be deleted'; end if; delete from vectors.evidence e - where e.kind = 'evidence' and e.metadata->>'vector_generation' = generation; + where e.kind = 'evidence' and e.metadata->>'vector_generation' = generation + and e.metadata->>'workspace_id' = workspace_id; get diagnostics affected = row_count; return jsonb_build_object('deleted', affected); end; @@ -143,8 +152,8 @@ $$; revoke all on function public._assert_vector_write_table(text, text[]) from public; revoke all on function public.existing_vector_hashes(text, text[]) from public; revoke all on function public.upsert_vector_records(text, jsonb) from public; -revoke all on function public.delete_vector_generation(text, text, text) from public; -revoke all on function public.list_evidence_generations(text, text) from public; +revoke all on function public.delete_vector_generation(text, text, text, text) from public; +revoke all on function public.list_evidence_generations(text, text, text) from public; -- Su alcuni progetti Supabase le funzioni in `public` ricevono grant automatici: revoca -- esplicitamente dai ruoli client generici, poi abilita solo il writer dedicato. @@ -153,20 +162,20 @@ begin if exists (select 1 from pg_roles where rolname = 'anon') then revoke all on function public.existing_vector_hashes(text, text[]) from anon; revoke all on function public.upsert_vector_records(text, jsonb) from anon; - revoke all on function public.delete_vector_generation(text, text, text) from anon; - revoke all on function public.list_evidence_generations(text, text) from anon; + revoke all on function public.delete_vector_generation(text, text, text, text) from anon; + revoke all on function public.list_evidence_generations(text, text, text) from anon; end if; if exists (select 1 from pg_roles where rolname = 'authenticated') then revoke all on function public.existing_vector_hashes(text, text[]) from authenticated; revoke all on function public.upsert_vector_records(text, jsonb) from authenticated; - revoke all on function public.delete_vector_generation(text, text, text) from authenticated; - revoke all on function public.list_evidence_generations(text, text) from authenticated; + revoke all on function public.delete_vector_generation(text, text, text, text) from authenticated; + revoke all on function public.list_evidence_generations(text, text, text) from authenticated; end if; if exists (select 1 from pg_roles where rolname = 'vector_writer') then grant execute on function public.existing_vector_hashes(text, text[]) to vector_writer; grant execute on function public.upsert_vector_records(text, jsonb) to vector_writer; - grant execute on function public.delete_vector_generation(text, text, text) to vector_writer; - grant execute on function public.list_evidence_generations(text, text) to vector_writer; + grant execute on function public.delete_vector_generation(text, text, text, text) to vector_writer; + grant execute on function public.list_evidence_generations(text, text, text) to vector_writer; end if; end $$; diff --git a/harness/tests/l0/test_pgvector_corpus_lifecycle.py b/harness/tests/l0/test_pgvector_corpus_lifecycle.py index bc988a32..9a06bc29 100644 --- a/harness/tests/l0/test_pgvector_corpus_lifecycle.py +++ b/harness/tests/l0/test_pgvector_corpus_lifecycle.py @@ -241,7 +241,8 @@ def test_real_pgvector_corpus_job_lifecycle(tmp_path, persistent_pgvector): record=VectorRecord( id="chunk:exact-vector-orphan", kind="evidence", ref="doc:orphan", title="orphan", content="orphan", - metadata={"document_id": "doc:orphan", "vector_generation": orphan}, + metadata={"document_id": "doc:orphan", "vector_generation": orphan, + "workspace_id": "pgvector-lifecycle"}, ), embedding=query, content_hash="sha256:" + "f" * 64, )]) @@ -251,5 +252,7 @@ def test_real_pgvector_corpus_job_lifecycle(tmp_path, persistent_pgvector): expected_fs = set(generations[-2:]) assert set(CorpusStore(tmp_path / "corpus").list_generations()) == expected_fs expected_vectors = expected_fs | {generations[0]} - assert set(recreated.list_evidence_generations("evidence")) == expected_vectors + assert set(recreated.list_evidence_generations( + "evidence", "pgvector-lifecycle" + )) == expected_vectors assert final_pipeline.gc(workspace_root=tmp_path)["evicted"] == [] diff --git a/harness/tests/l0/test_pgvector_store.py b/harness/tests/l0/test_pgvector_store.py index d01b3a3a..83c77b66 100644 --- a/harness/tests/l0/test_pgvector_store.py +++ b/harness/tests/l0/test_pgvector_store.py @@ -130,14 +130,29 @@ def test_pgvector_lists_and_deletes_exact_evidence_generation(store): value = VectorWriteRecord( record=VectorRecord( id="evidence-generation-a", kind="evidence", ref="doc:a", title="a", - content="content", metadata={"vector_generation": generation}, + 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") - assert store.delete_generation("evidence", generation) == 1 - assert generation not in store.list_evidence_generations("evidence") + 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): diff --git a/harness/tests/l0/test_vector_adapter_parity.py b/harness/tests/l0/test_vector_adapter_parity.py index 41fe3a93..ec5899ad 100644 --- a/harness/tests/l0/test_vector_adapter_parity.py +++ b/harness/tests/l0/test_vector_adapter_parity.py @@ -235,9 +235,10 @@ def test_http_delete_generation_uses_exact_allowlisted_rpc_payload(monkeypatch): lambda url, json, **kwargs: calls.append((url, json)) or Response({"deleted": 2}), ) client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) - assert client.delete_generation("evidence", "gen:" + "a" * 32) == 2 + assert client.delete_generation("evidence", "gen:" + "a" * 32, "default") == 2 assert calls == [("https://vectors.test/rpc/delete_vector_generation", { "table_name": "evidence", "kind": "evidence", "generation": "gen:" + "a" * 32, + "workspace_id": "default", })] @@ -248,7 +249,7 @@ def test_http_delete_generation_legacy_404_fails_closed_without_body_leak(monkey ) client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) with pytest.raises(VectorRestError, match="delete_vector_generation RPC is unavailable") as error: - client.delete_generation("evidence", "gen:" + "a" * 32) + client.delete_generation("evidence", "gen:" + "a" * 32, "default") assert "secret" not in str(error.value) @@ -261,9 +262,9 @@ def test_http_list_evidence_generations_exact_rpc_and_legacy_fail_closed(monkeyp ]), ) client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) - assert client.list_evidence_generations("evidence") == ["gen:" + "a" * 32] + assert client.list_evidence_generations("evidence", "default") == ["gen:" + "a" * 32] assert calls[0][0].endswith("/rpc/list_evidence_generations") - assert calls[0][1] == {"table_name": "evidence", "kind": "evidence"} + assert calls[0][1] == {"table_name": "evidence", "kind": "evidence", "workspace_id": "default"} @pytest.mark.parametrize("generation", ["gen:a", "gen:" + "A" * 32, "gen:" + "a" * 33]) @@ -274,7 +275,7 @@ def test_http_generation_operations_reject_noncanonical_values(monkeypatch, gene ) client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) with pytest.raises(ValueError, match="canonical"): - client.delete_generation("evidence", generation) + client.delete_generation("evidence", generation, "default") def test_http_inventory_rejects_malformed_rpc_output(monkeypatch): @@ -284,4 +285,4 @@ def test_http_inventory_rejects_malformed_rpc_output(monkeypatch): ) client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) with pytest.raises(VectorRestError, match="malformed"): - client.list_evidence_generations("evidence") + client.list_evidence_generations("evidence", "default") diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 92fb12f1..e4d067bb 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -1,7 +1,7 @@ import pytest from tht.corpus.chunk import ChunkPolicy -from tht.corpus.pipeline import CorpusPipeline, PipelineError +from tht.corpus.pipeline import CorpusPipeline, PipelineError, PipelineResult from tht.corpus.store import CorpusStore from tht.corpus.models import CorpusManifest from tht.ports.evidence import AcquiredDocument, SourceObject @@ -55,17 +55,19 @@ class Vectors: value.record.id: value.content_hash for value in self.records } - def delete_generation(self, collection, generation): + def delete_generation(self, collection, generation, workspace_id): self.records = [ value for value in self.records - if value.record.metadata["vector_generation"] != generation + if not (value.record.metadata["vector_generation"] == generation + and value.record.metadata.get("workspace_id") == workspace_id) ] return 0 - def list_evidence_generations(self, collection): + def list_evidence_generations(self, collection, workspace_id): return sorted({ value.record.metadata["vector_generation"] for value in self.records if value.record.kind == "evidence" + and value.record.metadata.get("workspace_id") == workspace_id }) @@ -128,10 +130,10 @@ def test_retention_keeps_filesystem_when_vector_purge_fails_then_retries(tmp_pat super().__init__() self.fail_delete = True - def delete_generation(self, collection, generation): + def delete_generation(self, collection, generation, workspace_id): if self.fail_delete: raise RuntimeError("credential secret") - return super().delete_generation(collection, generation) + return super().delete_generation(collection, generation, workspace_id) vectors = FailingDelete() for index in range(2): @@ -156,13 +158,13 @@ def test_gc_reconciles_vector_only_generation(tmp_path): from tht.vectorstore.records import VectorRecord vectors.records.append(VectorWriteRecord( record=VectorRecord(id="orphan", kind="evidence", ref="doc:x", title="", content="x", - metadata={"vector_generation": orphan}), + metadata={"vector_generation": orphan, "workspace_id": "default"}), embedding=[0.0, 0.0, 0.0], content_hash="sha256:" + "0" * 64, )) candidate = pipeline(tmp_path, Source([]), vectors=vectors, retain=1) report = candidate.gc(workspace_root=tmp_path) assert report["evicted"] == [orphan] - assert vectors.list_evidence_generations("evidence") == [] + assert vectors.list_evidence_generations("evidence", "default") == [] assert candidate.gc(workspace_root=tmp_path)["evicted"] == [] @@ -233,7 +235,7 @@ def test_gc_preserves_vector_dependencies_of_retained_manifests(tmp_path): vectors=vectors, retain=2, ).run().generation assert CorpusStore(tmp_path / "corpus").list_generations() == [second, third] - assert first in vectors.list_evidence_generations("evidence") + assert first in vectors.list_evidence_generations("evidence", "default") def test_active_searcher_without_active_fails_closed_for_evidence(tmp_path): @@ -319,6 +321,14 @@ def test_active_evidence_query_holds_lock_against_publish(tmp_path): assert published.is_set() +def test_pipeline_result_dump_does_not_deepcopy_frozen_metadata(): + manifest = CorpusManifest(metadata={"nested": {"value": ["safe"]}}) + payload = PipelineResult( + "succeeded", None, False, (), (), (), manifest, + ).model_dump(mode="json") + assert payload["manifest"]["metadata"] == {"nested": {"value": ["safe"]}} + + def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path): one = item("one", "a") first_source = Source([(one, "hello")]) diff --git a/harness/tht/adapters/vector/pgvector.py b/harness/tht/adapters/vector/pgvector.py index 756277df..d0db7495 100644 --- a/harness/tht/adapters/vector/pgvector.py +++ b/harness/tht/adapters/vector/pgvector.py @@ -273,16 +273,18 @@ class PgVectorStore: filter_params.append(collection_kinds) if metadata_filter is not None: if collection != "evidence" or set(metadata_filter) != { - "vector_generation", "document_ids" + "vector_generation", "document_ids", "workspace_id" }: raise VectorStoreError("Unsupported vector metadata filter") generation = metadata_filter["vector_generation"] document_ids = metadata_filter["document_ids"] - if not isinstance(generation, str) or not isinstance(document_ids, list): + workspace_id = metadata_filter["workspace_id"] + if not isinstance(generation, str) or not isinstance(document_ids, list) or not isinstance(workspace_id, str): raise VectorStoreError("Invalid vector metadata filter") clauses.append(sql.SQL("metadata->>'vector_generation' = %s")) clauses.append(sql.SQL("metadata->>'document_id' = ANY(%s)")) - filter_params.extend((generation, document_ids)) + clauses.append(sql.SQL("metadata->>'workspace_id' = %s")) + filter_params.extend((generation, document_ids, workspace_id)) where = ( sql.SQL(" WHERE ") + sql.SQL(" AND ").join(clauses) if clauses else sql.SQL("") @@ -404,9 +406,11 @@ class PgVectorStore: raw.close() return len(records) - def delete_generation(self, collection: str, generation: str) -> int: + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: if collection != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None: raise VectorStoreError("Only exact Evidence generations may be deleted") + if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None: + raise VectorStoreError("Invalid Evidence workspace namespace") raw = None try: raw = self._require_writer().raw_connection() @@ -414,9 +418,10 @@ class PgVectorStore: cursor.execute( sql.SQL( "DELETE FROM {} WHERE kind = 'evidence' " - "AND metadata->>'vector_generation' = %s" + "AND metadata->>'vector_generation' = %s " + "AND metadata->>'workspace_id' = %s" ).format(_collection(self._schema, collection)), - (generation,), + (generation, workspace_id), ) count = cursor.rowcount raw.commit() @@ -429,9 +434,11 @@ class PgVectorStore: if raw is not None: raw.close() - def list_evidence_generations(self, collection: str) -> list[str]: + def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]: if collection != "evidence": raise VectorStoreError("Only exact Evidence generations may be listed") + if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None: + raise VectorStoreError("Invalid Evidence workspace namespace") raw = None try: raw = self._require_writer().raw_connection() @@ -440,8 +447,9 @@ class PgVectorStore: sql.SQL( "SELECT DISTINCT metadata->>'vector_generation' FROM {} " "WHERE kind = 'evidence' AND metadata->>'vector_generation' " - "~ '^gen:[0-9a-f]{{32}}$' ORDER BY 1" - ).format(_collection(self._schema, collection)) + "~ '^gen:[0-9a-f]{{32}}$' AND metadata->>'workspace_id' = %s ORDER BY 1" + ).format(_collection(self._schema, collection)), + (workspace_id,), ) return [row[0] for row in cursor.fetchall()] except Exception as exc: diff --git a/harness/tht/adapters/vector/thoth_http.py b/harness/tht/adapters/vector/thoth_http.py index 4b5276b1..5a93b3c3 100644 --- a/harness/tht/adapters/vector/thoth_http.py +++ b/harness/tht/adapters/vector/thoth_http.py @@ -155,19 +155,19 @@ class ThothHttpVectorStore: except VectorRestError as exc: raise VectorStoreError(str(exc)) from exc - def delete_generation(self, collection: str, generation: str) -> int: + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: if collection != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None: raise VectorStoreError("Only exact Evidence generations may be deleted") try: - return self._require_writer().delete_generation(collection, generation) + return self._require_writer().delete_generation(collection, generation, workspace_id) except VectorRestError as exc: raise VectorStoreError(str(exc)) from exc - def list_evidence_generations(self, collection: str) -> list[str]: + def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]: if collection != "evidence": raise VectorStoreError("Only exact Evidence generations may be listed") try: - return self._require_writer().list_evidence_generations(collection) + return self._require_writer().list_evidence_generations(collection, workspace_id) except VectorRestError as exc: raise VectorWriteUnavailable("Vector generation inventory unavailable") from exc diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index 4c068a7c..0a516889 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -47,9 +47,17 @@ class PipelineResult: resumed_from: str | None = None def model_dump(self, mode=None): - value = asdict(self) - value["manifest"] = self.manifest.model_dump(mode="json") - return value + return { + "status": self.status, + "generation": self.generation, + "published": self.published, + "changed": list(self.changed), + "unchanged": list(self.unchanged), + "removed": list(self.removed), + "manifest": self.manifest.model_dump(mode="json"), + "run_id": self.run_id, + "resumed_from": self.resumed_from, + } def _fingerprint(value) -> str: @@ -62,6 +70,7 @@ class CorpusPipeline: self, *, store: CorpusStore, sources: list[EvidenceSource], embedder, vector_store: VectorStore, embedding_model: str, embedding_dimensions: int, chunk_policy: ChunkPolicy, pipeline_version: str, retain_published_generations: int = 3, + workspace_id: str = "default", ) -> None: self.store = store self.sources = sources @@ -74,6 +83,7 @@ class CorpusPipeline: if isinstance(retain_published_generations, bool) or retain_published_generations < 1: raise ValueError("retain_published_generations must be at least 1") self.retain_published_generations = retain_published_generations + self.workspace_id = workspace_id def _protected_generations(self, workspace_root: Path) -> set[str]: protected = {value for value in (self.store.active_generation(),) if value} @@ -92,9 +102,14 @@ class CorpusPipeline: return protected def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict: + active_manifest = self.store.active_manifest() + if active_manifest is not None: + persisted_workspace = active_manifest.metadata.get("workspace_id") + if isinstance(persisted_workspace, str): + self.workspace_id = persisted_workspace published = self.store.published_generations() list_vectors = getattr(self.vector_store, "list_evidence_generations", None) - vector_generations = set(list_vectors("evidence")) if list_vectors else set() + vector_generations = set(list_vectors("evidence", self.workspace_id)) if list_vectors else set() generations = sorted(set(published) | vector_generations) job_protected = self._protected_generations(workspace_root) active = self.store.active_generation() @@ -124,7 +139,7 @@ class CorpusPipeline: continue if purge_vector: try: - self.vector_store.delete_generation("evidence", generation) + self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: failures.append({"generation": generation, "error": "vector cleanup failed"}) continue @@ -154,6 +169,7 @@ class CorpusPipeline: return self._run(dry_run=dry_run, resume=resume) def run_as_job(self, **kwargs) -> PipelineResult: + self.workspace_id = kwargs["workspace_id"] with self.store.writer_lock(): return self._run_as_job(**kwargs) @@ -259,6 +275,7 @@ class CorpusPipeline: vector_generation=plan["generation"], documents=tuple(documents), chunks=tuple(chunks), metadata={ + "workspace_id": self.workspace_id, "compatibility_fingerprint": compatibility, "fingerprints": plan["fingerprints"], "removed": plan["removed"], @@ -289,7 +306,7 @@ class CorpusPipeline: changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]} parts = [part for part in manifest.chunks if part.document_id in changed_docs] embeddings = read(context, "embeddings.json") - return [self._vector_record(part, vector, plan["generation"]) + return [self._vector_record(part, vector, plan["generation"], self.workspace_id) for part, vector in zip(parts, embeddings, strict=True)] def compensate(context: JobContext) -> None: @@ -297,7 +314,7 @@ class CorpusPipeline: if self.store.active_generation() != generation: self.store.discard(generation) try: - self.vector_store.delete_generation("evidence", generation) + self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: pass write(context, "compensated.json", {"generation": generation}) @@ -495,6 +512,7 @@ class CorpusPipeline: vector_generation=generation, documents=tuple(documents), chunks=tuple(chunks), metadata={ + "workspace_id": self.workspace_id, "compatibility_fingerprint": compatibility, "fingerprints": fingerprints, "removed": list(removed), @@ -508,7 +526,7 @@ class CorpusPipeline: raise PipelineError("embedding count mismatch") if any(len(vector) != self.embedding_dimensions for vector in embeddings): raise PipelineError("embedding dimension mismatch") - records = [self._vector_record(part, vector, generation) for part, vector in zip(changed_chunks, embeddings, strict=True)] + records = [self._vector_record(part, vector, generation, self.workspace_id) for part, vector in zip(changed_chunks, embeddings, strict=True)] if records: written = self.vector_store.upsert("evidence", records) vector_written = True @@ -547,17 +565,21 @@ class CorpusPipeline: pass if vector_written: try: - self.vector_store.delete_generation("evidence", generation) + self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: pass @staticmethod - def _vector_record(chunk: CanonicalChunk, embedding: list[float], generation: str): + def _vector_record( + chunk: CanonicalChunk, embedding: list[float], generation: str, workspace_id: str, + ): record = VectorRecord( - id=f"{generation}:{chunk.chunk_id}", kind="evidence", ref=chunk.document_id, + id=f"{workspace_id}:{generation}:{chunk.chunk_id}", + kind="evidence", ref=chunk.document_id, title=str(chunk.metadata.get("title", "")), content=chunk.content, metadata={ **dict(chunk.metadata), "document_id": chunk.document_id, + "workspace_id": workspace_id, "source_uri": chunk.source_uri, "ordinal": chunk.ordinal, "vector_generation": generation, }, diff --git a/harness/tht/ports/vector.py b/harness/tht/ports/vector.py index 57da32a8..9fd1cd8d 100644 --- a/harness/tht/ports/vector.py +++ b/harness/tht/ports/vector.py @@ -80,9 +80,9 @@ class VectorStore(Protocol): def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int: ... - def delete_generation(self, collection: str, generation: str) -> int: ... + def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int: ... - def list_evidence_generations(self, collection: str) -> list[str]: ... + def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]: ... __all__ = [ diff --git a/harness/tht/search/evidence.py b/harness/tht/search/evidence.py index a2d24f93..47ff2718 100644 --- a/harness/tht/search/evidence.py +++ b/harness/tht/search/evidence.py @@ -30,7 +30,8 @@ class ActiveEvidenceSearcher: hits.extend(self.delegate.search(embedding, **kwargs)) if include_evidence: manifest = self.corpus.active_manifest() - if manifest is not None: + workspace_id = manifest.metadata.get("workspace_id") if manifest else None + if manifest is not None and isinstance(workspace_id, str): by_generation: dict[str, list[str]] = {} mapping = dict(manifest.metadata.get("document_generations", {})) for document in manifest.documents: @@ -43,6 +44,7 @@ class ActiveEvidenceSearcher: metadata_filter={ "vector_generation": generation, "document_ids": sorted(document_ids), + "workspace_id": workspace_id, }, )) return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:top_n] diff --git a/harness/tht/vectorstore/rest_client.py b/harness/tht/vectorstore/rest_client.py index 544e2e0d..dd7e1ac9 100644 --- a/harness/tht/vectorstore/rest_client.py +++ b/harness/tht/vectorstore/rest_client.py @@ -128,13 +128,16 @@ class VectorRestClient: return len(payload) return len(rows) - def delete_generation(self, table_name: str, generation: str) -> int: + def delete_generation(self, table_name: str, generation: str, workspace_id: str) -> int: if table_name != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None: raise ValueError("generation must be canonical") + if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None: + raise ValueError("workspace namespace must be canonical") try: payload = self._call( "delete_vector_generation", - {"table_name": table_name, "kind": "evidence", "generation": generation}, + {"table_name": table_name, "kind": "evidence", "generation": generation, + "workspace_id": workspace_id}, ) except VectorRestError as error: if "HTTP 404" in str(error): @@ -146,11 +149,13 @@ class VectorRestClient: return int(payload.get("deleted", 0)) return 0 - def list_evidence_generations(self, table_name: str) -> list[str]: + def list_evidence_generations(self, table_name: str, workspace_id: str) -> list[str]: + if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None: + raise ValueError("workspace namespace must be canonical") try: rows = self._call( "list_evidence_generations", - {"table_name": table_name, "kind": "evidence"}, + {"table_name": table_name, "kind": "evidence", "workspace_id": workspace_id}, ) or [] except VectorRestError as error: if "HTTP 404" in str(error):