feat: add durable addressed registry publication

This commit is contained in:
2026-08-11 13:39:38 +02:00
parent 475e94b662
commit ee87a59e1f
24 changed files with 380 additions and 55 deletions
@@ -0,0 +1,16 @@
from __future__ import annotations
import os, stat
import pytest
from tht.workspace_writer_lock import WorkspaceWriterConflict, verify_workspace_writer_fds
def test_verifier_rejects_missing_capability():
with pytest.raises(WorkspaceWriterConflict): verify_workspace_writer_fds(env={})
def test_verifier_checks_fd_identity(tmp_path):
root=tmp_path / "root"; root.mkdir(mode=0o700); writer=root / "writer.lock"; writer.touch(mode=0o600)
r=os.open(root, os.O_RDONLY); w=os.open(writer, os.O_RDWR)
try:
st=os.fstat(r); env={"THOTH_WORKSPACE_ID":"abc-workspace", "THOTH_WORKSPACE_REVISION":"a"*40, "THOTH_WORKSPACE_DEVICE":str(st.st_dev), "THOTH_WORKSPACE_INODE":str(st.st_ino)}
cap=verify_workspace_writer_fds(writer_fd=w, root_fd=r, env=env)
assert cap.inode == st.st_ino
finally: os.close(w); os.close(r)
+9
View File
@@ -20,6 +20,13 @@ from tht.vectorstore.embeddings import EmbeddingsError
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
logger = logging.getLogger(__name__)
def _require_writer_capability() -> None:
"""Mutating children opt into the backend-owned fd capability contract."""
import os
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1":
from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability()
_PREPROCESS_EXPECTED_ERRORS = (
OSError, RuntimeError, ValueError, TypeError, KeyError,
EvidenceSourceError, VectorStoreError, EmbeddingsError,
@@ -29,6 +36,7 @@ _PREPROCESS_EXPECTED_ERRORS = (
def run_dwh_from_config(
config: Path, *, steps: tuple[str, ...], resume: str | None = None,
):
_require_writer_capability()
from tht.cli.lsh_cmd import build_lsh_artifacts
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
from tht.jobs.dwh_pipeline import (
@@ -74,6 +82,7 @@ def _parse_dwh_steps(value: str) -> tuple[str, ...]:
def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None):
_require_writer_capability()
from tht.adapters.factory import build_evidence_sources, build_vector_store
from tht.cli.vector_cmd import make_embedder
from tht.corpus.chunk import ChunkPolicy
+7
View File
@@ -16,6 +16,12 @@ from tht.mschema.eligibility import classify_all
schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)")
logger = logging.getLogger(__name__)
def _require_writer_capability() -> None:
import os
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1":
from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability()
def _add_examples(dwh, phys, examples) -> None:
for table_name, table in phys.tables.items():
@@ -54,6 +60,7 @@ def annotations_path(cfg: Config) -> Path:
def refresh_catalog(cfg, *, dwh=None, output_path: Path | None = None):
_require_writer_capability()
"""Run the existing catalog algorithm and persist its canonical output."""
target = dwh if dwh is not None else build_dwh(cfg)
physical = target.introspect()
+7
View File
@@ -24,6 +24,12 @@ from tht.vectorstore.store import SyncStats, content_hash
vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)")
logger = logging.getLogger(__name__)
def _require_writer_capability() -> None:
import os
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1":
from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability()
def make_embedder(embeddings_cfg):
"""Factory del client embeddings (monkeypatchabile nei test)."""
@@ -64,6 +70,7 @@ def open_searcher(cfg):
def sync_canonical_records(collection, records, *, store, embedder):
_require_writer_capability()
kinds = sorted({record.kind for record in records})
existing = store.existing_hashes(collection, kinds)
pending = []
+51
View File
@@ -0,0 +1,51 @@
"""Capability verifier for mutating workspace children.
The backend passes writer.lock as fd 3 and the retained workspace directory as fd 4.
This module intentionally has no path fallback: callers either run with the capability
or fail closed before touching artifacts.
"""
from __future__ import annotations
import fcntl
import os
import stat
from dataclasses import dataclass
class WorkspaceWriterConflict(RuntimeError):
"preprocessing_conflict"
def __init__(self, message: str = "preprocessing_conflict") -> None:
super().__init__(message)
@dataclass(frozen=True)
class WorkspaceCapability:
workspace_id: str
revision: str
device: int
inode: int
writer_device: int
writer_inode: int
def _identity(env: dict[str, str]) -> tuple[str, str, int, int]:
wid, rev = env.get("THOTH_WORKSPACE_ID"), env.get("THOTH_WORKSPACE_REVISION")
if not wid or not rev or not __import__("re").fullmatch(r"[a-z][a-z0-9-]{2,62}", wid) or not __import__("re").fullmatch(r"[0-9a-f]{40}", rev):
raise WorkspaceWriterConflict()
try: device, inode = int(env["THOTH_WORKSPACE_DEVICE"]), int(env["THOTH_WORKSPACE_INODE"])
except (KeyError, ValueError): raise WorkspaceWriterConflict()
return wid, rev, device, inode
def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: dict[str, str] | None = None) -> WorkspaceCapability:
env = dict(os.environ if env is None else env)
wid, rev, device, inode = _identity(env)
try: root = os.fstat(root_fd); writer = os.fstat(writer_fd)
except OSError as exc: raise WorkspaceWriterConflict() from exc
if not stat.S_ISDIR(root.st_mode) or root.st_uid != os.getuid() or (root.st_mode & 0o777) != 0o700 or (root.st_dev, root.st_ino) != (device, inode): raise WorkspaceWriterConflict()
if not stat.S_ISREG(writer.st_mode) or writer.st_uid != os.getuid() or (writer.st_mode & 0o777) != 0o600: raise WorkspaceWriterConflict()
try: fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError as exc: raise WorkspaceWriterConflict() from exc
# Keep the OFD locked. A lock check is necessarily best effort on some BSDs; identity and
# descriptor ownership remain mandatory and no path-based lock is accepted.
return WorkspaceCapability(wid, rev, root.st_dev, root.st_ino, writer.st_dev, writer.st_ino)
def require_workspace_writer_capability() -> WorkspaceCapability:
return verify_workspace_writer_fds()