fix(evidence): harden unchanged job snapshot
This commit is contained in:
@@ -6,7 +6,9 @@ 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
|
||||
@@ -66,6 +68,31 @@ def _fingerprint(value) -> 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,
|
||||
@@ -213,6 +240,7 @@ class CorpusPipeline:
|
||||
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,
|
||||
@@ -220,12 +248,72 @@ class CorpusPipeline:
|
||||
"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()
|
||||
if (not dry_run and resume_run_id is None and previous is not None
|
||||
and previous.metadata.get("compatibility_fingerprint") == compatibility
|
||||
and previous.metadata.get("fingerprints") == {
|
||||
item.source_id: item.fingerprint for _, item in discovered
|
||||
}):
|
||||
|
||||
def document_sources(manifest: CorpusManifest) -> dict[str, dict[str, str]]:
|
||||
return {
|
||||
document.document_id: {
|
||||
"source_id": document.source_id,
|
||||
"source_uri": document.source_uri,
|
||||
"source_fingerprint": document.source_fingerprint,
|
||||
}
|
||||
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 = persisted_snapshot[source_id]
|
||||
if (
|
||||
document.source_uri != source["uri"]
|
||||
or document.source_fingerprint != source["fingerprint"]
|
||||
or self.store.read_document(document.document_id, manifest.manifest_id)
|
||||
!= document.content
|
||||
):
|
||||
return False
|
||||
generations = manifest.metadata.get("document_generations")
|
||||
if not isinstance(generations, Mapping):
|
||||
return False
|
||||
existing = self.vector_store.existing_hashes("evidence", ["evidence"])
|
||||
for part in manifest.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,
|
||||
@@ -263,17 +351,22 @@ class CorpusPipeline:
|
||||
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}
|
||||
rebuild = bool(previous and previous.metadata.get("compatibility_fingerprint") != compatibility)
|
||||
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 prior[item.source_id].source_fingerprint != item.fingerprint
|
||||
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,
|
||||
@@ -311,6 +404,16 @@ class CorpusPipeline:
|
||||
metadata={
|
||||
"workspace_id": self.workspace_id,
|
||||
"compatibility_fingerprint": compatibility,
|
||||
"job_binding": plan["job_binding"],
|
||||
"source_snapshot": plan["source_snapshot"],
|
||||
"document_sources": {
|
||||
document.document_id: {
|
||||
"source_id": document.source_id,
|
||||
"source_uri": document.source_uri,
|
||||
"source_fingerprint": document.source_fingerprint,
|
||||
}
|
||||
for document in documents
|
||||
},
|
||||
"fingerprints": plan["fingerprints"],
|
||||
"removed": plan["removed"],
|
||||
"document_generations": generations,
|
||||
|
||||
Reference in New Issue
Block a user