feat(evidence): build semantic fragments from typed units

This commit is contained in:
2026-08-24 21:33:48 +02:00
parent 8adc085746
commit 3420c57c8b
9 changed files with 642 additions and 14 deletions
+26 -4
View File
@@ -13,8 +13,9 @@ from datetime import UTC
from pathlib import Path
import tht.evidence.acquisition as evidence_acquisition
from tht.evidence.canonical import ReviewItem
from tht.evidence.contracts import EvidenceSource, SourceObject, canonical_provenance_uri
from tht.evidence.corpus.chunk import ChunkPolicy, chunk
from tht.evidence.corpus.chunk import AtomicContentTooLargeError, ChunkPolicy, chunk
from tht.evidence.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest
from tht.evidence.corpus.normalize import normalize
from tht.evidence.corpus.store import CorpusStore
@@ -51,6 +52,7 @@ class PipelineResult:
manifest: CorpusManifest = field(repr=False)
run_id: str | None = None
resumed_from: str | None = None
review_items: tuple[ReviewItem, ...] = ()
def __repr__(self) -> str:
counts = {
@@ -85,6 +87,7 @@ class PipelineResult:
"manifest_id": self.manifest.manifest_id,
"run_id": self.run_id,
"resumed_from": self.resumed_from,
"review_items": [item.model_dump(mode="json") for item in self.review_items],
}
@@ -466,7 +469,11 @@ class CorpusPipeline:
evidence_acquisition.acquire(source, item), self.pipeline_version,
))
documents.sort(key=lambda value: value.source_id)
chunks = [part for document in documents for part in chunk(document, self.chunk_policy)]
try:
chunks = [part for document in documents for part in chunk(document, self.chunk_policy)]
except AtomicContentTooLargeError as error:
write(context, "review-items.json", [error.review_item.model_dump(mode="json")])
raise
previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {}
changed = set(plan["changed"])
generations = {
@@ -667,6 +674,14 @@ class CorpusPipeline:
], after_stage_return=after_stage_return)
run_dir = workspace_root / ".tht-jobs" / "evidence" / "runs" / report.run_id
plan = json.loads((run_dir / "artifacts" / "plan.json").read_text())
review_items_path = run_dir / "artifacts" / "review-items.json"
try:
review_items = tuple(
ReviewItem.model_validate(item)
for item in json.loads(review_items_path.read_text(encoding="utf-8"))
) if review_items_path.exists() else ()
except (OSError, ValueError) as error:
raise PipelineError("preprocessing review items are corrupt") from error
if dry_run:
manifest = self.store.active_manifest() or CorpusManifest(pipeline_version=self.pipeline_version)
generation = None
@@ -682,9 +697,9 @@ class CorpusPipeline:
if manifest_path.exists() else CorpusManifest(pipeline_version=self.pipeline_version))
published = False
return PipelineResult(
report.status, generation, published, tuple(plan["changed"]),
"blocked" if review_items else report.status, generation, published, tuple(plan["changed"]),
tuple(plan["unchanged"]), tuple(plan["removed"]), manifest,
report.run_id, report.resumed_from,
report.run_id, report.resumed_from, review_items,
)
def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult:
@@ -778,6 +793,13 @@ class CorpusPipeline:
)
self.store.publish(staged)
self.gc(workspace_root=self.store.root.parent)
except AtomicContentTooLargeError as error:
self._compensate(generation, vector_written)
manifest = previous or CorpusManifest(pipeline_version=self.pipeline_version)
return PipelineResult(
"blocked", None, False, changed, unchanged, removed, manifest,
review_items=(error.review_item,),
)
except PipelineError:
self._compensate(generation, vector_written)
raise