diff --git a/harness/tests/test_dwh_preprocess_job.py b/harness/tests/test_dwh_preprocess_job.py index b3422b36..fb6858fe 100644 --- a/harness/tests/test_dwh_preprocess_job.py +++ b/harness/tests/test_dwh_preprocess_job.py @@ -6,6 +6,7 @@ from typer.testing import CliRunner from tht.cli import app from tht.jobs.dwh_pipeline import DwhPreprocessPipeline +from tht.jobs.dwh_pipeline import active_generation_dir, config_dwh_binding from tht.jobs.dwh_pipeline import resolve_dwh_snapshot from tht.jobs.dwh_pipeline import lease_dwh_snapshot from tht.jobs.locking import _lock_name @@ -30,6 +31,94 @@ def test_dwh_and_evidence_jobs_have_distinct_lock_names(): assert _lock_name("demo", "dwh") != _lock_name("demo", "evidence") +def test_unowned_reads_fail_closed_without_creating_any_files(tmp_path): + import pytest + + cfg = snapshot_config(tmp_path) + with pytest.raises(Exception, match="not initialized"): + resolve_dwh_snapshot(cfg) + with pytest.raises(Exception, match="not initialized"): + with lease_dwh_snapshot(cfg): + pass + assert not (tmp_path / ".tht-dwh").exists() + + +def test_writer_claim_allows_only_lock_and_empty_generations(tmp_path): + import pytest + + for name, make_entry in ( + ("unexpected", lambda root: (root / "unexpected").write_text("x")), + ("stale-temp", lambda root: (root / ".OWNER.json.stale.tmp").write_text("x")), + ("unexpected-dir", lambda root: (root / "other").mkdir()), + ): + root = tmp_path / name / ".tht-dwh" + root.mkdir(parents=True, mode=0o700) + make_entry(root) + calls = [] + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=root.parent, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: calls.append("called"), + build_lsh=lambda physical, output: None, + ) + with pytest.raises(Exception, match="unbound"): + pipeline.run() + assert calls == [] + assert not (root / "OWNER.json").exists() + + allowed = tmp_path / "allowed" + (allowed / ".tht-dwh" / "generations").mkdir(parents=True, mode=0o700) + report = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=allowed, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: output.write_text("catalog"), + build_lsh=lambda physical, output: _write_lsh([], physical, output), + ).run() + assert report.status == "succeeded" + + +def test_writer_rejects_unbound_legacy_artifacts_before_building(tmp_path): + import pytest + + legacy = tmp_path / "artifacts" / "mschema" / "physical.yaml" + legacy.parent.mkdir(parents=True) + legacy.write_text("legacy") + calls = [] + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: calls.append("called"), + build_lsh=lambda physical, output: None, + current_physical=legacy, + ) + with pytest.raises(Exception, match="legacy artifacts are unbound"): + pipeline.run() + assert calls == [] + assert not (tmp_path / ".tht-dwh" / "OWNER.json").exists() + + +def test_owner_requires_exact_read_only_owner_mode_and_active_requires_binding(tmp_path): + import pytest + + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: output.write_text("catalog"), + build_lsh=lambda physical, output: _write_lsh([], physical, output), + ) + pipeline.run() + cfg = snapshot_config(tmp_path) + binding = config_dwh_binding(cfg) + assert active_generation_dir(tmp_path, binding) is not None + with pytest.raises(Exception, match="different workspace configuration"): + active_generation_dir(tmp_path, {**binding, "workspace_id": "other"}) + + marker = tmp_path / ".tht-dwh" / "OWNER.json" + marker.chmod(0o440) + with pytest.raises(Exception, match="ownership marker"): + resolve_dwh_snapshot(cfg) + + def test_selected_dwh_stages_run_in_declared_order(tmp_path): calls = [] pipeline = DwhPreprocessPipeline( diff --git a/harness/tests/test_search_pack.py b/harness/tests/test_search_pack.py index 86174342..edcfc199 100644 --- a/harness/tests/test_search_pack.py +++ b/harness/tests/test_search_pack.py @@ -5,6 +5,8 @@ from types import SimpleNamespace from typer.testing import CliRunner from tht.cli import app +from tht.config import load_config +from tht.jobs.dwh_pipeline import DwhPreprocessPipeline, config_dwh_binding from tht.mschema.models import ColumnPhysical, PhysicalSchema, TablePhysical from tht.vectorstore.embeddings import EmbeddingsError @@ -43,11 +45,11 @@ class _FakeSearcher: def _workspace(tmp_path, with_session=None): - PhysicalSchema( + physical = PhysicalSchema( database="d", schema="s", introspected_at=datetime(2026, 1, 1), tables={"fact_ablazione": TablePhysical( comment="Ablazioni", columns={"cod_paz": ColumnPhysical(type="bigint")})}, - ).to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") + ) cfg = tmp_path / "workspace.yaml" cfg.write_text( "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" @@ -56,6 +58,18 @@ def _workspace(tmp_path, with_session=None): f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, " f"sessions: {tmp_path/'sessions'}}}\n" ) + binding = config_dwh_binding(load_config(cfg)) + DwhPreprocessPipeline( + workspace_id=binding["workspace_id"], workspace_root=tmp_path, + config_fingerprint=binding["config_fingerprint"], + input_fingerprint=binding["input_fingerprint"], + introspect=lambda output: physical.to_yaml(output), + build_lsh=lambda _physical, output: [ + (output / name).write_text("index") + for name in ("s_lsh.pkl", "s_minhashes.pkl", "s_meta.json") + ], + lsh_filenames=("s_lsh.pkl", "s_minhashes.pkl", "s_meta.json"), + ).run() if with_session: sdir = tmp_path / "sessions" / with_session sdir.mkdir(parents=True) diff --git a/harness/tht/cli/lsh_cmd.py b/harness/tht/cli/lsh_cmd.py index b4e62439..fbce161f 100644 --- a/harness/tht/cli/lsh_cmd.py +++ b/harness/tht/cli/lsh_cmd.py @@ -77,9 +77,9 @@ def build_cmd(config: Path = CONFIG_OPT) -> None: ) raise typer.Exit(code=1) typer.echo("Estrazione valori (i più frequenti) dalle colonne testuali eligible...") - from tht.jobs.dwh_pipeline import active_generation_dir + from tht.jobs.dwh_pipeline import resolve_dwh_snapshot - if active_generation_dir(cfg.paths.artifacts.parent) is None: + if resolve_dwh_snapshot(cfg).generation is None: minhashes, skipped, truncated, values = build_lsh_artifacts(cfg, verbose=True) n_values = sum(len(v) for table in values.values() for v in table.values()) n_columns = sum(len(table) for table in values.values()) diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 7e075bc6..f1aca27a 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -19,15 +19,14 @@ def run_dwh_from_config( config: Path, *, steps: tuple[str, ...], resume: str | None = None, ): from tht.cli.lsh_cmd import build_lsh_artifacts - from tht.cli.schema_cmd import _load_config_or_exit, physical_path, refresh_catalog + from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog from tht.jobs.dwh_pipeline import ( - DwhPreprocessPipeline, active_generation_dir, config_dwh_binding, + DwhPreprocessPipeline, config_dwh_binding, ) cfg = _load_config_or_exit(config) binding = config_dwh_binding(cfg) workspace_root = cfg.paths.artifacts.parent - active = active_generation_dir(workspace_root) lsh_names = ( f"{cfg.database.db_schema}_lsh.pkl", f"{cfg.database.db_schema}_minhashes.pkl", @@ -43,8 +42,8 @@ def run_dwh_from_config( cfg, physical_file=physical, output_dir=output ), lsh_filenames=lsh_names, - current_physical=physical_path(cfg), - current_lsh_dir=active if active is not None else cfg.paths.indexes / "lsh", + current_physical=cfg.paths.artifacts / "mschema" / "physical.yaml", + current_lsh_dir=cfg.paths.indexes / "lsh", ) return pipeline.run(steps, resume_run_id=resume) diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 04dcbbb9..7f4d4786 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -39,6 +39,8 @@ def _load_config_or_exit(config: Path): def physical_path(cfg) -> Path: from tht.jobs.dwh_pipeline import resolve_dwh_snapshot + if not (cfg.paths.artifacts.parent / ".tht-dwh").exists(): + return cfg.paths.artifacts / "mschema" / "physical.yaml" return resolve_dwh_snapshot(cfg).physical @@ -92,9 +94,9 @@ def introspect_cmd( ) return try: - from tht.jobs.dwh_pipeline import active_generation_dir + from tht.jobs.dwh_pipeline import resolve_dwh_snapshot - if active_generation_dir(cfg.paths.artifacts.parent) is None: + if resolve_dwh_snapshot(cfg).generation is None: phys = refresh_catalog(cfg) else: from tht.cli.preprocess_cmd import run_dwh_from_config diff --git a/harness/tht/jobs/dwh_pipeline.py b/harness/tht/jobs/dwh_pipeline.py index 234cfc99..cec4ad8e 100644 --- a/harness/tht/jobs/dwh_pipeline.py +++ b/harness/tht/jobs/dwh_pipeline.py @@ -51,7 +51,22 @@ def _binding_digest(binding: dict[str, str]) -> str: def _read_root_binding(workspace_root: Path) -> dict[str, str]: marker = workspace_root / ".tht-dwh" / OWNER_MARKER try: - payload = json.loads(_read_owned(marker, readonly=True).decode("utf-8")) + fd = os.open(marker, os.O_RDONLY | os.O_NOFOLLOW) + try: + info = os.fstat(fd) + if ( + not stat.S_ISREG(info.st_mode) + or info.st_uid != os.getuid() + or info.st_nlink != 1 + or stat.S_IMODE(info.st_mode) != 0o400 + ): + raise OSError("unsafe DWH ownership marker") + chunks = [] + while chunk := os.read(fd, 1024 * 1024): + chunks.append(chunk) + finally: + os.close(fd) + payload = json.loads(b"".join(chunks).decode("utf-8")) binding = payload["binding"] if ( payload.get("schema_version") != 1 @@ -81,6 +96,9 @@ def _validate_root_binding(workspace_root: Path, expected: dict[str, str]) -> No raise CorruptCheckpointError("DWH generations directory is invalid") from error else: try: + info = os.fstat(generations_fd) + if info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) != 0o700: + raise OSError("unsafe DWH generations directory") generation_entries = os.listdir(generations_fd) finally: os.close(generations_fd) @@ -94,25 +112,31 @@ def _claim_or_validate_root_binding(workspace_root: Path, binding: dict[str, str if marker.exists() or marker.is_symlink(): _validate_root_binding(workspace_root, binding) return - generations = root / "generations" - generations_nonempty = False + root_fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) try: - generations_fd = os.open( - generations, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW - ) - except FileNotFoundError: - pass + entries = set(os.listdir(root_fd)) + if not entries <= {"generation.lock", "generations"} or "generation.lock" not in entries: + raise CorruptCheckpointError( + "DWH artifacts are unbound; migrate them explicitly or use an empty root" + ) + if "generations" in entries: + generations_fd = os.open( + "generations", os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=root_fd + ) + try: + info = os.fstat(generations_fd) + if info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) != 0o700: + raise OSError("unsafe DWH generations directory") + if os.listdir(generations_fd): + raise CorruptCheckpointError( + "DWH artifacts are unbound; migrate them explicitly or use an empty root" + ) + finally: + os.close(generations_fd) except OSError as error: - raise CorruptCheckpointError("DWH generations directory is invalid") from error - else: - try: - generations_nonempty = bool(os.listdir(generations_fd)) - finally: - os.close(generations_fd) - if (root / "ACTIVE").exists() or generations_nonempty: - raise CorruptCheckpointError( - "DWH artifacts are unbound; migrate them explicitly or use an empty root" - ) + raise CorruptCheckpointError("DWH workspace root is invalid") from error + finally: + os.close(root_fd) payload = { "schema_version": 1, "binding": binding, @@ -148,13 +172,9 @@ class DwhSnapshotLease: def __enter__(self) -> DwhArtifactSnapshot: binding = config_dwh_binding(self.cfg) - claim_fd = _acquire_generation_lock(self.cfg.paths.artifacts.parent, exclusive=True) - try: - _claim_or_validate_root_binding(self.cfg.paths.artifacts.parent, binding) - finally: - fcntl.flock(claim_fd, fcntl.LOCK_UN) - os.close(claim_fd) - self._fd = _acquire_generation_lock(self.cfg.paths.artifacts.parent, exclusive=False) + self._fd = _acquire_existing_generation_lock( + self.cfg.paths.artifacts.parent, exclusive=False + ) try: self.snapshot = _resolve_dwh_snapshot_locked(self.cfg, binding) return self.snapshot @@ -179,7 +199,12 @@ def _acquire_generation_lock(workspace_root: Path, *, exclusive: bool) -> int: fd = os.open(root / "generation.lock", os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW, 0o600) try: info = os.fstat(fd) - 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 + or stat.S_IMODE(info.st_mode) != 0o600 + ): raise OSError("unsafe DWH generation lock") fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH) return fd @@ -188,6 +213,37 @@ def _acquire_generation_lock(workspace_root: Path, *, exclusive: bool) -> int: raise +def _acquire_existing_generation_lock(workspace_root: Path, *, exclusive: bool) -> int: + root = workspace_root / ".tht-dwh" + try: + root_fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + root_info = os.fstat(root_fd) + if root_info.st_uid != os.getuid() or stat.S_IMODE(root_info.st_mode) != 0o700: + raise OSError("unsafe DWH workspace root") + fd = os.open("generation.lock", os.O_RDWR | os.O_NOFOLLOW, dir_fd=root_fd) + finally: + os.close(root_fd) + try: + info = os.fstat(fd) + if ( + not stat.S_ISREG(info.st_mode) + or info.st_uid != os.getuid() + or info.st_nlink != 1 + or stat.S_IMODE(info.st_mode) != 0o600 + ): + raise OSError("unsafe DWH generation lock") + fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH) + return fd + except BaseException: + os.close(fd) + raise + except OSError as error: + raise CorruptCheckpointError( + "DWH workspace ownership is not initialized; run preprocessing first" + ) from error + + def _digest(path: Path) -> str: return hashlib.sha256(path.read_bytes()).hexdigest() @@ -289,13 +345,9 @@ def validate_generation( def resolve_dwh_snapshot(cfg) -> DwhArtifactSnapshot: binding = config_dwh_binding(cfg) - claim_fd = _acquire_generation_lock(cfg.paths.artifacts.parent, exclusive=True) - try: - _claim_or_validate_root_binding(cfg.paths.artifacts.parent, binding) - finally: - fcntl.flock(claim_fd, fcntl.LOCK_UN) - os.close(claim_fd) - lease_fd = _acquire_generation_lock(cfg.paths.artifacts.parent, exclusive=False) + lease_fd = _acquire_existing_generation_lock( + cfg.paths.artifacts.parent, exclusive=False + ) try: return _resolve_dwh_snapshot_locked(cfg, binding) finally: @@ -307,7 +359,7 @@ def _resolve_dwh_snapshot_locked( cfg, binding: dict[str, str], ) -> DwhArtifactSnapshot: _validate_root_binding(cfg.paths.artifacts.parent, binding) - target = active_generation_dir(cfg.paths.artifacts.parent) + target = _active_generation_dir_locked(cfg.paths.artifacts.parent, binding) if target is None: return DwhArtifactSnapshot( None, cfg.paths.artifacts / "mschema" / "physical.yaml", cfg.paths.indexes / "lsh" @@ -316,7 +368,21 @@ def _resolve_dwh_snapshot_locked( return DwhArtifactSnapshot(target.name, target / "physical.yaml", target) -def active_generation_dir(workspace_root: Path) -> Path | None: +def active_generation_dir( + workspace_root: Path, expected_binding: dict[str, str] +) -> Path | None: + lease_fd = _acquire_existing_generation_lock(workspace_root, exclusive=False) + try: + return _active_generation_dir_locked(workspace_root, expected_binding) + finally: + fcntl.flock(lease_fd, fcntl.LOCK_UN) + os.close(lease_fd) + + +def _active_generation_dir_locked( + workspace_root: Path, expected_binding: dict[str, str] +) -> Path | None: + _validate_root_binding(workspace_root, expected_binding) pointer = workspace_root / ".tht-dwh" / "ACTIVE" try: generation = _read_owned(pointer, readonly=False).decode("utf-8").strip() @@ -329,6 +395,7 @@ def active_generation_dir(workspace_root: Path) -> Path | None: target = pointer.parent / "generations" / generation if not target.is_dir() or target.is_symlink(): raise CorruptCheckpointError("active DWH generation is missing") + validate_generation(target, expected_binding) return target @@ -381,7 +448,7 @@ class DwhPreprocessPipeline: def _assert_active_binding(self) -> None: _validate_root_binding(self.workspace_root, self.binding) - active = active_generation_dir(self.workspace_root) + active = _active_generation_dir_locked(self.workspace_root, self.binding) if active is not None: validate_generation(active, self.binding) @@ -391,8 +458,13 @@ class DwhPreprocessPipeline: self._validate_steps(steps) lease_fd = _acquire_generation_lock(self.workspace_root, exclusive=True) try: + self._assert_no_legacy_artifacts() _claim_or_validate_root_binding(self.workspace_root, self.binding) self._assert_active_binding() + active = _active_generation_dir_locked(self.workspace_root, self.binding) + if active is not None: + self.current_physical = active / "physical.yaml" + self.current_lsh_dir = active finally: fcntl.flock(lease_fd, fcntl.LOCK_UN) os.close(lease_fd) @@ -454,7 +526,7 @@ class DwhPreprocessPipeline: return set() target = self.workspace_root / ".tht-dwh" / "generations" / source.run_id validate_generation(target, self.binding) - active = active_generation_dir(self.workspace_root) + active = active_generation_dir(self.workspace_root, self.binding) if active != target: raise CorruptCheckpointError("sealed DWH publication is not ACTIVE") self._validate_published(target, run_dir / "artifacts", running.artifact_files) @@ -484,6 +556,19 @@ class DwhPreprocessPipeline: raise CorruptCheckpointError("resume artifact manifest is invalid") from error self._validate_published(target, run_dir / "artifacts", required) + def _assert_no_legacy_artifacts(self) -> None: + marker = self.workspace_root / ".tht-dwh" / OWNER_MARKER + if marker.exists() or marker.is_symlink(): + return + legacy_physical = self.current_physical is not None and self.current_physical.exists() + legacy_lsh = self.current_lsh_dir is not None and any( + (self.current_lsh_dir / name).exists() for name in self.lsh_filenames + ) + if legacy_physical or legacy_lsh: + raise CorruptCheckpointError( + "DWH legacy artifacts are unbound; migrate them explicitly or use an empty root" + ) + @staticmethod def _artifacts(context: JobContext) -> Path: root = context.run_dir / "artifacts" @@ -630,7 +715,7 @@ class DwhPreprocessPipeline: if not root.exists(): return _validate_root_binding(self.workspace_root, self.binding) - active = active_generation_dir(self.workspace_root) + active = _active_generation_dir_locked(self.workspace_root, self.binding) if active is not None: validate_generation(active, self.binding) active_name = active.name if active else None