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

334 lines
13 KiB
Python

"""Evidence-owned runtime lookup bound to the atomically active corpus generation."""
import re
import unicodedata
from dataclasses import dataclass
from typing import Literal, Protocol
from pydantic import BaseModel, ConfigDict
from tht.evidence.canonical import EvidenceKind, EvidencePurpose
from tht.evidence.corpus.store import CorpusStore
from tht.ports.vector import VectorReadUnavailable, VectorStoreError
class CorpusWorkspaceMismatchError(RuntimeError):
"""The configured workspace does not own the persisted corpus."""
class EvidenceQueryEmbedder(Protocol):
def embed_query(self, query: str) -> list[float]: ...
class EvidenceSearchContext(BaseModel):
"""Optional query enrichment and explicit, server-enforced Evidence constraints."""
model_config = ConfigDict(frozen=True, extra="forbid")
concepts: tuple[str, ...] = ()
tables: tuple[str, ...] = ()
columns: tuple[str, ...] = ()
required_kinds: tuple[EvidenceKind, ...] = ()
required_concepts: tuple[str, ...] = ()
required_tables: tuple[str, ...] = ()
required_columns: tuple[str, ...] = ()
@dataclass(frozen=True)
class EvidenceResult:
evidence_id: str
title: str
kind: str
excerpts: tuple[str, ...]
provenance: dict
citation: str
document_id: str
score: float
@dataclass(frozen=True)
class EvidenceSearchOutcome:
status: Literal["available", "unavailable"]
vector_generation: str | None
results: tuple[EvidenceResult, ...] = ()
code: str | None = None
message: str | None = None
@classmethod
def unavailable(cls, code: str, message: str) -> "EvidenceSearchOutcome":
return cls("unavailable", None, (), code, message[:240])
def _normalized_values(values: tuple[str, ...]) -> tuple[str, ...]:
normalized = {
unicodedata.normalize("NFC", value).strip()
for value in values
if isinstance(value, str) and unicodedata.normalize("NFC", value).strip()
}
return tuple(sorted(normalized))
def render_evidence_query(query: str, context: EvidenceSearchContext) -> str:
"""Build the one exact query text shared by dense and BM25 retrieval."""
question = unicodedata.normalize("NFC", query).replace("\r\n", "\n").replace("\r", "\n").strip()
if not question:
raise ValueError("Evidence query must not be empty")
sections = [("Domanda", question)]
for label, values in (
("Concetti", _normalized_values(context.concepts)),
("Tabelle", _normalized_values(context.tables)),
("Colonne", _normalized_values(context.columns)),
):
if values:
sections.append((label, ", ".join(values)))
return "\n".join(f"{label}: {value}" for label, value in sections)
def _required_metadata_filter(purpose: EvidencePurpose, context: EvidenceSearchContext) -> dict[str, object]:
return {
"purpose": purpose,
"required_kinds": list(_normalized_values(context.required_kinds)),
"required_concepts": list(_normalized_values(context.required_concepts)),
"required_tables": list(_normalized_values(context.required_tables)),
"required_columns": list(_normalized_values(context.required_columns)),
}
def _active_vector_generation(searcher) -> str | None:
supplied = getattr(searcher, "vector_generation", None)
if isinstance(supplied, str) and supplied:
return supplied
corpus = getattr(searcher, "corpus", None)
if corpus is None:
return None
with corpus.writer_lock():
manifest = corpus.active_manifest()
return manifest.vector_generation if manifest is not None else None
def _group_evidence_fragments(hits) -> tuple[EvidenceResult, ...]:
grouped: dict[str, list] = {}
for hit in hits:
metadata = getattr(hit, "metadata", {})
evidence_id = metadata.get("evidence_id") if isinstance(metadata, dict) else None
if not isinstance(evidence_id, str) or not evidence_id:
raise VectorStoreError("Evidence search returned malformed payload")
grouped.setdefault(evidence_id, []).append(hit)
results = []
for evidence_id, fragments in grouped.items():
ordered = sorted(
fragments,
key=lambda item: (-float(item.similarity), int(item.metadata.get("ordinal", 0)), item.id),
)
first = ordered[0]
metadata = first.metadata
citation = metadata.get("source_uri")
document_id = metadata.get("document_id")
evidence_kind = metadata.get("evidence_kind")
if not all(isinstance(value, str) and value for value in (citation, document_id, evidence_kind)):
raise VectorStoreError("Evidence search returned malformed payload")
results.append(EvidenceResult(
evidence_id=evidence_id,
title=str(first.title),
kind=evidence_kind,
excerpts=tuple(str(item.content) for item in ordered),
provenance=dict(metadata.get("provenance", {})),
citation=citation,
document_id=document_id,
score=float(first.similarity),
))
return tuple(sorted(results, key=lambda item: (-item.score, item.evidence_id)))
def search_evidence(
query: str,
purpose: EvidencePurpose,
context: EvidenceSearchContext,
*,
searcher,
embedder: EvidenceQueryEmbedder,
top_n: int = 10,
) -> EvidenceSearchOutcome:
"""Search the active generation once, with no stale-generation or purpose fallback."""
try:
rendered = render_evidence_query(query, context)
generation = _active_vector_generation(searcher)
if generation is None:
return EvidenceSearchOutcome.unavailable("active_corpus_unavailable", "Active Evidence corpus is unavailable")
query_embedding = embedder.embed_query(rendered)
hits = searcher.search(
query_embedding,
top_n=top_n,
kinds=["evidence"],
query_text=rendered,
metadata_filter=_required_metadata_filter(purpose, context),
)
return EvidenceSearchOutcome("available", generation, _group_evidence_fragments(hits))
except VectorReadUnavailable:
return EvidenceSearchOutcome.unavailable("vector_unavailable", "Evidence vector search is unavailable")
except (VectorStoreError, CorpusWorkspaceMismatchError):
return EvidenceSearchOutcome.unavailable("evidence_search_unavailable", "Evidence search is unavailable")
class ActiveEvidenceSearcher:
"""Searcher facade that enforces ACTIVE generation predicates before LIMIT."""
def __init__(
self,
corpus: CorpusStore,
delegate,
expected_workspace_id: str | None = None,
evidence_language: str = "italian",
):
self.corpus = corpus
self.delegate = delegate
self.expected_workspace_id = expected_workspace_id
self.evidence_language = evidence_language
def search(
self,
embedding,
top_n=10,
kinds=None,
metadata_filter=None,
query_text=None,
query_language=None,
):
requested = set(kinds) if kinds is not None else {
"schema_table", "schema_column", "schema_relationship", "evidence", "memory", "solved_question",
}
include_evidence = "evidence" in requested
other_kinds = sorted(requested - {"evidence"})
with self.corpus.writer_lock():
manifest = self.corpus.active_manifest()
persisted_workspace = manifest.metadata.get("workspace_id") if manifest else None
if manifest is not None and (
not isinstance(persisted_workspace, str)
or re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", persisted_workspace) is None
):
raise CorpusWorkspaceMismatchError(
"corpus workspace ownership is missing or invalid; use a new corpus root or rebuild"
)
if manifest is not None and self.expected_workspace_id is not None and (
persisted_workspace != self.expected_workspace_id
):
raise CorpusWorkspaceMismatchError(
"corpus belongs to a different workspace; use a new corpus root or rebuild"
)
if not include_evidence:
kwargs = {"top_n": top_n, "kinds": kinds}
if metadata_filter is not None:
kwargs["metadata_filter"] = metadata_filter
return self.delegate.search(embedding, **kwargs)
hits = []
if other_kinds:
kwargs = {"top_n": top_n, "kinds": other_kinds}
if metadata_filter is not None:
kwargs["metadata_filter"] = metadata_filter
hits.extend(self.delegate.search(embedding, **kwargs))
if include_evidence:
workspace_id = manifest.metadata.get("workspace_id") if manifest else None
if manifest is not None and isinstance(workspace_id, str):
by_generation: dict[str, list[str]] = {}
mapping = dict(manifest.metadata.get("document_generations", {}))
for document in manifest.documents:
generation = mapping.get(document.document_id, manifest.vector_generation)
if generation:
by_generation.setdefault(generation, []).append(document.document_id)
if by_generation and (not isinstance(query_text, str) or query_text.strip() == ""):
raise VectorStoreError("Evidence hybrid query text is required")
for generation, document_ids in sorted(by_generation.items()):
filters = dict(metadata_filter or {})
filters.update({
"vector_generation": generation,
"document_ids": sorted(document_ids),
"workspace_id": workspace_id,
})
hits.extend(self.delegate.search(
embedding, top_n=top_n, kinds=["evidence"],
query_text=query_text,
query_language=query_language or self.evidence_language,
metadata_filter=filters,
))
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:top_n]
def active_searcher(cfg, delegate, *, workspace_id: str | None = None):
corpus_root = cfg.paths.artifacts.parent / "corpus"
languages = {"en": "english", "it": "italian"}
language = languages.get(getattr(cfg, "language", "en"))
if language is None:
raise VectorStoreError("workspace language is unsupported for Qdrant BM25")
return ActiveEvidenceSearcher(CorpusStore(corpus_root), delegate, workspace_id, language)
def validate_corpus_workspace(cfg, workspace_id: str) -> None:
"""Fail before downstream retrieval setup when configured corpus ownership differs."""
corpus = CorpusStore(cfg.paths.artifacts.parent / "corpus")
with corpus.writer_lock():
manifest = corpus.active_manifest()
if manifest is None:
return
persisted = manifest.metadata.get("workspace_id")
if not isinstance(persisted, str) or re.fullmatch(
r"[a-z][a-z0-9_-]{0,63}", persisted
) is None:
raise CorpusWorkspaceMismatchError(
"corpus workspace ownership is missing or invalid; use a new corpus root or rebuild"
)
if persisted != workspace_id:
raise CorpusWorkspaceMismatchError(
"corpus belongs to a different workspace; use a new corpus root or rebuild"
)
def resolve_citation(
store: CorpusStore, evidence_id: str, *, materialized_root=None,
) -> str:
with store.writer_lock():
manifest = store.active_manifest()
if manifest is None:
return ""
for document in manifest.documents:
frontmatter = document.metadata.get("frontmatter", {})
identifiers = {document.document_id, document.source_id, str(frontmatter.get("id", ""))}
if evidence_id in identifiers:
root = materialized_root or (store.root / "runtime")
filename = document.document_id.removeprefix("doc:") + ".md"
path = store.materialize_document(
document.document_id, root / filename, generation=manifest.manifest_id,
)
return str(path) if path else ""
return ""
def build_retrieval_entries(results, *, excerpt_chars: int) -> list[dict]:
"""Project ordered Evidence search hits into the retrieval-pack shape."""
return [
{
"title": getattr(result, "title", getattr(result, "label", "")),
"status": getattr(result, "status", None),
"excerpt": (
result.excerpts[0] if isinstance(result, EvidenceResult) and result.excerpts
else getattr(result, "content", "")
)[:excerpt_chars],
}
for result in results
]
__all__ = [
"ActiveEvidenceSearcher",
"CorpusWorkspaceMismatchError",
"EvidenceQueryEmbedder",
"EvidenceResult",
"EvidenceSearchContext",
"EvidenceSearchOutcome",
"active_searcher",
"build_retrieval_entries",
"render_evidence_query",
"resolve_citation",
"search_evidence",
"validate_corpus_workspace",
]