"""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 and evidence.local_archive_root and ( evidence.local_archive_root / "evidence/.local/state.yaml" ).exists(): return False 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, local_snapshot: Path | 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) from tht.evidence.adapters import FilesystemEvidenceSource sources = [FilesystemEvidenceSource(local_snapshot, patterns=("curated/**/*.md",))] \ if local_snapshot is not None else build_sources(cfg.evidence) pipeline = build_preprocessing_pipeline( store=CorpusStore(corpus_root), sources=sources, 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"), consolidate: bool = typer.Option(False, "--consolidate"), ) -> None: if action is not None and action != "gc": raise typer.BadParameter("only the optional 'gc' action is supported") if consolidate and (dry_run or resume is not None or action is not None): message = "Consolidation cannot be combined with dry-run, resume or gc" if json_output: typer.echo(json.dumps({"status": "failed", "code": "invalid_consolidation", "error": message})) else: typer.secho(message, fg=typer.colors.RED, err=True) raise typer.Exit(code=2) 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: if consolidate: from tht.evidence.administration import consolidate_from_config result = consolidate_from_config(config) else: result = run_from_config(config, dry_run=dry_run, resume=resume) except Exception as error: # noqa: BLE001 from tht.evidence.administration import ConsolidationError detail = str(error) if isinstance(error, ConsolidationError) else "preprocessing failed" payload = {"status": "failed"} if isinstance(error, ConsolidationError): payload["saved"] = error.saved if json_output: typer.echo(json.dumps( _evidence_json_payload( cfg, payload, code="preprocessing_failed", error=detail, ), sort_keys=True, )) else: typer.secho(f"ERRORE: {detail}", 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")