251 lines
10 KiB
Python
251 lines
10 KiB
Python
"""One-shot preprocessing commands."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import re
|
|
from pathlib import Path
|
|
|
|
import typer
|
|
|
|
from tht.cli.config_cmd import CONFIG_OPT
|
|
from tht.cli.schema_cmd import _load_config_or_exit
|
|
from tht.config import workspace_id_for_config
|
|
from tht.ports.evidence import EvidenceSourceError
|
|
from tht.ports.vector import VectorStoreError
|
|
from tht.vectorstore.embeddings import EmbeddingsError
|
|
|
|
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_PREPROCESS_EXPECTED_ERRORS = (
|
|
OSError, RuntimeError, ValueError, TypeError, KeyError,
|
|
EvidenceSourceError, VectorStoreError, EmbeddingsError,
|
|
)
|
|
|
|
|
|
def run_dwh_from_config(
|
|
config: Path, *, steps: tuple[str, ...], resume: str | None = None,
|
|
):
|
|
from tht.cli.lsh_cmd import build_lsh_artifacts
|
|
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
|
|
from tht.jobs.dwh_pipeline import (
|
|
DwhPreprocessPipeline,
|
|
config_dwh_binding,
|
|
)
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
binding = config_dwh_binding(cfg)
|
|
workspace_root = cfg.paths.artifacts.parent
|
|
lsh_names = (
|
|
f"{cfg.database.db_schema}_lsh.pkl",
|
|
f"{cfg.database.db_schema}_minhashes.pkl",
|
|
f"{cfg.database.db_schema}_meta.json",
|
|
)
|
|
pipeline = DwhPreprocessPipeline(
|
|
workspace_id=binding["workspace_id"],
|
|
workspace_root=workspace_root,
|
|
config_fingerprint=binding["config_fingerprint"],
|
|
input_fingerprint=binding["input_fingerprint"],
|
|
introspect=lambda output: refresh_catalog(cfg, output_path=output),
|
|
build_lsh=lambda physical, output: build_lsh_artifacts(
|
|
cfg, physical_file=physical, output_dir=output
|
|
),
|
|
lsh_filenames=lsh_names,
|
|
current_physical=cfg.paths.artifacts / "mschema" / "physical.yaml",
|
|
current_lsh_dir=cfg.paths.indexes / "lsh",
|
|
)
|
|
return pipeline.run(steps, resume_run_id=resume)
|
|
|
|
|
|
def _parse_dwh_steps(value: str) -> tuple[str, ...]:
|
|
allowed = ("introspect", "lsh")
|
|
steps = tuple(part.strip() for part in value.split(",") if part.strip())
|
|
if (
|
|
not steps
|
|
or len(steps) != len(set(steps))
|
|
or any(step not in allowed for step in steps)
|
|
or tuple(sorted(steps, key=allowed.index)) != steps
|
|
):
|
|
raise ValueError("steps must be a unique ordered subset of introspect,lsh")
|
|
return steps
|
|
|
|
|
|
def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None):
|
|
from tht.adapters.factory import build_evidence_sources, build_vector_store
|
|
from tht.cli.vector_cmd import make_embedder
|
|
from tht.corpus.chunk import ChunkPolicy
|
|
from tht.corpus.pipeline import CorpusPipeline
|
|
from tht.corpus.store import CorpusStore
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
if cfg.embeddings is None:
|
|
raise RuntimeError("embeddings are not configured")
|
|
corpus_root = cfg.paths.artifacts.parent / "corpus"
|
|
pipeline = CorpusPipeline(
|
|
store=CorpusStore(corpus_root), sources=build_evidence_sources(cfg),
|
|
embedder=make_embedder(cfg.embeddings),
|
|
vector_store=build_vector_store(cfg, require_write=True),
|
|
embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim,
|
|
chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=cfg.vector.max_chunk_chars),
|
|
pipeline_version="evidence-v1",
|
|
retain_published_generations=cfg.vector.retain_published_generations,
|
|
)
|
|
def fingerprint(value: str) -> str:
|
|
return "sha256:" + hashlib.sha256(value.encode()).hexdigest()
|
|
|
|
return pipeline.run_as_job(
|
|
workspace_id=workspace_id_for_config(cfg, config),
|
|
workspace_root=corpus_root.parent,
|
|
config_fingerprint=fingerprint(cfg.model_dump_json()),
|
|
input_fingerprint=fingerprint(config.resolve().as_posix()),
|
|
dry_run=dry_run,
|
|
resume_run_id=resume,
|
|
)
|
|
|
|
|
|
def gc_from_config(config: Path, *, dry_run: bool = False):
|
|
from tht.adapters.factory import build_evidence_sources, build_vector_store
|
|
from tht.cli.vector_cmd import make_embedder
|
|
from tht.corpus.chunk import ChunkPolicy
|
|
from tht.corpus.pipeline import CorpusPipeline
|
|
from tht.corpus.store import CorpusStore
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
if cfg.embeddings is None:
|
|
raise RuntimeError("embeddings are not configured")
|
|
corpus_root = cfg.paths.artifacts.parent / "corpus"
|
|
pipeline = CorpusPipeline(
|
|
store=CorpusStore(corpus_root), sources=build_evidence_sources(cfg),
|
|
embedder=make_embedder(cfg.embeddings), vector_store=build_vector_store(cfg, require_write=True),
|
|
embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim,
|
|
chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=cfg.vector.max_chunk_chars),
|
|
pipeline_version="evidence-v1",
|
|
retain_published_generations=cfg.vector.retain_published_generations,
|
|
)
|
|
pipeline.workspace_id = workspace_id_for_config(cfg, config)
|
|
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
|
|
|
|
|
|
@preprocess_app.command("evidence")
|
|
def evidence_cmd(
|
|
action: str | None = typer.Argument(None),
|
|
config: Path = CONFIG_OPT,
|
|
dry_run: bool = typer.Option(False, "--dry-run"),
|
|
resume: str | None = typer.Option(None, "--resume"),
|
|
json_output: bool = typer.Option(False, "--json"),
|
|
) -> None:
|
|
if action is not None and action != "gc":
|
|
raise typer.BadParameter("only the optional 'gc' action is supported")
|
|
if action == "gc":
|
|
try:
|
|
payload = gc_from_config(config, dry_run=dry_run)
|
|
except _PREPROCESS_EXPECTED_ERRORS:
|
|
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: evidence cleanup failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
except Exception:
|
|
if not json_output:
|
|
logger.exception("Evidence cleanup failed")
|
|
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: evidence cleanup failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
|
else:
|
|
typer.echo(f"OK: evicted={len(payload['evicted'])} failures={len(payload['failures'])}")
|
|
return
|
|
if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None:
|
|
payload = {"status": "failed", "error": "resume requires a preprocessing run id"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: resume requires a preprocessing run id", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=2)
|
|
try:
|
|
result = run_from_config(config, dry_run=dry_run, resume=resume)
|
|
except _PREPROCESS_EXPECTED_ERRORS:
|
|
payload = {"status": "failed", "error": "preprocessing failed"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
except Exception:
|
|
if not json_output:
|
|
logger.exception("Evidence preprocessing failed")
|
|
payload = {"status": "failed", "error": "preprocessing failed"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
payload = result.model_dump(mode="json")
|
|
if payload.get("status") != "succeeded":
|
|
payload["error"] = "preprocessing job failed"
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: preprocessing job failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1)
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
|
else:
|
|
counts = payload["counts"]
|
|
typer.echo(
|
|
f"OK: run={payload['run_id']} generation={payload['generation']} "
|
|
f"changed={counts['changed']} unchanged={counts['unchanged']} "
|
|
f"removed={counts['removed']}"
|
|
)
|
|
|
|
|
|
@preprocess_app.command("dwh")
|
|
def dwh_cmd(
|
|
config: Path = CONFIG_OPT,
|
|
steps: str = typer.Option("introspect,lsh", "--steps"),
|
|
resume: str | None = typer.Option(None, "--resume"),
|
|
json_output: bool = typer.Option(False, "--json"),
|
|
) -> None:
|
|
try:
|
|
selected = _parse_dwh_steps(steps)
|
|
except ValueError:
|
|
payload = {"status": "failed", "error": "invalid DWH preprocessing steps"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: invalid DWH preprocessing steps", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=2) from None
|
|
if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None:
|
|
payload = {"status": "failed", "error": "resume requires a preprocessing run id"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: resume requires a preprocessing run id", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=2)
|
|
try:
|
|
result = run_dwh_from_config(config, steps=selected, resume=resume)
|
|
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
|
|
payload = {"status": "failed", "error": "DWH preprocessing failed"}
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, sort_keys=True))
|
|
else:
|
|
typer.secho("ERRORE: DWH preprocessing failed", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
payload = result.model_dump(mode="json")
|
|
if json_output:
|
|
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
|
elif result.status == "succeeded":
|
|
typer.echo(f"OK: run={result.run_id} stages={','.join(selected)}")
|
|
else:
|
|
typer.secho(f"ERRORE: run={result.run_id} DWH preprocessing failed", fg=typer.colors.RED, err=True)
|
|
if result.status != "succeeded":
|
|
raise typer.Exit(code=1)
|