Files
ThothII/harness/tht/jobs/dwh_pipeline.py
T

66 lines
2.1 KiB
Python

"""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()