fix(evidence): bind corpus to workspace identity
This commit is contained in:
@@ -122,3 +122,16 @@ and the final full harness plus scoped Ruff/diff invocation completed with exit
|
||||
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.
|
||||
|
||||
## Immutable corpus/workspace binding
|
||||
|
||||
- A corpus root becomes bound to the workspace id persisted in its ACTIVE manifest. Job, non-job,
|
||||
explicit GC, and ACTIVE search entry points compare the configured namespace before discovery,
|
||||
vector access, staging, deletion, or ACTIVE mutation.
|
||||
- Reusing the same paths after renaming a workspace now fails closed with a typed/sanitized message:
|
||||
use a new corpus root or perform an intentional explicit rebuild. Unscoped legacy manifests also
|
||||
fail this ownership check.
|
||||
- Tests prove unchanged-document reuse cannot silently mix workspace A vectors into a workspace B
|
||||
manifest, and that mismatched job, GC, and search paths perform no vector/filesystem mutations.
|
||||
|
||||
Focused workspace-binding, search-pack, preprocess JSON, and scoped Ruff/diff tests pass.
|
||||
|
||||
@@ -329,6 +329,66 @@ def test_pipeline_result_dump_does_not_deepcopy_frozen_metadata():
|
||||
assert payload["manifest"]["metadata"] == {"nested": {"value": ["safe"]}}
|
||||
|
||||
|
||||
def test_reused_corpus_root_rejects_workspace_rename_before_any_mutation(tmp_path):
|
||||
vectors = Vectors()
|
||||
first = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
|
||||
first.run_as_job(
|
||||
workspace_id="workspace-a", workspace_root=tmp_path,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + "2" * 64,
|
||||
)
|
||||
active = first.store.active_generation()
|
||||
records = list(vectors.records)
|
||||
renamed_source = Source([(item("one", "a"), "stable")])
|
||||
renamed = pipeline(tmp_path, renamed_source, vectors=vectors)
|
||||
with pytest.raises(PipelineError, match="different workspace"):
|
||||
renamed.run_as_job(
|
||||
workspace_id="workspace-b", workspace_root=tmp_path,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + "2" * 64,
|
||||
)
|
||||
assert renamed_source.acquire_calls == []
|
||||
assert renamed.store.active_generation() == active
|
||||
assert vectors.records == records
|
||||
|
||||
|
||||
def test_gc_rejects_workspace_mismatch_without_deleting(tmp_path):
|
||||
vectors = Vectors()
|
||||
owner = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
|
||||
owner.run_as_job(
|
||||
workspace_id="workspace-a", workspace_root=tmp_path,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + "2" * 64,
|
||||
)
|
||||
generations = owner.store.list_generations()
|
||||
wrong = pipeline(tmp_path, Source([]), vectors=vectors)
|
||||
wrong.workspace_id = "workspace-b"
|
||||
with pytest.raises(PipelineError, match="different workspace"):
|
||||
wrong.gc(workspace_root=tmp_path)
|
||||
assert wrong.store.list_generations() == generations
|
||||
|
||||
|
||||
def test_active_search_rejects_workspace_mismatch_before_delegate(tmp_path):
|
||||
from tht.search.evidence import ActiveEvidenceSearcher, CorpusWorkspaceMismatchError
|
||||
|
||||
vectors = Vectors()
|
||||
owner = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
|
||||
owner.run_as_job(
|
||||
workspace_id="workspace-a", workspace_root=tmp_path,
|
||||
config_fingerprint="sha256:" + "1" * 64,
|
||||
input_fingerprint="sha256:" + "2" * 64,
|
||||
)
|
||||
|
||||
class Delegate:
|
||||
def search(self, *args, **kwargs):
|
||||
raise AssertionError("workspace mismatch reached vector delegate")
|
||||
|
||||
with pytest.raises(CorpusWorkspaceMismatchError, match="different workspace"):
|
||||
ActiveEvidenceSearcher(
|
||||
owner.store, Delegate(), expected_workspace_id="workspace-b",
|
||||
).search([1.0], kinds=["evidence"])
|
||||
|
||||
|
||||
def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path):
|
||||
one = item("one", "a")
|
||||
first_source = Source([(one, "hello")])
|
||||
|
||||
@@ -113,6 +113,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False):
|
||||
pipeline_version="evidence-v1",
|
||||
retain_published_generations=cfg.vector.retain_published_generations,
|
||||
)
|
||||
pipeline.workspace_id = config.stem.lower().replace(".", "-").replace("_", "-")
|
||||
with pipeline.store.writer_lock():
|
||||
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
|
||||
|
||||
|
||||
@@ -53,7 +53,10 @@ def search_cmd(
|
||||
require_vector_cfg(cfg)
|
||||
from tht.search.evidence import active_searcher
|
||||
|
||||
runtime_searcher = active_searcher(cfg, open_searcher(cfg))
|
||||
runtime_searcher = active_searcher(
|
||||
cfg, open_searcher(cfg),
|
||||
workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"),
|
||||
)
|
||||
if kind is not None and kind not in KIND_MAP:
|
||||
typer.secho(
|
||||
f"ERRORE: --kind sconosciuto: {kind} (validi: {', '.join(KIND_MAP)})",
|
||||
@@ -260,7 +263,10 @@ def pack_cmd(
|
||||
try:
|
||||
from tht.search.evidence import active_searcher
|
||||
|
||||
searcher = active_searcher(cfg, open_searcher(cfg))
|
||||
searcher = active_searcher(
|
||||
cfg, open_searcher(cfg),
|
||||
workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"),
|
||||
)
|
||||
embedder = make_embedder(cfg.embeddings)
|
||||
vec = embedder.embed_query(question)
|
||||
except degrade as e:
|
||||
|
||||
@@ -85,6 +85,16 @@ class CorpusPipeline:
|
||||
self.retain_published_generations = retain_published_generations
|
||||
self.workspace_id = workspace_id
|
||||
|
||||
def _assert_workspace_binding(self) -> None:
|
||||
manifest = self.store.active_manifest()
|
||||
if manifest is None:
|
||||
return
|
||||
persisted = manifest.metadata.get("workspace_id")
|
||||
if persisted != self.workspace_id:
|
||||
raise PipelineError(
|
||||
"corpus belongs to a different workspace; use a new corpus root or explicit rebuild"
|
||||
)
|
||||
|
||||
def _protected_generations(self, workspace_root: Path) -> set[str]:
|
||||
protected = {value for value in (self.store.active_generation(),) if value}
|
||||
runs = workspace_root / ".tht-jobs" / "evidence" / "runs"
|
||||
@@ -102,11 +112,7 @@ 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
|
||||
self._assert_workspace_binding()
|
||||
published = self.store.published_generations()
|
||||
list_vectors = getattr(self.vector_store, "list_evidence_generations", None)
|
||||
vector_generations = set(list_vectors("evidence", self.workspace_id)) if list_vectors else set()
|
||||
@@ -166,11 +172,13 @@ class CorpusPipeline:
|
||||
|
||||
def run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult:
|
||||
with self.store.writer_lock():
|
||||
self._assert_workspace_binding()
|
||||
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():
|
||||
self._assert_workspace_binding()
|
||||
return self._run_as_job(**kwargs)
|
||||
|
||||
def _run_as_job(
|
||||
|
||||
@@ -3,12 +3,17 @@
|
||||
from tht.corpus.store import CorpusStore
|
||||
|
||||
|
||||
class CorpusWorkspaceMismatchError(RuntimeError):
|
||||
"""The configured workspace does not own the persisted corpus."""
|
||||
|
||||
|
||||
class ActiveEvidenceSearcher:
|
||||
"""Searcher facade that enforces ACTIVE generation predicates before LIMIT."""
|
||||
|
||||
def __init__(self, corpus: CorpusStore, delegate):
|
||||
def __init__(self, corpus: CorpusStore, delegate, expected_workspace_id: str | None = None):
|
||||
self.corpus = corpus
|
||||
self.delegate = delegate
|
||||
self.expected_workspace_id = expected_workspace_id
|
||||
|
||||
def search(self, embedding, top_n=10, kinds=None, metadata_filter=None):
|
||||
requested = set(kinds) if kinds is not None else {
|
||||
@@ -31,6 +36,12 @@ class ActiveEvidenceSearcher:
|
||||
if include_evidence:
|
||||
manifest = self.corpus.active_manifest()
|
||||
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):
|
||||
by_generation: dict[str, list[str]] = {}
|
||||
mapping = dict(manifest.metadata.get("document_generations", {}))
|
||||
@@ -50,9 +61,9 @@ class ActiveEvidenceSearcher:
|
||||
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:top_n]
|
||||
|
||||
|
||||
def active_searcher(cfg, delegate):
|
||||
def active_searcher(cfg, delegate, *, workspace_id: str | None = None):
|
||||
corpus_root = cfg.paths.artifacts.parent / "corpus"
|
||||
return ActiveEvidenceSearcher(CorpusStore(corpus_root), delegate)
|
||||
return ActiveEvidenceSearcher(CorpusStore(corpus_root), delegate, workspace_id)
|
||||
|
||||
|
||||
def resolve_evidence_file(
|
||||
|
||||
Reference in New Issue
Block a user