fix(evidence): validate ownership before all searches

This commit is contained in:
2026-07-12 06:04:00 +02:00
parent c6966f3d15
commit d21aac9151
4 changed files with 68 additions and 15 deletions
@@ -142,3 +142,9 @@ ACTIVE owner (or `default` only for a brand-new direct corpus), preserving safe
the real pgvector lifecycle. Explicit config/job identities still fail closed on any mismatch. The the real pgvector lifecycle. Explicit config/job identities still fail closed on any mismatch. The
two reported regressions, workspace mismatch guards, real Docker lifecycle, scoped Ruff/diff, and two reported regressions, workspace mismatch guards, real Docker lifecycle, scoped Ruff/diff, and
the full harness suite all pass. the full harness suite all pass.
Final fail-closed follow-up: persisted ACTIVE ownership is now validated under the corpus lock before
every configured search delegate, including default, mixed, pack, and non-Evidence-only operations.
Malformed or missing `metadata.workspace_id` is intrinsically rejected even for unbound direct
callers; source discovery, vector operations, GC, files, and ACTIVE remain untouched. Focused tests,
real Docker lifecycle, scoped Ruff/diff, and the full harness regression run pass.
+33 -3
View File
@@ -257,7 +257,10 @@ def test_active_searcher_splits_default_and_mixed_kinds_before_global_limit(tmp_
from tht.search.evidence import ActiveEvidenceSearcher from tht.search.evidence import ActiveEvidenceSearcher
store = CorpusStore(tmp_path / "corpus") store = CorpusStore(tmp_path / "corpus")
generation = store.stage(CorpusManifest(), {}, generation="gen:" + "a" * 32) generation = store.stage(
CorpusManifest(metadata={"workspace_id": "default"}), {},
generation="gen:" + "a" * 32,
)
store.publish(generation) store.publish(generation)
calls = [] calls = []
@@ -368,7 +371,8 @@ def test_gc_rejects_workspace_mismatch_without_deleting(tmp_path):
assert wrong.store.list_generations() == generations assert wrong.store.list_generations() == generations
def test_active_search_rejects_workspace_mismatch_before_delegate(tmp_path): @pytest.mark.parametrize("kinds", [None, ["evidence", "memory"], ["memory"]])
def test_active_search_rejects_workspace_mismatch_before_delegate(tmp_path, kinds):
from tht.search.evidence import ActiveEvidenceSearcher, CorpusWorkspaceMismatchError from tht.search.evidence import ActiveEvidenceSearcher, CorpusWorkspaceMismatchError
vectors = Vectors() vectors = Vectors()
@@ -386,7 +390,33 @@ def test_active_search_rejects_workspace_mismatch_before_delegate(tmp_path):
with pytest.raises(CorpusWorkspaceMismatchError, match="different workspace"): with pytest.raises(CorpusWorkspaceMismatchError, match="different workspace"):
ActiveEvidenceSearcher( ActiveEvidenceSearcher(
owner.store, Delegate(), expected_workspace_id="workspace-b", owner.store, Delegate(), expected_workspace_id="workspace-b",
).search([1.0], kinds=["evidence"]) ).search([1.0], kinds=kinds)
def test_unscoped_active_manifest_is_never_adopted_by_direct_run_or_gc(tmp_path):
store = CorpusStore(tmp_path / "corpus")
generation = store.stage(CorpusManifest(), {})
store.publish(generation)
class ForbiddenSource:
def discover(self):
raise AssertionError("invalid corpus reached source discovery")
class ForbiddenVectors:
def __getattr__(self, name):
raise AssertionError(f"invalid corpus reached vector operation {name}")
candidate = CorpusPipeline(
store=store, sources=[ForbiddenSource()], embedder=Embedder(),
vector_store=ForbiddenVectors(), embedding_model="model", embedding_dimensions=3,
chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=100),
pipeline_version="evidence-v1",
)
with pytest.raises(PipelineError, match="missing or invalid"):
candidate.run()
with pytest.raises(PipelineError, match="missing or invalid"):
candidate.gc(workspace_root=tmp_path)
assert store.active_generation() == generation
def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path): def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path):
+7
View File
@@ -4,6 +4,7 @@ from __future__ import annotations
import hashlib import hashlib
import json import json
import re
import uuid import uuid
from dataclasses import asdict, dataclass from dataclasses import asdict, dataclass
from pathlib import Path from pathlib import Path
@@ -92,6 +93,12 @@ class CorpusPipeline:
self.workspace_id = "default" self.workspace_id = "default"
return return
persisted = manifest.metadata.get("workspace_id") persisted = manifest.metadata.get("workspace_id")
if not isinstance(persisted, str) or re.fullmatch(
r"[a-z][a-z0-9_-]{0,63}", persisted
) is None:
raise PipelineError(
"corpus workspace ownership is missing or invalid; use a new corpus root or explicit rebuild"
)
if self.workspace_id is None and isinstance(persisted, str): if self.workspace_id is None and isinstance(persisted, str):
self.workspace_id = persisted self.workspace_id = persisted
return return
+22 -12
View File
@@ -1,5 +1,7 @@
"""Runtime Evidence lookup bound to the atomically active corpus generation.""" """Runtime Evidence lookup bound to the atomically active corpus generation."""
import re
from tht.corpus.store import CorpusStore from tht.corpus.store import CorpusStore
@@ -21,12 +23,27 @@ class ActiveEvidenceSearcher:
} }
include_evidence = "evidence" in requested include_evidence = "evidence" in requested
other_kinds = sorted(requested - {"evidence"}) other_kinds = sorted(requested - {"evidence"})
if not include_evidence:
kwargs = {"top_n": top_n, "kinds": kinds}
if metadata_filter is not None:
kwargs["metadata_filter"] = metadata_filter
return self.delegate.search(embedding, **kwargs)
with self.corpus.writer_lock(): with self.corpus.writer_lock():
manifest = self.corpus.active_manifest()
persisted_workspace = manifest.metadata.get("workspace_id") if manifest else None
if manifest is not None and (
not isinstance(persisted_workspace, str)
or re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", persisted_workspace) is None
):
raise CorpusWorkspaceMismatchError(
"corpus workspace ownership is missing or invalid; use a new corpus root or rebuild"
)
if manifest is not None and self.expected_workspace_id is not None and (
persisted_workspace != self.expected_workspace_id
):
raise CorpusWorkspaceMismatchError(
"corpus belongs to a different workspace; use a new corpus root or rebuild"
)
if not include_evidence:
kwargs = {"top_n": top_n, "kinds": kinds}
if metadata_filter is not None:
kwargs["metadata_filter"] = metadata_filter
return self.delegate.search(embedding, **kwargs)
hits = [] hits = []
if other_kinds: if other_kinds:
kwargs = {"top_n": top_n, "kinds": other_kinds} kwargs = {"top_n": top_n, "kinds": other_kinds}
@@ -34,14 +51,7 @@ class ActiveEvidenceSearcher:
kwargs["metadata_filter"] = metadata_filter kwargs["metadata_filter"] = metadata_filter
hits.extend(self.delegate.search(embedding, **kwargs)) hits.extend(self.delegate.search(embedding, **kwargs))
if include_evidence: if include_evidence:
manifest = self.corpus.active_manifest()
workspace_id = manifest.metadata.get("workspace_id") if manifest else None workspace_id = manifest.metadata.get("workspace_id") if manifest else None
if manifest is not None and self.expected_workspace_id is not None and (
workspace_id != self.expected_workspace_id
):
raise CorpusWorkspaceMismatchError(
"corpus belongs to a different workspace; use a new corpus root or rebuild"
)
if manifest is not None and isinstance(workspace_id, str): if manifest is not None and isinstance(workspace_id, str):
by_generation: dict[str, list[str]] = {} by_generation: dict[str, list[str]] = {}
mapping = dict(manifest.metadata.get("document_generations", {})) mapping = dict(manifest.metadata.get("document_generations", {}))