"""Resumable DWH catalog and LSH preprocessing stages.""" from __future__ import annotations import hashlib from collections.abc import Callable from pathlib import Path from tht.jobs.models import JobReport, JobSpec from tht.jobs.runner import run_job DWH_STAGE_IDS = ("introspect", "lsh") class DwhPreprocessPipeline: """Adapt existing DWH preprocessing operations to the shared job envelope.""" def __init__( self, *, workspace_id: str, workspace_root: Path, config_fingerprint: str, input_fingerprint: str, introspect: Callable[[], object], build_lsh: Callable[[], object], ) -> None: self.workspace_id = workspace_id self.workspace_root = workspace_root self.config_fingerprint = config_fingerprint self.input_fingerprint = input_fingerprint self._operations = {"introspect": introspect, "lsh": build_lsh} def run( self, steps: tuple[str, ...] = DWH_STAGE_IDS, *, resume_run_id: str | None = None ) -> JobReport: if not steps or len(steps) != len(set(steps)) or any( step not in DWH_STAGE_IDS for step in steps ): raise ValueError("DWH preprocessing steps must be unique introspect/lsh stages") if tuple(sorted(steps, key=DWH_STAGE_IDS.index)) != steps: raise ValueError("DWH preprocessing steps must follow introspect,lsh order") spec = JobSpec( workspace_id=self.workspace_id, job_type="dwh", workspace_root=self.workspace_root, spec_version="jobs-v1", pipeline_version="dwh-v1", config_fingerprint=self.config_fingerprint, input_fingerprint=self.input_fingerprint, stage_ids=steps, resume_run_id=resume_run_id, ) def stage(operation): def execute(_context): return operation() return execute return run_job(spec, tuple(stage(self._operations[step]) for step in steps)) def fingerprint(value: str) -> str: return "sha256:" + hashlib.sha256(value.encode("utf-8")).hexdigest()