fix(preprocess): restrict DWH root claims to writers
This commit is contained in:
@@ -6,6 +6,7 @@ from typer.testing import CliRunner
|
|||||||
|
|
||||||
from tht.cli import app
|
from tht.cli import app
|
||||||
from tht.jobs.dwh_pipeline import DwhPreprocessPipeline
|
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 resolve_dwh_snapshot
|
||||||
from tht.jobs.dwh_pipeline import lease_dwh_snapshot
|
from tht.jobs.dwh_pipeline import lease_dwh_snapshot
|
||||||
from tht.jobs.locking import _lock_name
|
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")
|
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):
|
def test_selected_dwh_stages_run_in_declared_order(tmp_path):
|
||||||
calls = []
|
calls = []
|
||||||
pipeline = DwhPreprocessPipeline(
|
pipeline = DwhPreprocessPipeline(
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ from types import SimpleNamespace
|
|||||||
from typer.testing import CliRunner
|
from typer.testing import CliRunner
|
||||||
|
|
||||||
from tht.cli import app
|
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.mschema.models import ColumnPhysical, PhysicalSchema, TablePhysical
|
||||||
from tht.vectorstore.embeddings import EmbeddingsError
|
from tht.vectorstore.embeddings import EmbeddingsError
|
||||||
|
|
||||||
@@ -43,11 +45,11 @@ class _FakeSearcher:
|
|||||||
|
|
||||||
|
|
||||||
def _workspace(tmp_path, with_session=None):
|
def _workspace(tmp_path, with_session=None):
|
||||||
PhysicalSchema(
|
physical = PhysicalSchema(
|
||||||
database="d", schema="s", introspected_at=datetime(2026, 1, 1),
|
database="d", schema="s", introspected_at=datetime(2026, 1, 1),
|
||||||
tables={"fact_ablazione": TablePhysical(
|
tables={"fact_ablazione": TablePhysical(
|
||||||
comment="Ablazioni", columns={"cod_paz": ColumnPhysical(type="bigint")})},
|
comment="Ablazioni", columns={"cod_paz": ColumnPhysical(type="bigint")})},
|
||||||
).to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml")
|
)
|
||||||
cfg = tmp_path / "workspace.yaml"
|
cfg = tmp_path / "workspace.yaml"
|
||||||
cfg.write_text(
|
cfg.write_text(
|
||||||
"database: {database: d, schema: s, user: u, password: p, transport: direct}\n"
|
"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"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, "
|
||||||
f"sessions: {tmp_path/'sessions'}}}\n"
|
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:
|
if with_session:
|
||||||
sdir = tmp_path / "sessions" / with_session
|
sdir = tmp_path / "sessions" / with_session
|
||||||
sdir.mkdir(parents=True)
|
sdir.mkdir(parents=True)
|
||||||
|
|||||||
@@ -77,9 +77,9 @@ def build_cmd(config: Path = CONFIG_OPT) -> None:
|
|||||||
)
|
)
|
||||||
raise typer.Exit(code=1)
|
raise typer.Exit(code=1)
|
||||||
typer.echo("Estrazione valori (i più frequenti) dalle colonne testuali eligible...")
|
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)
|
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_values = sum(len(v) for table in values.values() for v in table.values())
|
||||||
n_columns = sum(len(table) for table in values.values())
|
n_columns = sum(len(table) for table in values.values())
|
||||||
|
|||||||
@@ -19,15 +19,14 @@ def run_dwh_from_config(
|
|||||||
config: Path, *, steps: tuple[str, ...], resume: str | None = None,
|
config: Path, *, steps: tuple[str, ...], resume: str | None = None,
|
||||||
):
|
):
|
||||||
from tht.cli.lsh_cmd import build_lsh_artifacts
|
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 (
|
from tht.jobs.dwh_pipeline import (
|
||||||
DwhPreprocessPipeline, active_generation_dir, config_dwh_binding,
|
DwhPreprocessPipeline, config_dwh_binding,
|
||||||
)
|
)
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
cfg = _load_config_or_exit(config)
|
||||||
binding = config_dwh_binding(cfg)
|
binding = config_dwh_binding(cfg)
|
||||||
workspace_root = cfg.paths.artifacts.parent
|
workspace_root = cfg.paths.artifacts.parent
|
||||||
active = active_generation_dir(workspace_root)
|
|
||||||
lsh_names = (
|
lsh_names = (
|
||||||
f"{cfg.database.db_schema}_lsh.pkl",
|
f"{cfg.database.db_schema}_lsh.pkl",
|
||||||
f"{cfg.database.db_schema}_minhashes.pkl",
|
f"{cfg.database.db_schema}_minhashes.pkl",
|
||||||
@@ -43,8 +42,8 @@ def run_dwh_from_config(
|
|||||||
cfg, physical_file=physical, output_dir=output
|
cfg, physical_file=physical, output_dir=output
|
||||||
),
|
),
|
||||||
lsh_filenames=lsh_names,
|
lsh_filenames=lsh_names,
|
||||||
current_physical=physical_path(cfg),
|
current_physical=cfg.paths.artifacts / "mschema" / "physical.yaml",
|
||||||
current_lsh_dir=active if active is not None else cfg.paths.indexes / "lsh",
|
current_lsh_dir=cfg.paths.indexes / "lsh",
|
||||||
)
|
)
|
||||||
return pipeline.run(steps, resume_run_id=resume)
|
return pipeline.run(steps, resume_run_id=resume)
|
||||||
|
|
||||||
|
|||||||
@@ -39,6 +39,8 @@ def _load_config_or_exit(config: Path):
|
|||||||
def physical_path(cfg) -> Path:
|
def physical_path(cfg) -> Path:
|
||||||
from tht.jobs.dwh_pipeline import resolve_dwh_snapshot
|
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
|
return resolve_dwh_snapshot(cfg).physical
|
||||||
|
|
||||||
|
|
||||||
@@ -92,9 +94,9 @@ def introspect_cmd(
|
|||||||
)
|
)
|
||||||
return
|
return
|
||||||
try:
|
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)
|
phys = refresh_catalog(cfg)
|
||||||
else:
|
else:
|
||||||
from tht.cli.preprocess_cmd import run_dwh_from_config
|
from tht.cli.preprocess_cmd import run_dwh_from_config
|
||||||
|
|||||||
@@ -51,7 +51,22 @@ def _binding_digest(binding: dict[str, str]) -> str:
|
|||||||
def _read_root_binding(workspace_root: Path) -> dict[str, str]:
|
def _read_root_binding(workspace_root: Path) -> dict[str, str]:
|
||||||
marker = workspace_root / ".tht-dwh" / OWNER_MARKER
|
marker = workspace_root / ".tht-dwh" / OWNER_MARKER
|
||||||
try:
|
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"]
|
binding = payload["binding"]
|
||||||
if (
|
if (
|
||||||
payload.get("schema_version") != 1
|
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
|
raise CorruptCheckpointError("DWH generations directory is invalid") from error
|
||||||
else:
|
else:
|
||||||
try:
|
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)
|
generation_entries = os.listdir(generations_fd)
|
||||||
finally:
|
finally:
|
||||||
os.close(generations_fd)
|
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():
|
if marker.exists() or marker.is_symlink():
|
||||||
_validate_root_binding(workspace_root, binding)
|
_validate_root_binding(workspace_root, binding)
|
||||||
return
|
return
|
||||||
generations = root / "generations"
|
root_fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW)
|
||||||
generations_nonempty = False
|
|
||||||
try:
|
try:
|
||||||
generations_fd = os.open(
|
entries = set(os.listdir(root_fd))
|
||||||
generations, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW
|
if not entries <= {"generation.lock", "generations"} or "generation.lock" not in entries:
|
||||||
)
|
raise CorruptCheckpointError(
|
||||||
except FileNotFoundError:
|
"DWH artifacts are unbound; migrate them explicitly or use an empty root"
|
||||||
pass
|
)
|
||||||
|
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:
|
except OSError as error:
|
||||||
raise CorruptCheckpointError("DWH generations directory is invalid") from error
|
raise CorruptCheckpointError("DWH workspace root is invalid") from error
|
||||||
else:
|
finally:
|
||||||
try:
|
os.close(root_fd)
|
||||||
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"
|
|
||||||
)
|
|
||||||
payload = {
|
payload = {
|
||||||
"schema_version": 1,
|
"schema_version": 1,
|
||||||
"binding": binding,
|
"binding": binding,
|
||||||
@@ -148,13 +172,9 @@ class DwhSnapshotLease:
|
|||||||
|
|
||||||
def __enter__(self) -> DwhArtifactSnapshot:
|
def __enter__(self) -> DwhArtifactSnapshot:
|
||||||
binding = config_dwh_binding(self.cfg)
|
binding = config_dwh_binding(self.cfg)
|
||||||
claim_fd = _acquire_generation_lock(self.cfg.paths.artifacts.parent, exclusive=True)
|
self._fd = _acquire_existing_generation_lock(
|
||||||
try:
|
self.cfg.paths.artifacts.parent, exclusive=False
|
||||||
_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)
|
|
||||||
try:
|
try:
|
||||||
self.snapshot = _resolve_dwh_snapshot_locked(self.cfg, binding)
|
self.snapshot = _resolve_dwh_snapshot_locked(self.cfg, binding)
|
||||||
return self.snapshot
|
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)
|
fd = os.open(root / "generation.lock", os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW, 0o600)
|
||||||
try:
|
try:
|
||||||
info = os.fstat(fd)
|
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")
|
raise OSError("unsafe DWH generation lock")
|
||||||
fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH)
|
fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH)
|
||||||
return fd
|
return fd
|
||||||
@@ -188,6 +213,37 @@ def _acquire_generation_lock(workspace_root: Path, *, exclusive: bool) -> int:
|
|||||||
raise
|
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:
|
def _digest(path: Path) -> str:
|
||||||
return hashlib.sha256(path.read_bytes()).hexdigest()
|
return hashlib.sha256(path.read_bytes()).hexdigest()
|
||||||
|
|
||||||
@@ -289,13 +345,9 @@ def validate_generation(
|
|||||||
|
|
||||||
def resolve_dwh_snapshot(cfg) -> DwhArtifactSnapshot:
|
def resolve_dwh_snapshot(cfg) -> DwhArtifactSnapshot:
|
||||||
binding = config_dwh_binding(cfg)
|
binding = config_dwh_binding(cfg)
|
||||||
claim_fd = _acquire_generation_lock(cfg.paths.artifacts.parent, exclusive=True)
|
lease_fd = _acquire_existing_generation_lock(
|
||||||
try:
|
cfg.paths.artifacts.parent, exclusive=False
|
||||||
_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)
|
|
||||||
try:
|
try:
|
||||||
return _resolve_dwh_snapshot_locked(cfg, binding)
|
return _resolve_dwh_snapshot_locked(cfg, binding)
|
||||||
finally:
|
finally:
|
||||||
@@ -307,7 +359,7 @@ def _resolve_dwh_snapshot_locked(
|
|||||||
cfg, binding: dict[str, str],
|
cfg, binding: dict[str, str],
|
||||||
) -> DwhArtifactSnapshot:
|
) -> DwhArtifactSnapshot:
|
||||||
_validate_root_binding(cfg.paths.artifacts.parent, binding)
|
_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:
|
if target is None:
|
||||||
return DwhArtifactSnapshot(
|
return DwhArtifactSnapshot(
|
||||||
None, cfg.paths.artifacts / "mschema" / "physical.yaml", cfg.paths.indexes / "lsh"
|
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)
|
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"
|
pointer = workspace_root / ".tht-dwh" / "ACTIVE"
|
||||||
try:
|
try:
|
||||||
generation = _read_owned(pointer, readonly=False).decode("utf-8").strip()
|
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
|
target = pointer.parent / "generations" / generation
|
||||||
if not target.is_dir() or target.is_symlink():
|
if not target.is_dir() or target.is_symlink():
|
||||||
raise CorruptCheckpointError("active DWH generation is missing")
|
raise CorruptCheckpointError("active DWH generation is missing")
|
||||||
|
validate_generation(target, expected_binding)
|
||||||
return target
|
return target
|
||||||
|
|
||||||
|
|
||||||
@@ -381,7 +448,7 @@ class DwhPreprocessPipeline:
|
|||||||
|
|
||||||
def _assert_active_binding(self) -> None:
|
def _assert_active_binding(self) -> None:
|
||||||
_validate_root_binding(self.workspace_root, self.binding)
|
_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:
|
if active is not None:
|
||||||
validate_generation(active, self.binding)
|
validate_generation(active, self.binding)
|
||||||
|
|
||||||
@@ -391,8 +458,13 @@ class DwhPreprocessPipeline:
|
|||||||
self._validate_steps(steps)
|
self._validate_steps(steps)
|
||||||
lease_fd = _acquire_generation_lock(self.workspace_root, exclusive=True)
|
lease_fd = _acquire_generation_lock(self.workspace_root, exclusive=True)
|
||||||
try:
|
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)
|
||||||
self._assert_active_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:
|
finally:
|
||||||
fcntl.flock(lease_fd, fcntl.LOCK_UN)
|
fcntl.flock(lease_fd, fcntl.LOCK_UN)
|
||||||
os.close(lease_fd)
|
os.close(lease_fd)
|
||||||
@@ -454,7 +526,7 @@ class DwhPreprocessPipeline:
|
|||||||
return set()
|
return set()
|
||||||
target = self.workspace_root / ".tht-dwh" / "generations" / source.run_id
|
target = self.workspace_root / ".tht-dwh" / "generations" / source.run_id
|
||||||
validate_generation(target, self.binding)
|
validate_generation(target, self.binding)
|
||||||
active = active_generation_dir(self.workspace_root)
|
active = active_generation_dir(self.workspace_root, self.binding)
|
||||||
if active != target:
|
if active != target:
|
||||||
raise CorruptCheckpointError("sealed DWH publication is not ACTIVE")
|
raise CorruptCheckpointError("sealed DWH publication is not ACTIVE")
|
||||||
self._validate_published(target, run_dir / "artifacts", running.artifact_files)
|
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
|
raise CorruptCheckpointError("resume artifact manifest is invalid") from error
|
||||||
self._validate_published(target, run_dir / "artifacts", required)
|
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
|
@staticmethod
|
||||||
def _artifacts(context: JobContext) -> Path:
|
def _artifacts(context: JobContext) -> Path:
|
||||||
root = context.run_dir / "artifacts"
|
root = context.run_dir / "artifacts"
|
||||||
@@ -630,7 +715,7 @@ class DwhPreprocessPipeline:
|
|||||||
if not root.exists():
|
if not root.exists():
|
||||||
return
|
return
|
||||||
_validate_root_binding(self.workspace_root, self.binding)
|
_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:
|
if active is not None:
|
||||||
validate_generation(active, self.binding)
|
validate_generation(active, self.binding)
|
||||||
active_name = active.name if active else None
|
active_name = active.name if active else None
|
||||||
|
|||||||
Reference in New Issue
Block a user