Files
ThothII/harness/tht/cli/preprocess_cmd.py
T
Codex cffa60772e
Publish documentation / publish (push) Successful in 2m12s
feat: complete catalog-driven preprocessing
2026-09-06 17:49:35 +02:00

565 lines
22 KiB
Python

"""One-shot preprocessing commands."""
from __future__ import annotations
import hashlib
import json
import os
import re
import shutil
import stat
from pathlib import Path
import typer
from tht.cli.config_cmd import CONFIG_OPT
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
def _bm25_language(workspace_language: str) -> str:
languages = {"en": "english", "it": "italian"}
try:
return languages[workspace_language]
except KeyError as exc:
raise ValueError("workspace language is unsupported for Qdrant BM25") from exc
def _evidence_json_context(config: Path):
from tht.cli.schema_cmd import _load_config_or_exit
return _load_config_or_exit(config)
def _evaluation_workspace_root(cfg) -> Path:
evidence = cfg.evidence
if evidence is None:
raise RuntimeError("Evidence evaluation fixture is unavailable")
if evidence.source_root is not None:
return evidence.source_root
filesystem_roots = [source.root for source in evidence.sources if source.type == "filesystem"]
if len(filesystem_roots) != 1:
raise RuntimeError("Evidence evaluation requires one filesystem workspace source")
root = filesystem_roots[0]
if root.name == "curated":
root = root.parent
if root.name == "evidence":
return root.parent
return root
def _requires_candidate_evaluation(cfg) -> bool:
evidence = cfg.evidence
if evidence is None or evidence.schema_version != 2:
return False
return evidence.source_root is not None or any(
source.type == "filesystem" for source in evidence.sources
)
def _candidate_evaluator(cfg, *, vector_store, embedder):
"""Bind candidate publication to the same read-only retrieval evaluator as the CLI."""
if not _requires_candidate_evaluation(cfg):
return None
from tht.evidence.canonical import load_curated_tree
from tht.evidence.evaluation import evaluate_retrieval, load_evaluation_fixture
workspace_root = _evaluation_workspace_root(cfg)
language = _bm25_language(cfg.language)
def evaluate(manifest):
document_generations = manifest.metadata.get("document_generations")
if not isinstance(document_generations, dict):
raise TypeError("candidate Evidence generation is invalid")
return evaluate_retrieval(
load_evaluation_fixture(workspace_root / "evidence" / "evaluation.yaml"),
workspace_revision=cfg._workspace_revision,
vector_generation=manifest.vector_generation,
document_generations=document_generations,
workspace_id=cfg._workspace_id,
language=language,
searcher=vector_store,
embedder=embedder,
expected_kinds={
evidence.id: evidence.kind
for evidence in load_curated_tree(workspace_root / "evidence" / "curated")
},
)
return evaluate
def _validate_materialized_curated_corpus(cfg) -> None:
"""Fail closed on a v2 pinned filesystem corpus before any vector write is possible."""
if not _requires_candidate_evaluation(cfg):
return
from tht.evidence import validate_workspace_evidence
if not validate_workspace_evidence(_evaluation_workspace_root(cfg)).publishable:
raise RuntimeError("curated Evidence corpus is invalid")
def _evidence_json_payload(cfg, payload: dict, *, code: str, error: str | None = None) -> dict:
value = {
**payload,
"schemaVersion": 1,
"status": payload.get("status", "failed"),
"code": code,
"operation": "preprocess_evidence",
"workspaceId": cfg._workspace_id,
"workspaceRevision": cfg._workspace_revision,
}
if error is not None:
value["error"] = error
return value
def _simple_json_payload(*, code: str, error: str) -> dict:
return {
"schemaVersion": 1,
"status": "failed",
"code": code,
"operation": "preprocess_evidence",
"error": error,
}
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"],
catalog_database_id=binding.get("catalog_database_id"),
metadata_content_revision=binding.get("metadata_content_revision"),
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 run_catalog_dwh_from_config(cfg):
"""Publish Catalog-derived physical schema and LSH as one bound generation."""
from tht.cli.lsh_cmd import build_lsh_artifacts
from tht.jobs.dwh_pipeline import DwhPreprocessPipeline, config_dwh_binding
from tht.mschema.catalog_snapshot import load_catalog_metadata_snapshot
snapshot = load_catalog_metadata_snapshot(
cfg.paths.catalog_metadata_snapshot, cfg._workspace_id
)
binding = config_dwh_binding(cfg)
workspace_root = cfg.paths.artifacts.parent
# Catalog preprocessing is a replace-in-place operation. Its DWH/LSH output is wholly
# derived, there is no supported concurrent runtime, and a failed Clear may have left an
# older binding behind. Start from an empty owned generation root so a retry can always
# rebuild the current Catalog revision instead of deadlocking on the stale OWNER marker.
_remove_owned_derived_path(workspace_root, workspace_root / ".tht-dwh")
_remove_owned_derived_path(workspace_root, cfg.paths.artifacts / "mschema" / "physical.yaml")
_remove_owned_derived_path(workspace_root, cfg.paths.indexes / "lsh")
observed: dict[str, object] = {}
def materialize_physical(output: Path):
physical, _annotations, _relationships = snapshot.to_schema_inputs()
physical.to_yaml(output)
return physical
def materialize_lsh(physical: Path, output: Path):
result = build_lsh_artifacts(cfg, physical_file=physical, output_dir=output)
observed["result"] = result
return result
report = DwhPreprocessPipeline(
workspace_id=str(binding["workspace_id"]),
workspace_root=workspace_root,
config_fingerprint=str(binding["config_fingerprint"]),
input_fingerprint=str(binding["input_fingerprint"]),
catalog_database_id=str(binding["catalog_database_id"]),
metadata_content_revision=int(binding["metadata_content_revision"]),
introspect=materialize_physical,
build_lsh=materialize_lsh,
lsh_filenames=(
f"{cfg.database.db_schema}_lsh.pkl",
f"{cfg.database.db_schema}_minhashes.pkl",
f"{cfg.database.db_schema}_meta.json",
),
).run(("introspect", "lsh"))
result = observed.get("result")
if not isinstance(result, tuple) or len(result) != 4:
raise RuntimeError("Catalog LSH publication did not complete")
return report, result
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_vector_store
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli.vector_cmd import make_embedder
from tht.evidence import build_preprocessing_pipeline, build_sources
from tht.evidence.corpus.chunk import ChunkPolicy
from tht.evidence.corpus.store import CorpusStore
cfg = _load_config_or_exit(config)
if cfg.embeddings is None:
raise RuntimeError("embeddings are not configured")
_validate_materialized_curated_corpus(cfg)
corpus_root = cfg.paths.artifacts.parent / "corpus"
vector_store = build_vector_store(cfg, require_write=True)
embedder = make_embedder(cfg.embeddings)
pipeline = build_preprocessing_pipeline(
store=CorpusStore(corpus_root), sources=build_sources(cfg.evidence),
embedder=embedder,
vector_store=vector_store,
embedding_id=cfg.embeddings.id or f"ollama/{cfg.embeddings.model}",
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,
sparse_language=_bm25_language(cfg.language),
candidate_evaluator=_candidate_evaluator(cfg, vector_store=vector_store, embedder=embedder),
)
def fingerprint(value: str) -> str:
return "sha256:" + hashlib.sha256(value.encode()).hexdigest()
return pipeline.run_as_job(
workspace_id=cfg._workspace_id,
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_vector_store
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli.vector_cmd import make_embedder
from tht.evidence import build_preprocessing_pipeline, build_sources
from tht.evidence.corpus.chunk import ChunkPolicy
from tht.evidence.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 = build_preprocessing_pipeline(
store=CorpusStore(corpus_root), sources=build_sources(cfg.evidence),
embedder=make_embedder(cfg.embeddings), vector_store=build_vector_store(cfg, require_write=True),
embedding_id=cfg.embeddings.id or f"ollama/{cfg.embeddings.model}",
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,
sparse_language=_bm25_language(cfg.language),
)
pipeline.workspace_id = cfg._workspace_id
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
def _remove_owned_derived_path(workspace_root: Path, target: Path) -> bool:
"""Remove one generated path without following links or escaping the workspace."""
root = workspace_root.resolve()
resolved = target.resolve(strict=False)
if not resolved.is_relative_to(root):
raise RuntimeError("derived cleanup target escapes the workspace")
try:
info = target.lstat()
except FileNotFoundError:
return False
if stat.S_ISLNK(info.st_mode) or info.st_uid != os.getuid():
raise RuntimeError("derived cleanup target is unsafe")
if stat.S_ISDIR(info.st_mode):
shutil.rmtree(target)
elif stat.S_ISREG(info.st_mode):
target.unlink()
else:
raise RuntimeError("derived cleanup target is unsafe")
return True
def clear_from_config(config: Path) -> dict[str, int]:
"""Clear workspace reference vectors and local preprocessing derivatives."""
from tht.adapters.factory import build_vector_store
from tht.cli._guards import require_vector_write_allowed
from tht.cli.schema_cmd import _load_config_or_exit
cfg = _load_config_or_exit(config)
require_vector_write_allowed(cfg, "preprocess clear")
workspace_root = cfg.paths.artifacts.parent
vector_store = build_vector_store(cfg, require_write=True)
reference_deleted = int(vector_store.clear_reference())
paths = [
workspace_root / ".tht-dwh",
workspace_root / ".tht-jobs",
workspace_root / "corpus",
cfg.paths.artifacts / "mschema" / "physical.yaml",
cfg.paths.indexes / "lsh",
]
if cfg.paths.catalog_metadata_snapshot is not None:
paths.append(cfg.paths.catalog_metadata_snapshot)
removed = sum(int(_remove_owned_derived_path(workspace_root, path)) for path in paths)
return {"referenceCollections": reference_deleted, "derivedPaths": removed}
@preprocess_app.command("clear", hidden=True)
def clear_cmd(
config: Path = CONFIG_OPT,
json_output: bool = typer.Option(False, "--json"),
) -> None:
"""Clear replaceable preprocessing output while preserving workspace memory."""
try:
counts = clear_from_config(config)
except Exception: # noqa: BLE001 - do not disclose paths, endpoints, or credentials
payload = {
"schemaVersion": 1,
"status": "failed",
"code": "preprocessing_clear_failed",
"operation": "preprocess_clear",
"error": "Preprocessing clear failed",
}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: preprocessing clear failed", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None
payload = {
"schemaVersion": 1,
"status": "succeeded",
"code": "ok",
"operation": "preprocess_clear",
"counts": counts,
}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.echo("OK: reference vectors and LSH cleared; memory preserved")
@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 Exception: # noqa: BLE001
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
cfg = _evidence_json_context(config) if json_output else None
if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None:
payload = _simple_json_payload(
code="invalid_resume",
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 Exception: # noqa: BLE001
payload = {"status": "failed"}
if json_output:
typer.echo(json.dumps(
_evidence_json_payload(
cfg,
payload,
code="preprocessing_failed",
error="preprocessing failed",
),
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":
if json_output:
typer.echo(json.dumps(
_evidence_json_payload(
cfg,
payload,
code="preprocessing_failed",
error="preprocessing job failed",
),
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(
_evidence_json_payload(cfg, payload, code="ok"),
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 Exception: # noqa: BLE001
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)
@preprocess_app.command("catalog", hidden=True)
def catalog_cmd(
catalog_metadata: Path = typer.Option(..., "--catalog-metadata"),
config: Path = CONFIG_OPT,
json_output: bool = typer.Option(False, "--json"),
) -> None:
"""Build current LSH and schema vectors from one immutable Catalog snapshot."""
from tht.cli._guards import require_vector_write_allowed
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli.vector_cmd import index_catalog_schema, require_vector_cfg
cfg = _load_config_or_exit(config)
cfg = cfg.model_copy(
update={
"paths": cfg.paths.model_copy(
update={"catalog_metadata_snapshot": catalog_metadata}
)
}
)
require_vector_write_allowed(cfg, "preprocess catalog")
require_vector_cfg(cfg)
try:
report, (minhashes, skipped, truncated, _values) = run_catalog_dwh_from_config(cfg)
if report.status != "succeeded":
raise RuntimeError("Catalog DWH preprocessing failed")
stats, counts, snapshot_path = index_catalog_schema(cfg)
except Exception: # noqa: BLE001 - public output must never disclose endpoints or SQL
payload = {
"schemaVersion": 1,
"status": "failed",
"code": "catalog_preprocessing_failed",
"operation": "preprocess_catalog",
"error": "Catalog preprocessing failed",
}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
else:
typer.secho("ERRORE: Catalog preprocessing failed", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None
payload = {
"schemaVersion": 1,
"status": "succeeded",
"code": "ok",
"operation": "preprocess_catalog",
"artifactIdentities": [
{
"kind": "catalog_metadata_snapshot",
"digest": "sha256:" + hashlib.sha256(snapshot_path.read_bytes()).hexdigest(),
}
],
"counts": {
**counts,
"lshEntries": len(minhashes),
"lshSkippedColumns": len(skipped),
"lshTruncatedColumns": len(truncated),
"added": stats.added,
"deleted": stats.deleted,
"unchanged": stats.unchanged,
"updated": stats.updated,
},
}
if json_output:
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
else:
typer.echo("OK: Catalog metadata, LSH and schema vectors rebuilt")