fix(evidence): lock GC and preflight search ownership
This commit is contained in:
@@ -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
|
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,
|
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.
|
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.
|
||||||
|
|||||||
@@ -208,8 +208,7 @@ def test_explicit_gc_blocks_while_job_holds_corpus_writer_lock(tmp_path):
|
|||||||
assert entered.wait(5)
|
assert entered.wait(5)
|
||||||
|
|
||||||
def collect():
|
def collect():
|
||||||
with candidate.store.writer_lock():
|
candidate.gc(workspace_root=tmp_path)
|
||||||
candidate.gc(workspace_root=tmp_path)
|
|
||||||
gc_finished.set()
|
gc_finished.set()
|
||||||
|
|
||||||
gc_thread = threading.Thread(target=collect)
|
gc_thread = threading.Thread(target=collect)
|
||||||
|
|||||||
@@ -114,8 +114,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False):
|
|||||||
retain_published_generations=cfg.vector.retain_published_generations,
|
retain_published_generations=cfg.vector.retain_published_generations,
|
||||||
)
|
)
|
||||||
pipeline.workspace_id = config.stem.lower().replace(".", "-").replace("_", "-")
|
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")
|
@preprocess_app.command("evidence")
|
||||||
|
|||||||
@@ -57,13 +57,17 @@ def search_cmd(
|
|||||||
from tht.search import combined_search
|
from tht.search import combined_search
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
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)
|
dwh_snapshot = _leased_dwh_snapshot(cfg, ctx)
|
||||||
require_vector_cfg(cfg)
|
require_vector_cfg(cfg)
|
||||||
from tht.search.evidence import active_searcher
|
from tht.search.evidence import active_searcher
|
||||||
|
|
||||||
runtime_searcher = active_searcher(
|
runtime_searcher = active_searcher(
|
||||||
cfg, open_searcher(cfg),
|
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:
|
if kind is not None and kind not in KIND_MAP:
|
||||||
typer.secho(
|
typer.secho(
|
||||||
@@ -256,6 +260,10 @@ def pack_cmd(
|
|||||||
from tht.vectorstore.rest_client import VectorRestError
|
from tht.vectorstore.rest_client import VectorRestError
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
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)
|
dwh_snapshot = _leased_dwh_snapshot(cfg, ctx)
|
||||||
require_vector_cfg(cfg)
|
require_vector_cfg(cfg)
|
||||||
|
|
||||||
@@ -272,7 +280,7 @@ def pack_cmd(
|
|||||||
|
|
||||||
searcher = active_searcher(
|
searcher = active_searcher(
|
||||||
cfg, open_searcher(cfg),
|
cfg, open_searcher(cfg),
|
||||||
workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"),
|
workspace_id=workspace_id,
|
||||||
)
|
)
|
||||||
embedder = make_embedder(cfg.embeddings)
|
embedder = make_embedder(cfg.embeddings)
|
||||||
vec = embedder.embed_query(question)
|
vec = embedder.embed_query(question)
|
||||||
|
|||||||
@@ -124,6 +124,10 @@ class CorpusPipeline:
|
|||||||
return protected
|
return protected
|
||||||
|
|
||||||
def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict:
|
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()
|
self._assert_workspace_binding()
|
||||||
published = self.store.published_generations()
|
published = self.store.published_generations()
|
||||||
list_vectors = getattr(self.vector_store, "list_evidence_generations", None)
|
list_vectors = getattr(self.vector_store, "list_evidence_generations", None)
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import stat
|
|||||||
import shutil
|
import shutil
|
||||||
import uuid
|
import uuid
|
||||||
import hashlib
|
import hashlib
|
||||||
|
import threading
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
@@ -49,6 +50,7 @@ class CorpusStore:
|
|||||||
self.active_path = self.root / "ACTIVE"
|
self.active_path = self.root / "ACTIVE"
|
||||||
self._replace = os.replace
|
self._replace = os.replace
|
||||||
self._fsync_directory = self._sync_root
|
self._fsync_directory = self._sync_root
|
||||||
|
self._lock_state = threading.local()
|
||||||
self._ensure_root()
|
self._ensure_root()
|
||||||
|
|
||||||
def _ensure_root(self) -> None:
|
def _ensure_root(self) -> None:
|
||||||
@@ -61,6 +63,14 @@ class CorpusStore:
|
|||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
def writer_lock(self):
|
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"
|
lock_path = self.root / ".writer.lock"
|
||||||
fd = os.open(lock_path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC, 0o600)
|
fd = os.open(lock_path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC, 0o600)
|
||||||
try:
|
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:
|
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")
|
raise UnsafeCorpusPath("corpus writer lock is unsafe")
|
||||||
fcntl.flock(fd, fcntl.LOCK_EX)
|
fcntl.flock(fd, fcntl.LOCK_EX)
|
||||||
|
self._lock_state.depth = 1
|
||||||
yield
|
yield
|
||||||
finally:
|
finally:
|
||||||
|
self._lock_state.depth = 0
|
||||||
fcntl.flock(fd, fcntl.LOCK_UN)
|
fcntl.flock(fd, fcntl.LOCK_UN)
|
||||||
os.close(fd)
|
os.close(fd)
|
||||||
|
|
||||||
|
|||||||
@@ -76,6 +76,26 @@ def active_searcher(cfg, delegate, *, workspace_id: str | None = None):
|
|||||||
return ActiveEvidenceSearcher(CorpusStore(corpus_root), delegate, workspace_id)
|
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(
|
def resolve_evidence_file(
|
||||||
store: CorpusStore, evidence_id: str, *, materialized_root=None,
|
store: CorpusStore, evidence_id: str, *, materialized_root=None,
|
||||||
) -> str:
|
) -> str:
|
||||||
|
|||||||
Reference in New Issue
Block a user