785 lines
37 KiB
Python
785 lines
37 KiB
Python
"""Incremental Evidence preprocessing with generation-isolated vector writes."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import re
|
|
import uuid
|
|
from collections.abc import Mapping, Sequence
|
|
from dataclasses import asdict, dataclass
|
|
from datetime import UTC
|
|
from pathlib import Path
|
|
|
|
from tht.corpus.chunk import ChunkPolicy, chunk
|
|
from tht.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest
|
|
from tht.corpus.normalize import normalize
|
|
from tht.corpus.store import CorpusStore
|
|
from tht.ports.evidence import EvidenceSource, SourceObject, canonical_provenance_uri
|
|
from tht.ports.vector import VectorStore, VectorWriteRecord
|
|
from tht.vectorstore.records import VectorRecord
|
|
from tht.jobs.models import JobSpec
|
|
from tht.jobs.runner import JobContext, StageArtifacts, run_job, seal_stage_artifacts
|
|
|
|
|
|
EVIDENCE_STAGE_IDS = (
|
|
"discover",
|
|
"acquire_normalize_chunk",
|
|
"embed",
|
|
"vector_upsert",
|
|
"stage_validate",
|
|
"publish",
|
|
"retention_cleanup",
|
|
)
|
|
|
|
|
|
class PipelineError(RuntimeError):
|
|
"""Credential-free failure at the preprocessing boundary."""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class PipelineResult:
|
|
status: str
|
|
generation: str | None
|
|
published: bool
|
|
changed: tuple[str, ...]
|
|
unchanged: tuple[str, ...]
|
|
removed: tuple[str, ...]
|
|
manifest: CorpusManifest
|
|
run_id: str | None = None
|
|
resumed_from: str | None = None
|
|
|
|
def model_dump(self, mode=None):
|
|
return {
|
|
"status": self.status,
|
|
"generation": self.generation,
|
|
"published": self.published,
|
|
"changed": list(self.changed),
|
|
"unchanged": list(self.unchanged),
|
|
"removed": list(self.removed),
|
|
"manifest": self.manifest.model_dump(mode="json"),
|
|
"run_id": self.run_id,
|
|
"resumed_from": self.resumed_from,
|
|
}
|
|
|
|
|
|
def _fingerprint(value) -> str:
|
|
payload = json.dumps(value, sort_keys=True, separators=(",", ":"), default=str)
|
|
return "sha256:" + hashlib.sha256(payload.encode()).hexdigest()
|
|
|
|
|
|
def _canonical_json(value):
|
|
if isinstance(value, Mapping):
|
|
return {str(key): _canonical_json(value[key]) for key in sorted(value)}
|
|
if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)):
|
|
return [_canonical_json(child) for child in value]
|
|
return value
|
|
|
|
|
|
def _source_snapshot(discovered) -> dict[str, dict]:
|
|
snapshot = {}
|
|
for _, item in discovered:
|
|
modified_at = item.modified_at.astimezone(UTC) if item.modified_at else None
|
|
metadata = _canonical_json(item.metadata)
|
|
snapshot[item.source_id] = {
|
|
"source_id": item.source_id,
|
|
"uri": item.uri,
|
|
"fingerprint": item.fingerprint,
|
|
"modified_at": modified_at.isoformat().replace("+00:00", "Z") if modified_at else None,
|
|
"metadata": metadata,
|
|
"media_type": metadata.get("media_type"),
|
|
"size": metadata.get("size"),
|
|
}
|
|
return snapshot
|
|
|
|
|
|
class CorpusPipeline:
|
|
def __init__(
|
|
self, *, store: CorpusStore, sources: list[EvidenceSource], embedder,
|
|
vector_store: VectorStore, embedding_model: str, embedding_dimensions: int,
|
|
chunk_policy: ChunkPolicy, pipeline_version: str, retain_published_generations: int = 3,
|
|
workspace_id: str | None = None,
|
|
) -> None:
|
|
self.store = store
|
|
self.sources = sources
|
|
self.embedder = embedder
|
|
self.vector_store = vector_store
|
|
self.embedding_model = embedding_model
|
|
self.embedding_dimensions = embedding_dimensions
|
|
self.chunk_policy = chunk_policy
|
|
self.pipeline_version = pipeline_version
|
|
if isinstance(retain_published_generations, bool) or retain_published_generations < 1:
|
|
raise ValueError("retain_published_generations must be at least 1")
|
|
self.retain_published_generations = retain_published_generations
|
|
self.workspace_id = workspace_id
|
|
|
|
def _assert_workspace_binding(self) -> None:
|
|
manifest = self.store.active_manifest()
|
|
if manifest is None:
|
|
if self.workspace_id is None:
|
|
self.workspace_id = "default"
|
|
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 PipelineError(
|
|
"corpus workspace ownership is missing or invalid; use a new corpus root or explicit rebuild"
|
|
)
|
|
if self.workspace_id is None and isinstance(persisted, str):
|
|
self.workspace_id = persisted
|
|
return
|
|
if persisted != self.workspace_id:
|
|
raise PipelineError(
|
|
"corpus belongs to a different workspace; use a new corpus root or explicit rebuild"
|
|
)
|
|
|
|
def _protected_generations(self, workspace_root: Path) -> set[str]:
|
|
protected = {value for value in (self.store.active_generation(),) if value}
|
|
runs = workspace_root / ".tht-jobs" / "evidence" / "runs"
|
|
for checkpoint in runs.glob("*/checkpoint.json") if runs.exists() else ():
|
|
try:
|
|
state = json.loads(checkpoint.read_text(encoding="utf-8"))
|
|
if state.get("status") not in {"running", "failed"}:
|
|
continue
|
|
plan = checkpoint.parent / "artifacts" / "plan.json"
|
|
generation = json.loads(plan.read_text(encoding="utf-8")).get("generation")
|
|
if isinstance(generation, str):
|
|
protected.add(generation)
|
|
except (OSError, ValueError):
|
|
continue
|
|
return protected
|
|
|
|
def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict:
|
|
with self.store.writer_lock():
|
|
return self._gc(workspace_root=workspace_root, dry_run=dry_run)
|
|
|
|
def _gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict:
|
|
self._assert_workspace_binding()
|
|
published = self.store.published_generations()
|
|
list_vectors = getattr(self.vector_store, "list_evidence_generations", None)
|
|
vector_generations = set(list_vectors("evidence", self.workspace_id)) if list_vectors else set()
|
|
generations = sorted(set(published) | vector_generations)
|
|
job_protected = self._protected_generations(workspace_root)
|
|
active = self.store.active_generation()
|
|
rollback_count = self.retain_published_generations - 1
|
|
rollback = [generation for generation in published if generation != active]
|
|
keep = ({active} if active else set()) | set(rollback[-rollback_count:] if rollback_count else ())
|
|
fs_keep = keep | job_protected
|
|
vector_protected = set(fs_keep)
|
|
for generation in fs_keep:
|
|
try:
|
|
manifest = self.store.manifest(generation)
|
|
except (OSError, ValueError):
|
|
continue
|
|
vector_protected.update(
|
|
value for value in manifest.metadata.get("document_generations", {}).values()
|
|
if isinstance(value, str)
|
|
)
|
|
evicted, failures = [], []
|
|
filesystem_generations = set(self.store.list_generations())
|
|
for generation in generations:
|
|
purge_vector = generation not in vector_protected
|
|
purge_filesystem = generation in filesystem_generations and generation not in fs_keep
|
|
if not purge_vector and not purge_filesystem:
|
|
continue
|
|
if dry_run:
|
|
evicted.append(generation)
|
|
continue
|
|
if purge_vector:
|
|
try:
|
|
self.vector_store.delete_generation("evidence", generation, self.workspace_id)
|
|
except Exception:
|
|
failures.append({"generation": generation, "error": "vector cleanup failed"})
|
|
continue
|
|
try:
|
|
if purge_filesystem:
|
|
self.store.discard(generation)
|
|
evicted.append(generation)
|
|
except Exception:
|
|
failures.append({"generation": generation, "error": "filesystem cleanup failed"})
|
|
return {"status": "partial" if failures else "succeeded", "dry_run": dry_run,
|
|
"active_generation": self.store.active_generation(), "evicted": evicted,
|
|
"protected": sorted(vector_protected), "failures": failures}
|
|
|
|
def _discover(self) -> list[tuple[EvidenceSource, SourceObject]]:
|
|
discovered = []
|
|
seen = set()
|
|
for source in self.sources:
|
|
for item in source.discover():
|
|
if item.source_id in seen:
|
|
raise PipelineError("duplicate Evidence source identity")
|
|
seen.add(item.source_id)
|
|
discovered.append((source, item))
|
|
return sorted(discovered, key=lambda pair: pair[1].source_id)
|
|
|
|
def run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult:
|
|
with self.store.writer_lock():
|
|
self._assert_workspace_binding()
|
|
return self._run(dry_run=dry_run, resume=resume)
|
|
|
|
def run_as_job(self, **kwargs) -> PipelineResult:
|
|
self.workspace_id = kwargs["workspace_id"]
|
|
with self.store.writer_lock():
|
|
self._assert_workspace_binding()
|
|
return self._run_as_job(**kwargs)
|
|
|
|
def _run_as_job(
|
|
self,
|
|
*,
|
|
workspace_id: str,
|
|
workspace_root: Path,
|
|
config_fingerprint: str,
|
|
input_fingerprint: str,
|
|
dry_run: bool = False,
|
|
resume_run_id: str | None = None,
|
|
after_stage_return=None,
|
|
) -> PipelineResult:
|
|
"""Execute preprocessing through the durable shared job envelope."""
|
|
discovered = self._discover()
|
|
discovered_fingerprint = _fingerprint(
|
|
{item.source_id: item.fingerprint for _, item in discovered}
|
|
)
|
|
source_snapshot = _source_snapshot(discovered)
|
|
source_by_id = {item.source_id: (source, item) for source, item in discovered}
|
|
compatibility = _fingerprint({
|
|
"pipeline": self.pipeline_version,
|
|
"model": self.embedding_model,
|
|
"dimensions": self.embedding_dimensions,
|
|
"chunk_policy": asdict(self.chunk_policy),
|
|
})
|
|
job_binding = {
|
|
"config_fingerprint": config_fingerprint,
|
|
"input_fingerprint": input_fingerprint,
|
|
"compatibility_fingerprint": compatibility,
|
|
"pipeline_version": self.pipeline_version,
|
|
"chunk_policy_version": self.chunk_policy.version,
|
|
"embedding_model": self.embedding_model,
|
|
"embedding_dimensions": self.embedding_dimensions,
|
|
}
|
|
previous = self.store.active_manifest()
|
|
|
|
def document_sources(manifest: CorpusManifest) -> dict[str, dict]:
|
|
return {
|
|
document.document_id: {
|
|
"document_id": document.document_id,
|
|
"source_id": document.source_id,
|
|
"source_uri": document.source_uri,
|
|
"source_fingerprint": document.source_fingerprint,
|
|
"modified_at": (
|
|
document.modified_at.isoformat().replace("+00:00", "Z")
|
|
if document.modified_at else None
|
|
),
|
|
"source_metadata": _canonical_json(document.metadata.get("source")),
|
|
"media_type": document.media_type,
|
|
"content_hash": document.content_hash,
|
|
"pipeline_version": document.pipeline_version,
|
|
}
|
|
for document in manifest.documents
|
|
}
|
|
|
|
def active_assets_are_valid(manifest: CorpusManifest | None) -> bool:
|
|
if manifest is None or manifest.metadata.get("workspace_id") != workspace_id:
|
|
return False
|
|
actual_documents = {document.source_id: document for document in manifest.documents}
|
|
persisted_snapshot = _canonical_json(manifest.metadata.get("source_snapshot"))
|
|
if (
|
|
not isinstance(persisted_snapshot, dict)
|
|
or set(persisted_snapshot) != set(actual_documents)
|
|
or manifest.metadata.get("compatibility_fingerprint") != compatibility
|
|
or _canonical_json(manifest.metadata.get("document_sources"))
|
|
!= document_sources(manifest)
|
|
):
|
|
return False
|
|
for source_id, document in actual_documents.items():
|
|
source_payload = persisted_snapshot[source_id]
|
|
source = SourceObject.model_validate({
|
|
"source_id": source_payload["source_id"],
|
|
"uri": source_payload["uri"],
|
|
"fingerprint": source_payload["fingerprint"],
|
|
"modified_at": source_payload["modified_at"],
|
|
"metadata": source_payload["metadata"],
|
|
})
|
|
content = self.store.read_document(document.document_id, manifest.manifest_id)
|
|
expected_uri = canonical_provenance_uri(source.uri)
|
|
expected_id = "doc:" + hashlib.sha256(
|
|
f"{source.source_id}\n{expected_uri}".encode()
|
|
).hexdigest()
|
|
expected_media_type = source_payload.get("media_type")
|
|
if (
|
|
content != document.content
|
|
or document.document_id != expected_id
|
|
or document.source_id != source.source_id
|
|
or document.source_uri != expected_uri
|
|
or document.source_fingerprint != source.fingerprint
|
|
or document.modified_at != source.modified_at
|
|
or _canonical_json(document.metadata.get("source"))
|
|
!= _canonical_json(source.metadata)
|
|
or (
|
|
isinstance(expected_media_type, str)
|
|
and document.media_type != expected_media_type
|
|
)
|
|
or document.pipeline_version != self.pipeline_version
|
|
):
|
|
return False
|
|
expected_chunks = tuple(
|
|
part for document in manifest.documents for part in chunk(document, self.chunk_policy)
|
|
)
|
|
if any(document.content and not chunk(document, self.chunk_policy)
|
|
for document in manifest.documents):
|
|
return False
|
|
if _canonical_json([part.model_dump(mode="json") for part in manifest.chunks]) != (
|
|
_canonical_json([part.model_dump(mode="json") for part in expected_chunks])
|
|
):
|
|
return False
|
|
generations = manifest.metadata.get("document_generations")
|
|
if not isinstance(generations, Mapping):
|
|
return False
|
|
health = self.vector_store.health()
|
|
if (
|
|
not health.ok
|
|
or health.dimension_compatible is not True
|
|
or health.expected_dimension != self.embedding_dimensions
|
|
or health.observed_dimensions != (self.embedding_dimensions,)
|
|
):
|
|
return False
|
|
existing = self.vector_store.existing_hashes("evidence", ["evidence"])
|
|
for part in expected_chunks:
|
|
generation = generations.get(part.document_id)
|
|
if not isinstance(generation, str):
|
|
return False
|
|
record_id = f"{workspace_id}:{generation}:{part.chunk_id}"
|
|
if existing.get(record_id) != part.content_hash:
|
|
return False
|
|
return True
|
|
|
|
try:
|
|
active_assets_valid = active_assets_are_valid(previous)
|
|
except Exception:
|
|
active_assets_valid = False
|
|
reusable = (
|
|
active_assets_valid
|
|
and _canonical_json(previous.metadata.get("source_snapshot")) == source_snapshot
|
|
and _canonical_json(previous.metadata.get("job_binding")) == job_binding
|
|
)
|
|
if not dry_run and resume_run_id is None and reusable:
|
|
return PipelineResult(
|
|
"succeeded", previous.manifest_id, False, (),
|
|
tuple(sorted(item.source_id for _, item in discovered)), (), previous,
|
|
)
|
|
spec = JobSpec(
|
|
workspace_id=workspace_id,
|
|
job_type="evidence",
|
|
workspace_root=workspace_root,
|
|
spec_version="jobs-v1",
|
|
pipeline_version=self.pipeline_version,
|
|
config_fingerprint=config_fingerprint,
|
|
input_fingerprint=_fingerprint([input_fingerprint, discovered_fingerprint]),
|
|
stage_ids=EVIDENCE_STAGE_IDS,
|
|
dry_run=dry_run,
|
|
resume_run_id=resume_run_id,
|
|
)
|
|
|
|
def artifact(context: JobContext, name: str) -> Path:
|
|
root = context.run_dir / "artifacts"
|
|
root.mkdir(exist_ok=True)
|
|
return root / name
|
|
|
|
def write(context: JobContext, name: str, value) -> None:
|
|
artifact(context, name).write_text(
|
|
json.dumps(value, sort_keys=True, separators=(",", ":")), encoding="utf-8"
|
|
)
|
|
|
|
def read(context: JobContext, name: str):
|
|
try:
|
|
return json.loads(artifact(context, name).read_text(encoding="utf-8"))
|
|
except (OSError, ValueError) as error:
|
|
raise PipelineError("preprocessing checkpoint artifact is corrupt") from error
|
|
|
|
def discover_stage(context: JobContext) -> None:
|
|
previous = self.store.active_manifest()
|
|
prior = {doc.source_id: doc for doc in previous.documents} if previous else {}
|
|
fingerprints = {item.source_id: item.fingerprint for _, item in discovered}
|
|
previous_snapshot = (
|
|
_canonical_json(previous.metadata.get("source_snapshot")) if previous else {}
|
|
)
|
|
rebuild = bool(previous and not active_assets_valid)
|
|
changed = sorted(
|
|
item.source_id for _, item in discovered
|
|
if rebuild or item.source_id not in prior
|
|
or previous_snapshot.get(item.source_id) != source_snapshot[item.source_id]
|
|
)
|
|
unchanged = sorted(set(fingerprints) - set(changed))
|
|
removed = sorted(set(prior) - set(fingerprints))
|
|
write(context, "plan.json", {
|
|
"generation": f"gen:{context.run_id}",
|
|
"compatibility": compatibility,
|
|
"job_binding": job_binding,
|
|
"source_snapshot": source_snapshot,
|
|
"fingerprints": fingerprints,
|
|
"changed": changed,
|
|
"unchanged": unchanged,
|
|
"removed": removed,
|
|
"previous": previous.model_dump(mode="json") if previous else None,
|
|
})
|
|
return StageArtifacts(("plan.json",))
|
|
|
|
def acquire_stage(context: JobContext) -> None:
|
|
if context.dry_run:
|
|
return StageArtifacts()
|
|
plan = read(context, "plan.json")
|
|
previous = CorpusManifest.model_validate(plan["previous"]) if plan["previous"] else None
|
|
prior = {doc.source_id: doc for doc in previous.documents} if previous else {}
|
|
documents = [prior[source_id] for source_id in plan["unchanged"]]
|
|
for source_id in plan["changed"]:
|
|
source, item = source_by_id[source_id]
|
|
documents.append(normalize(source.acquire(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)]
|
|
previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {}
|
|
changed = set(plan["changed"])
|
|
generations = {
|
|
document.document_id: (
|
|
plan["generation"] if document.source_id in changed
|
|
else previous_generations.get(document.document_id, previous.vector_generation)
|
|
) for document in documents
|
|
}
|
|
manifest = CorpusManifest(
|
|
pipeline_version=self.pipeline_version,
|
|
embedding_model=self.embedding_model,
|
|
embedding_dimensions=self.embedding_dimensions,
|
|
vector_generation=plan["generation"],
|
|
documents=tuple(documents), chunks=tuple(chunks),
|
|
metadata={
|
|
"workspace_id": self.workspace_id,
|
|
"compatibility_fingerprint": compatibility,
|
|
"job_binding": plan["job_binding"],
|
|
"source_snapshot": plan["source_snapshot"],
|
|
"document_sources": {
|
|
document.document_id: {
|
|
"document_id": document.document_id,
|
|
"source_id": document.source_id,
|
|
"source_uri": document.source_uri,
|
|
"source_fingerprint": document.source_fingerprint,
|
|
"modified_at": (
|
|
document.modified_at.isoformat().replace("+00:00", "Z")
|
|
if document.modified_at else None
|
|
),
|
|
"source_metadata": _canonical_json(
|
|
document.metadata.get("source")
|
|
),
|
|
"media_type": document.media_type,
|
|
"content_hash": document.content_hash,
|
|
"pipeline_version": document.pipeline_version,
|
|
}
|
|
for document in documents
|
|
},
|
|
"fingerprints": plan["fingerprints"],
|
|
"removed": plan["removed"],
|
|
"document_generations": generations,
|
|
},
|
|
)
|
|
write(context, "manifest.json", manifest.model_dump(mode="json"))
|
|
return StageArtifacts(("manifest.json",))
|
|
|
|
def embed_stage(context: JobContext) -> None:
|
|
if context.dry_run:
|
|
return StageArtifacts()
|
|
plan = read(context, "plan.json")
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]}
|
|
parts = [part for part in manifest.chunks if part.document_id in changed_docs]
|
|
embeddings = self.embedder.embed_documents([part.content for part in parts])
|
|
if len(embeddings) != len(parts) or any(
|
|
len(vector) != self.embedding_dimensions for vector in embeddings
|
|
):
|
|
raise PipelineError("embedding output is incompatible")
|
|
write(context, "embeddings.json", embeddings)
|
|
return StageArtifacts(("embeddings.json",))
|
|
|
|
def records(context: JobContext):
|
|
plan = read(context, "plan.json")
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]}
|
|
parts = [part for part in manifest.chunks if part.document_id in changed_docs]
|
|
embeddings = read(context, "embeddings.json")
|
|
return [self._vector_record(part, vector, plan["generation"], self.workspace_id)
|
|
for part, vector in zip(parts, embeddings, strict=True)]
|
|
|
|
def compensate(context: JobContext) -> None:
|
|
generation = read(context, "plan.json")["generation"]
|
|
if self.store.active_generation() != generation:
|
|
self.store.discard(generation)
|
|
try:
|
|
self.vector_store.delete_generation("evidence", generation, self.workspace_id)
|
|
except Exception:
|
|
pass
|
|
write(context, "compensated.json", {"generation": generation})
|
|
|
|
def rotate_compensated_generation(context: JobContext) -> None:
|
|
marker = artifact(context, "compensated.json")
|
|
if not marker.exists():
|
|
return
|
|
plan = read(context, "plan.json")
|
|
old = plan["generation"]
|
|
plan["generation"] = f"gen:{uuid.uuid4().hex}"
|
|
write(context, "plan.json", plan)
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
changed = set(plan["changed"])
|
|
generations = dict(manifest.metadata["document_generations"])
|
|
for document in manifest.documents:
|
|
if document.source_id in changed and generations.get(document.document_id) == old:
|
|
generations[document.document_id] = plan["generation"]
|
|
manifest_payload = manifest.model_dump(mode="json")
|
|
manifest_payload["metadata"]["document_generations"] = generations
|
|
manifest_payload["vector_generation"] = plan["generation"]
|
|
manifest = CorpusManifest.model_validate(manifest_payload)
|
|
write(context, "manifest.json", manifest.model_dump(mode="json"))
|
|
marker.unlink()
|
|
|
|
def vector_stage(context: JobContext) -> None:
|
|
if context.dry_run:
|
|
return StageArtifacts()
|
|
rotate_compensated_generation(context)
|
|
values = records(context)
|
|
write(context, "vector-intent.json", {
|
|
"generation": read(context, "plan.json")["generation"],
|
|
"records": {value.record.id: value.content_hash for value in values},
|
|
})
|
|
seal_stage_artifacts(
|
|
context, "vector_upsert",
|
|
("plan.json", "manifest.json", "vector-intent.json"), spec,
|
|
)
|
|
try:
|
|
existing = self.vector_store.existing_hashes("evidence", ["evidence"])
|
|
missing = [
|
|
value for value in values
|
|
if existing.get(value.record.id) != value.content_hash
|
|
]
|
|
if missing and self.vector_store.upsert("evidence", missing) != len(missing):
|
|
raise PipelineError("vector write count mismatch")
|
|
except Exception:
|
|
compensate(context)
|
|
raise
|
|
return StageArtifacts(("plan.json", "manifest.json", "vector-intent.json"))
|
|
|
|
def stage_stage(context: JobContext) -> None:
|
|
if context.dry_run:
|
|
return StageArtifacts()
|
|
plan = read(context, "plan.json")
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
recovered = False
|
|
try:
|
|
if artifact(context, "compensated.json").exists():
|
|
recovered = True
|
|
rotate_compensated_generation(context)
|
|
values = records(context)
|
|
existing = self.vector_store.existing_hashes("evidence", ["evidence"])
|
|
missing = [value for value in values if existing.get(value.record.id) != value.content_hash]
|
|
if missing and self.vector_store.upsert("evidence", missing) != len(missing):
|
|
raise PipelineError("vector write count mismatch")
|
|
plan = read(context, "plan.json")
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
if not self.store.generation_path(plan["generation"]).exists():
|
|
self.store.stage(
|
|
manifest, {doc.document_id: doc.content for doc in manifest.documents},
|
|
generation=plan["generation"],
|
|
)
|
|
self.store.manifest(plan["generation"])
|
|
except Exception:
|
|
compensate(context)
|
|
raise
|
|
return StageArtifacts(
|
|
("plan.json", "manifest.json", "vector-intent.json") if recovered else ()
|
|
)
|
|
|
|
def publish_stage(context: JobContext) -> None:
|
|
if context.dry_run:
|
|
return StageArtifacts()
|
|
if artifact(context, "compensated.json").exists():
|
|
rotate_compensated_generation(context)
|
|
values = records(context)
|
|
try:
|
|
existing = self.vector_store.existing_hashes("evidence", ["evidence"])
|
|
missing = [value for value in values if existing.get(value.record.id) != value.content_hash]
|
|
if missing and self.vector_store.upsert("evidence", missing) != len(missing):
|
|
raise PipelineError("vector write count mismatch")
|
|
except Exception:
|
|
compensate(context)
|
|
raise
|
|
manifest = CorpusManifest.model_validate(read(context, "manifest.json"))
|
|
generation = read(context, "plan.json")["generation"]
|
|
try:
|
|
if not self.store.generation_path(generation).exists():
|
|
self.store.stage(
|
|
manifest, {doc.document_id: doc.content for doc in manifest.documents},
|
|
generation=generation,
|
|
)
|
|
except Exception:
|
|
compensate(context)
|
|
raise
|
|
generation = read(context, "plan.json")["generation"]
|
|
try:
|
|
self.store.publish(generation)
|
|
except Exception:
|
|
compensate(context)
|
|
raise
|
|
return StageArtifacts(("plan.json", "manifest.json", "vector-intent.json"))
|
|
|
|
def retention_stage(context: JobContext) -> None:
|
|
if not context.dry_run:
|
|
self.gc(workspace_root=workspace_root)
|
|
|
|
report = run_job(spec, [
|
|
discover_stage, acquire_stage, embed_stage, vector_stage,
|
|
stage_stage, publish_stage, retention_stage,
|
|
], 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())
|
|
if dry_run:
|
|
manifest = self.store.active_manifest() or CorpusManifest(pipeline_version=self.pipeline_version)
|
|
generation = None
|
|
published = False
|
|
elif report.status == "succeeded":
|
|
generation = plan["generation"]
|
|
manifest = self.store.manifest(generation)
|
|
published = True
|
|
else:
|
|
generation = plan["generation"]
|
|
manifest_path = run_dir / "artifacts" / "manifest.json"
|
|
manifest = (CorpusManifest.model_validate_json(manifest_path.read_text())
|
|
if manifest_path.exists() else CorpusManifest(pipeline_version=self.pipeline_version))
|
|
published = False
|
|
return PipelineResult(
|
|
report.status, generation, published, tuple(plan["changed"]),
|
|
tuple(plan["unchanged"]), tuple(plan["removed"]), manifest,
|
|
report.run_id, report.resumed_from,
|
|
)
|
|
|
|
def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult:
|
|
generation = None
|
|
vector_written = False
|
|
previous = self.store.active_manifest()
|
|
try:
|
|
discovered = self._discover()
|
|
except Exception as error:
|
|
raise PipelineError("Evidence discovery failed") from error
|
|
prior_documents = {doc.source_id: doc for doc in previous.documents} if previous else {}
|
|
fingerprints = {item.source_id: item.fingerprint for _, item in discovered}
|
|
compatibility = _fingerprint({
|
|
"pipeline": self.pipeline_version, "model": self.embedding_model,
|
|
"dimensions": self.embedding_dimensions, "chunk_policy": asdict(self.chunk_policy),
|
|
})
|
|
previous_compatibility = previous.metadata.get("compatibility_fingerprint") if previous else None
|
|
rebuild = previous is not None and compatibility != previous_compatibility
|
|
changed = tuple(item.source_id for _, item in discovered if rebuild or prior_documents.get(item.source_id) is None or prior_documents[item.source_id].source_fingerprint != item.fingerprint)
|
|
unchanged = tuple(item.source_id for _, item in discovered if item.source_id not in changed)
|
|
removed = tuple(sorted(set(prior_documents) - set(fingerprints)))
|
|
if dry_run:
|
|
manifest = previous or CorpusManifest(pipeline_version=self.pipeline_version)
|
|
return PipelineResult("succeeded", None, False, changed, unchanged, removed, manifest)
|
|
if previous is not None and not changed and not removed:
|
|
return PipelineResult(
|
|
"succeeded", previous.manifest_id, False, changed, unchanged, removed, previous
|
|
)
|
|
|
|
documents: list[CanonicalDocument] = [prior_documents[source_id] for source_id in unchanged]
|
|
changed_set = set(changed)
|
|
try:
|
|
for source, item in discovered:
|
|
if item.source_id in changed_set:
|
|
documents.append(normalize(source.acquire(item), self.pipeline_version))
|
|
documents.sort(key=lambda document: document.source_id)
|
|
chunks: list[CanonicalChunk] = []
|
|
for document in documents:
|
|
chunks.extend(chunk(document, self.chunk_policy))
|
|
generation = resume or f"gen:{uuid.uuid4().hex}"
|
|
previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {}
|
|
document_generations = {
|
|
document.document_id: (
|
|
generation if document.source_id in changed_set
|
|
else previous_generations.get(document.document_id, previous.vector_generation)
|
|
)
|
|
for document in documents
|
|
}
|
|
manifest = CorpusManifest(
|
|
pipeline_version=self.pipeline_version,
|
|
embedding_model=self.embedding_model,
|
|
embedding_dimensions=self.embedding_dimensions,
|
|
vector_generation=generation,
|
|
documents=tuple(documents), chunks=tuple(chunks),
|
|
metadata={
|
|
"workspace_id": self.workspace_id,
|
|
"compatibility_fingerprint": compatibility,
|
|
"fingerprints": fingerprints,
|
|
"removed": list(removed),
|
|
"document_generations": document_generations,
|
|
},
|
|
)
|
|
changed_documents = {document.document_id for document in documents if document.source_id in changed_set}
|
|
changed_chunks = [part for part in chunks if part.document_id in changed_documents]
|
|
embeddings = self.embedder.embed_documents([part.content for part in changed_chunks])
|
|
if len(embeddings) != len(changed_chunks):
|
|
raise PipelineError("embedding count mismatch")
|
|
if any(len(vector) != self.embedding_dimensions for vector in embeddings):
|
|
raise PipelineError("embedding dimension mismatch")
|
|
records = [self._vector_record(part, vector, generation, self.workspace_id) for part, vector in zip(changed_chunks, embeddings, strict=True)]
|
|
if records:
|
|
written = self.vector_store.upsert("evidence", records)
|
|
vector_written = True
|
|
if written != len(records):
|
|
raise PipelineError("vector write count mismatch")
|
|
generation_path = self.store.generation_path(generation)
|
|
if resume is not None and generation_path.exists():
|
|
staged_manifest = self.store.manifest(generation)
|
|
expected = manifest.model_dump(mode="json", exclude={"created_at", "manifest_id"})
|
|
actual = staged_manifest.model_dump(mode="json", exclude={"created_at", "manifest_id"})
|
|
actual["metadata"].pop("files", None)
|
|
if actual != expected:
|
|
raise PipelineError("resume generation is incompatible")
|
|
staged = generation
|
|
else:
|
|
staged = self.store.stage(
|
|
manifest, {document.document_id: document.content for document in documents},
|
|
generation=generation,
|
|
)
|
|
self.store.publish(staged)
|
|
self.gc(workspace_root=self.store.root.parent)
|
|
except PipelineError:
|
|
self._compensate(generation, vector_written)
|
|
raise
|
|
except Exception as error:
|
|
self._compensate(generation, vector_written)
|
|
raise PipelineError("Evidence preprocessing failed") from error
|
|
return PipelineResult("succeeded", generation, True, changed, unchanged, removed, self.store.manifest(generation))
|
|
|
|
def _compensate(self, generation: str | None, vector_written: bool) -> None:
|
|
if generation is None:
|
|
return
|
|
try:
|
|
self.store.discard(generation)
|
|
except Exception:
|
|
pass
|
|
if vector_written:
|
|
try:
|
|
self.vector_store.delete_generation("evidence", generation, self.workspace_id)
|
|
except Exception:
|
|
pass
|
|
|
|
@staticmethod
|
|
def _vector_record(
|
|
chunk: CanonicalChunk, embedding: list[float], generation: str, workspace_id: str,
|
|
):
|
|
record = VectorRecord(
|
|
id=f"{workspace_id}:{generation}:{chunk.chunk_id}",
|
|
kind="evidence", ref=chunk.document_id,
|
|
title=str(chunk.metadata.get("title", "")), content=chunk.content,
|
|
metadata={
|
|
**dict(chunk.metadata), "document_id": chunk.document_id,
|
|
"workspace_id": workspace_id,
|
|
"source_uri": chunk.source_uri, "ordinal": chunk.ordinal,
|
|
"vector_generation": generation,
|
|
},
|
|
)
|
|
return VectorWriteRecord(record=record, embedding=embedding, content_hash=chunk.content_hash)
|