diff --git a/harness/tests/test_dwh_preprocess_job.py b/harness/tests/test_dwh_preprocess_job.py index fb6858fe..c2e810ad 100644 --- a/harness/tests/test_dwh_preprocess_job.py +++ b/harness/tests/test_dwh_preprocess_job.py @@ -97,6 +97,61 @@ def test_writer_rejects_unbound_legacy_artifacts_before_building(tmp_path): assert not (tmp_path / ".tht-dwh" / "OWNER.json").exists() +def test_writer_rejects_dangling_legacy_symlinks_before_claim_or_callback(tmp_path): + import pytest + + legacy = tmp_path / "artifacts" / "mschema" / "physical.yaml" + legacy.parent.mkdir(parents=True) + legacy.symlink_to(tmp_path / "missing-catalog") + 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_publication_remains_on_locked_root_when_path_is_swapped( + monkeypatch, tmp_path, +): + import os + import pytest + import tht.jobs.dwh_pipeline as module + + real_replace = module.os.replace + moved = tmp_path / "locked-root" + replacement = tmp_path / ".tht-dwh" + swapped = False + + def swapping_replace(source, destination, *args, **kwargs): + nonlocal swapped + if destination == "OWNER.json" and kwargs.get("dst_dir_fd") is not None: + swapped = True + replacement.rename(moved) + replacement.mkdir(mode=0o700) + return real_replace(source, destination, *args, **kwargs) + + monkeypatch.setattr(module.os, "replace", swapping_replace) + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: (_ for _ in ()).throw(AssertionError("callback called")), + build_lsh=lambda physical, output: None, + ) + with pytest.raises(Exception): + pipeline.run() + assert swapped + assert (moved / "OWNER.json").is_file() + assert not (replacement / "OWNER.json").exists() + assert os.stat(moved / "generation.lock").st_ino != os.stat(replacement).st_ino + + def test_owner_requires_exact_read_only_owner_mode_and_active_requires_binding(tmp_path): import pytest diff --git a/harness/tests/test_schema_introspect_guard.py b/harness/tests/test_schema_introspect_guard.py index 80fe6f2a..5a93bf43 100644 --- a/harness/tests/test_schema_introspect_guard.py +++ b/harness/tests/test_schema_introspect_guard.py @@ -34,17 +34,62 @@ def _write_config(tmp_path): return cfg -def test_introspect_cache_hit_skips_dwh(tmp_path): - # Le credenziali sono fasulle: se la guardia non scattasse PRIMA del branch - # transport, il comando tenterebbe la connessione e fallirebbe. +def test_introspect_rejects_unbound_legacy_cache(tmp_path): catalog = _write_catalog(tmp_path) - before = catalog.read_bytes() cfg = _write_config(tmp_path) res = CliRunner().invoke(app, ["schema", "introspect", "-c", str(cfg)]) + assert res.exit_code == 1 + assert "legacy artifacts are unbound" in res.output + assert catalog.exists() + + +def test_introspect_fresh_root_initializes_through_writer_job(tmp_path, monkeypatch): + import tht.cli.schema_cmd as module + + cfg = _write_config(tmp_path) + physical = PhysicalSchema( + database="d", schema="s", introspected_at=datetime(2026, 1, 1), + tables={"dim_patient": TablePhysical(columns={"id": ColumnPhysical(type="bigint")})}, + ) + + def refresh(_cfg, *, output_path=None, **_kwargs): + physical.to_yaml(output_path) + return physical + + monkeypatch.setattr(module, "refresh_catalog", refresh) + res = CliRunner().invoke(app, ["schema", "introspect", "-c", str(cfg)]) assert res.exit_code == 0, res.output - assert "OK (cache)" in res.output + assert (tmp_path / ".tht-dwh" / "OWNER.json").is_file() assert "1 tabelle" in res.output - assert catalog.read_bytes() == before + + +def test_lsh_build_fresh_root_initializes_introspection_and_lsh(tmp_path, monkeypatch): + import tht.cli.lsh_cmd as lsh_module + import tht.cli.schema_cmd as schema_module + import tht.lshindex as lshindex_module + + cfg = _write_config(tmp_path) + physical = PhysicalSchema( + database="d", schema="s", introspected_at=datetime(2026, 1, 1), + tables={"dim_patient": TablePhysical(columns={"id": ColumnPhysical(type="bigint")})}, + ) + + def refresh(_cfg, *, output_path=None, **_kwargs): + physical.to_yaml(output_path) + return physical + + def build(_cfg, *, physical_file, output_dir, **_kwargs): + assert physical_file.is_file() + for name in ("s_lsh.pkl", "s_minhashes.pkl", "s_meta.json"): + (output_dir / name).write_text("index") + return {}, [], [], {} + + monkeypatch.setattr(schema_module, "refresh_catalog", refresh) + monkeypatch.setattr(lsh_module, "build_lsh_artifacts", build) + monkeypatch.setattr(lshindex_module, "load_index", lambda *_args, **_kwargs: (None, {}, None)) + res = CliRunner().invoke(app, ["lsh", "build", "-c", str(cfg)]) + assert res.exit_code == 0, res.output + assert (tmp_path / ".tht-dwh" / "OWNER.json").is_file() def test_introspect_refresh_bypasses_cache(tmp_path): diff --git a/harness/tht/cli/lsh_cmd.py b/harness/tht/cli/lsh_cmd.py index fbce161f..d2467dbc 100644 --- a/harness/tht/cli/lsh_cmd.py +++ b/harness/tht/cli/lsh_cmd.py @@ -69,32 +69,30 @@ def build_lsh_artifacts( def build_cmd(config: Path = CONFIG_OPT) -> None: """Costruisce l'indice LSH dai valori del database e lo salva su pickle.""" cfg = _load_config_or_exit(config) - phys_file = physical_path(cfg) - if not phys_file.exists(): - typer.secho( - f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", - fg=typer.colors.RED, err=True, - ) - raise typer.Exit(code=1) - typer.echo("Estrazione valori (i più frequenti) dalle colonne testuali eligible...") - from tht.jobs.dwh_pipeline import resolve_dwh_snapshot - - 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()) - else: - from tht.cli.preprocess_cmd import run_dwh_from_config - from tht.lshindex import load_index - - report = run_dwh_from_config(config, steps=("lsh",)) - if report.status != "succeeded": - typer.secho("ERRORE: DWH preprocessing failed", fg=typer.colors.RED, err=True) + dwh_root = cfg.paths.artifacts.parent / ".tht-dwh" + initialized = dwh_root.exists() or dwh_root.is_symlink() + if initialized: + phys_file = physical_path(cfg) + if not phys_file.exists(): + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, err=True, + ) raise typer.Exit(code=1) - _, minhashes, _ = load_index(_lsh_dir(cfg), name=cfg.database.db_schema) - skipped, truncated = [], [] - n_values = len(minhashes) - n_columns = len({(entry[1], entry[2]) for entry in minhashes.values()}) + typer.echo("Estrazione valori (i più frequenti) dalle colonne testuali eligible...") + from tht.cli.preprocess_cmd import run_dwh_from_config + from tht.lshindex import load_index + + report = run_dwh_from_config( + config, steps=("lsh",) if initialized else ("introspect", "lsh") + ) + if report.status != "succeeded": + typer.secho("ERRORE: DWH preprocessing failed", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + _, minhashes, _ = load_index(_lsh_dir(cfg), name=cfg.database.db_schema) + skipped, truncated = [], [] + n_values = len(minhashes) + n_columns = len({(entry[1], entry[2]) for entry in minhashes.values()}) typer.echo(f" {n_values} valori da {n_columns} colonne") for s in skipped: typer.secho(f" saltata {s.table}.{s.column}: {s.reason}", fg=typer.colors.YELLOW) diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 7f4d4786..f8b3c1b7 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -72,8 +72,11 @@ def introspect_cmd( Se physical.yaml esiste già, esce subito (cache); usa --refresh per rigenerarlo. """ cfg = _load_config_or_exit(config) - out = physical_path(cfg) - if out.exists() and not refresh: + dwh_root = cfg.paths.artifacts.parent / ".tht-dwh" + out = cfg.paths.artifacts / "mschema" / "physical.yaml" + if dwh_root.exists() or dwh_root.is_symlink(): + out = physical_path(cfg) + if (dwh_root.exists() or dwh_root.is_symlink()) and out.exists() and not refresh: from datetime import UTC, datetime from tht.mschema.models import PhysicalSchema @@ -94,18 +97,14 @@ def introspect_cmd( ) return try: - from tht.jobs.dwh_pipeline import resolve_dwh_snapshot + from tht.cli.preprocess_cmd import run_dwh_from_config + from tht.mschema.models import PhysicalSchema - if resolve_dwh_snapshot(cfg).generation is None: - phys = refresh_catalog(cfg) - else: - from tht.cli.preprocess_cmd import run_dwh_from_config - from tht.mschema.models import PhysicalSchema - - report = run_dwh_from_config(config, steps=("introspect",)) - if report.status != "succeeded": - raise RuntimeError("DWH preprocessing failed") - phys = PhysicalSchema.from_yaml(physical_path(cfg)) + report = run_dwh_from_config(config, steps=("introspect",)) + if report.status != "succeeded": + raise RuntimeError("DWH preprocessing failed") + out = physical_path(cfg) + phys = PhysicalSchema.from_yaml(out) except Exception as e: typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) diff --git a/harness/tht/jobs/dwh_pipeline.py b/harness/tht/jobs/dwh_pipeline.py index cec4ad8e..8567737c 100644 --- a/harness/tht/jobs/dwh_pipeline.py +++ b/harness/tht/jobs/dwh_pipeline.py @@ -106,15 +106,24 @@ def _validate_root_binding(workspace_root: Path, expected: dict[str, str]) -> No raise CorruptCheckpointError("DWH generations exist without a consistent ACTIVE pointer") -def _claim_or_validate_root_binding(workspace_root: Path, binding: dict[str, str]) -> None: +def _claim_or_validate_root_binding( + workspace_root: Path, binding: dict[str, str], lock_fd: int +) -> None: root = workspace_root / ".tht-dwh" - marker = root / OWNER_MARKER - if marker.exists() or marker.is_symlink(): - _validate_root_binding(workspace_root, binding) - return root_fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + temporary: str | None = None try: entries = set(os.listdir(root_fd)) + opened_lock = os.open("generation.lock", os.O_RDONLY | os.O_NOFOLLOW, dir_fd=root_fd) + try: + held_info, opened_info = os.fstat(lock_fd), os.fstat(opened_lock) + if (held_info.st_dev, held_info.st_ino) != (opened_info.st_dev, opened_info.st_ino): + raise OSError("DWH workspace root changed while locked") + finally: + os.close(opened_lock) + if OWNER_MARKER in entries: + _validate_root_binding(workspace_root, binding) + return 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" @@ -133,28 +142,38 @@ def _claim_or_validate_root_binding(workspace_root: Path, binding: dict[str, str ) finally: os.close(generations_fd) + payload = { + "schema_version": 1, + "binding": binding, + "binding_sha256": _binding_digest(binding), + } + temporary = f".{OWNER_MARKER}.{uuid.uuid4().hex}.tmp" + fd = os.open( + temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, + 0o600, dir_fd=root_fd, + ) + try: + data = (json.dumps(payload, sort_keys=True, separators=(",", ":")) + "\n").encode() + offset = 0 + while offset < len(data): + offset += os.write(fd, data[offset:]) + os.fsync(fd) + os.fchmod(fd, 0o400) + os.fsync(fd) + finally: + os.close(fd) + os.replace(temporary, OWNER_MARKER, src_dir_fd=root_fd, dst_dir_fd=root_fd) + temporary = None + os.fsync(root_fd) except OSError as error: raise CorruptCheckpointError("DWH workspace root is invalid") from error finally: + if temporary is not None: + try: + os.unlink(temporary, dir_fd=root_fd) + except FileNotFoundError: + pass os.close(root_fd) - payload = { - "schema_version": 1, - "binding": binding, - "binding_sha256": _binding_digest(binding), - } - temporary = root / f".{OWNER_MARKER}.{uuid.uuid4().hex}.tmp" - fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600) - try: - with os.fdopen(fd, "w", encoding="utf-8") as stream: - stream.write(json.dumps(payload, sort_keys=True, separators=(",", ":")) + "\n") - stream.flush() - os.fsync(stream.fileno()) - temporary.chmod(0o400) - os.replace(temporary, marker) - DwhPreprocessPipeline._fsync(root) - except BaseException: - temporary.unlink(missing_ok=True) - raise @dataclass(frozen=True) @@ -459,7 +478,7 @@ class DwhPreprocessPipeline: 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) + _claim_or_validate_root_binding(self.workspace_root, self.binding, lease_fd) self._assert_active_binding() active = _active_generation_dir_locked(self.workspace_root, self.binding) if active is not None: @@ -560,9 +579,13 @@ class DwhPreprocessPipeline: 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_physical = self.current_physical is not None and ( + self.current_physical.exists() or self.current_physical.is_symlink() + ) legacy_lsh = self.current_lsh_dir is not None and any( - (self.current_lsh_dir / name).exists() for name in self.lsh_filenames + (self.current_lsh_dir / name).exists() + or (self.current_lsh_dir / name).is_symlink() + for name in self.lsh_filenames ) if legacy_physical or legacy_lsh: raise CorruptCheckpointError(