fix(preprocess): initialize fresh DWH writers safely

This commit is contained in:
2026-07-12 07:07:31 +02:00
parent 22e806a41b
commit 4cf0adbdd2
5 changed files with 190 additions and 70 deletions
+23 -25
View File
@@ -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)
+12 -13
View File
@@ -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)
+49 -26
View File
@@ -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(