Files
ThothII/harness/tht/cli/preprocess_cmd.py
T
Codex 82e2c91f42
Publish documentation / publish (push) Successful in 1m27s
feat: implement memory and evidence administration with guided repairs
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation.

Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
2026-09-10 10:31:34 +02:00

589 lines
23 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 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")