From efcb0deb31213824a4dafda90645e7fe0cdcad24 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 06:07:30 +0200 Subject: [PATCH] fix(evidence): lock GC and preflight search ownership --- .superpowers/sdd/evidence-task-5c-report.md | 7 +++++++ harness/tests/test_corpus_pipeline.py | 3 +-- harness/tht/cli/preprocess_cmd.py | 3 +-- harness/tht/cli/search_cmd.py | 12 ++++++++++-- harness/tht/corpus/pipeline.py | 4 ++++ harness/tht/corpus/store.py | 12 ++++++++++++ harness/tht/search/evidence.py | 20 ++++++++++++++++++++ 7 files changed, 55 insertions(+), 6 deletions(-) diff --git a/.superpowers/sdd/evidence-task-5c-report.md b/.superpowers/sdd/evidence-task-5c-report.md index d6c72e5f..b0c12cae 100644 --- a/.superpowers/sdd/evidence-task-5c-report.md +++ b/.superpowers/sdd/evidence-task-5c-report.md @@ -148,3 +148,10 @@ every configured search delegate, including default, mixed, pack, and non-Eviden 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. + +Final lock/preflight follow-up: `CorpusPipeline.gc()` now acquires the corpus writer lock itself for +ownership validation through vector/filesystem cleanup. The store lock is thread-reentrant so nested +job retention is safe without weakening cross-thread/process exclusion; the CLI wrapper no longer +double-locks. Search find/pack performs locked corpus ownership preflight immediately after config +load, before DWH leasing, vector/searcher factories, embeddings, or schema work. Focused concurrency +and fail-closed tests, real Docker lifecycle, scoped Ruff/diff, and the full harness pass. diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 12d963f6..fa9e2357 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -208,8 +208,7 @@ def test_explicit_gc_blocks_while_job_holds_corpus_writer_lock(tmp_path): assert entered.wait(5) def collect(): - with candidate.store.writer_lock(): - candidate.gc(workspace_root=tmp_path) + candidate.gc(workspace_root=tmp_path) gc_finished.set() gc_thread = threading.Thread(target=collect) diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 92676154..5c94ee95 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -114,8 +114,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False): 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) + return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run) @preprocess_app.command("evidence") diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index 7500e25a..2b100542 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -57,13 +57,17 @@ def search_cmd( from tht.search import combined_search cfg = _load_config_or_exit(config) + from tht.search.evidence import validate_corpus_workspace + + workspace_id = config.stem.lower().replace(".", "-").replace("_", "-") + validate_corpus_workspace(cfg, workspace_id) dwh_snapshot = _leased_dwh_snapshot(cfg, ctx) require_vector_cfg(cfg) from tht.search.evidence import active_searcher runtime_searcher = active_searcher( cfg, open_searcher(cfg), - workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"), + workspace_id=workspace_id, ) if kind is not None and kind not in KIND_MAP: typer.secho( @@ -256,6 +260,10 @@ def pack_cmd( from tht.vectorstore.rest_client import VectorRestError cfg = _load_config_or_exit(config) + from tht.search.evidence import validate_corpus_workspace + + workspace_id = config.stem.lower().replace(".", "-").replace("_", "-") + validate_corpus_workspace(cfg, workspace_id) dwh_snapshot = _leased_dwh_snapshot(cfg, ctx) require_vector_cfg(cfg) @@ -272,7 +280,7 @@ def pack_cmd( searcher = active_searcher( cfg, open_searcher(cfg), - workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"), + workspace_id=workspace_id, ) embedder = make_embedder(cfg.embeddings) vec = embedder.embed_query(question) diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index d8ae0e62..c5a66649 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -124,6 +124,10 @@ class CorpusPipeline: return protected def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict: + with self.store.writer_lock(): + return self._gc(workspace_root=workspace_root, dry_run=dry_run) + + def _gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict: self._assert_workspace_binding() published = self.store.published_generations() list_vectors = getattr(self.vector_store, "list_evidence_generations", None) diff --git a/harness/tht/corpus/store.py b/harness/tht/corpus/store.py index dfcd65c3..0202e566 100644 --- a/harness/tht/corpus/store.py +++ b/harness/tht/corpus/store.py @@ -10,6 +10,7 @@ import stat import shutil import uuid import hashlib +import threading from datetime import UTC, datetime from pathlib import Path from contextlib import contextmanager @@ -49,6 +50,7 @@ class CorpusStore: self.active_path = self.root / "ACTIVE" self._replace = os.replace self._fsync_directory = self._sync_root + self._lock_state = threading.local() self._ensure_root() def _ensure_root(self) -> None: @@ -61,6 +63,14 @@ class CorpusStore: @contextmanager def writer_lock(self): + depth = getattr(self._lock_state, "depth", 0) + if depth: + self._lock_state.depth = depth + 1 + try: + yield + finally: + self._lock_state.depth -= 1 + return lock_path = self.root / ".writer.lock" fd = os.open(lock_path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC, 0o600) try: @@ -68,8 +78,10 @@ class CorpusStore: if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or info.st_nlink != 1: raise UnsafeCorpusPath("corpus writer lock is unsafe") fcntl.flock(fd, fcntl.LOCK_EX) + self._lock_state.depth = 1 yield finally: + self._lock_state.depth = 0 fcntl.flock(fd, fcntl.LOCK_UN) os.close(fd) diff --git a/harness/tht/search/evidence.py b/harness/tht/search/evidence.py index fc2c5cff..077b3355 100644 --- a/harness/tht/search/evidence.py +++ b/harness/tht/search/evidence.py @@ -76,6 +76,26 @@ def active_searcher(cfg, delegate, *, workspace_id: str | None = None): return ActiveEvidenceSearcher(CorpusStore(corpus_root), delegate, workspace_id) +def validate_corpus_workspace(cfg, workspace_id: str) -> None: + """Fail before downstream retrieval setup when configured corpus ownership differs.""" + corpus = CorpusStore(cfg.paths.artifacts.parent / "corpus") + with corpus.writer_lock(): + manifest = corpus.active_manifest() + if manifest is None: + return + 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 CorpusWorkspaceMismatchError( + "corpus workspace ownership is missing or invalid; use a new corpus root or rebuild" + ) + if persisted != workspace_id: + raise CorpusWorkspaceMismatchError( + "corpus belongs to a different workspace; use a new corpus root or rebuild" + ) + + def resolve_evidence_file( store: CorpusStore, evidence_id: str, *, materialized_root=None, ) -> str: