igerasimov's picture
Deploy Phase 2 dataset classifier (part 19)
58c2da3 verified
Raw History Blame Contribute Delete
20.5 kB
"""Verified post-blind expert comparison with no model or mutation capability."""
from __future__ import annotations
import hashlib
import json
import logging
from collections import Counter, defaultdict
from typing import Protocol
from gcmd_classifier.datasets.comparison_models import (
ComparisonCategoryCounts,
ComparisonComponentHashes,
ComparisonRelationship,
ComparisonRunStatus,
DatasetComparisonItem,
DatasetComparisonResult,
ExpertOnlyAssignment,
HierarchyRelationshipKind,
HumanReviewRecord,
PersistedComparisonOutcome,
ProposalType,
)
from gcmd_classifier.datasets.errors import DatasetComparisonError
from gcmd_classifier.datasets.expert_resolution import resolve_expert_assignments
from gcmd_classifier.datasets.models import (
ComparisonCategory,
DatasetClassificationOutcome,
DatasetClassificationResult,
DatasetReviewStatus,
)
from gcmd_classifier.datasets.runtime_models import (
ArtifactType,
ExpertSourceBinding,
SealedExpertArtifactEnvelope,
VerifiedPersistedBlindResultReference,
)
from gcmd_classifier.datasets.runtime_persistence import DatasetArtifactStore
from gcmd_classifier.logging_config import get_logger, log_event
from gcmd_classifier.vocabulary.index import VocabularyIndex
COMPARISON_POLICY = "dataset-comparison-policy-v1"
RESOLUTION_POLICY = "expert-resolution-policy-v1"
NORMALIZATION_POLICY = "case-whitespace-normalization-v1"
REVIEW_POLICY = "dataset-initial-review-v1"
class ExpertUnsealer(Protocol):
def unseal(self, persisted_bytes: bytes) -> SealedExpertArtifactEnvelope:
"""Open typed expert bytes only after prerequisite verification."""
class TypedExpertUnsealer:
def unseal(self, persisted_bytes: bytes) -> SealedExpertArtifactEnvelope:
return SealedExpertArtifactEnvelope.model_validate_json(persisted_bytes)
def _json_default(value: object) -> object:
if hasattr(value, "model_dump"):
return value.model_dump(mode="json")
if hasattr(value, "value"):
return value.value
raise TypeError(f"Unsupported canonical comparison value: {type(value).__name__}")
def _canonical(value: object) -> bytes:
return json.dumps(
value,
sort_keys=True,
separators=(",", ":"),
ensure_ascii=False,
default=_json_default,
).encode()
def _hash(value: object) -> str:
return hashlib.sha256(_canonical(value)).hexdigest()
def _distance(vocabulary: VocabularyIndex, ancestor: str, descendant: str) -> int:
ancestors = vocabulary.ancestors_of(descendant)
return ancestors.index(ancestor) + 1
def _verified_inputs(
*,
store: DatasetArtifactStore,
session_id: str,
verified: VerifiedPersistedBlindResultReference,
vocabulary: VocabularyIndex,
unsealer: ExpertUnsealer,
) -> tuple[DatasetClassificationResult, SealedExpertArtifactEnvelope]:
if not isinstance(verified, VerifiedPersistedBlindResultReference):
raise DatasetComparisonError(
"COMPARISON_PREREQUISITE_INVALID", "A verified blind reference is required."
)
run = verified.run
blind_bytes = store.load_verified(run, session_id, verified.artifact)
cmr_bytes = store.load_verified(run, session_id, verified.cmr_source_artifact)
sealed_bytes = store.load_verified(run, session_id, verified.sealed_expert_artifact)
binding_bytes = store.load_verified(run, session_id, verified.expert_source_binding_artifact)
try:
blind = DatasetClassificationResult.model_validate_json(blind_bytes)
binding = ExpertSourceBinding.model_validate_json(binding_bytes)
except Exception as exc:
raise DatasetComparisonError(
"COMPARISON_PREREQUISITE_INVALID", "Persisted blind or binding schema is invalid."
) from exc
if (
verified.blind_result_sha256 != verified.artifact.sha256
or verified.cmr_source_sha256 != verified.cmr_source_artifact.sha256
or verified.sealed_expert_artifact_sha256 != verified.sealed_expert_artifact.sha256
or verified.expert_source_binding_sha256 != verified.expert_source_binding_artifact.sha256
or blind.identity != verified.identity
or blind.dataset_key != verified.identity.concept_id
or blind.evidence_packet_sha256 != verified.evidence_packet_sha256
or blind.processing_status.value != "completed"
or blind.classification_outcome is None
or blind.classification_outcome.value != verified.completed_blind_status
or binding.identity != verified.identity
or binding.run != verified.run
or binding.cmr_source_sha256 != verified.cmr_source_sha256
or binding.sealed_expert_artifact_sha256 != verified.sealed_expert_artifact_sha256
or binding.sealed_payload_sha256 != verified.sealed_payload_sha256
or vocabulary.vocabulary_version != verified.vocabulary_hash
):
raise DatasetComparisonError(
"COMPARISON_PROVENANCE_INVALID",
"Blind, expert, or hierarchy provenance does not agree.",
)
try:
source = json.loads(cmr_bytes)
item = source["items"][0]
exact = (
item["meta"]["concept-id"],
item["meta"]["native-id"],
item["umm"]["ShortName"],
item["umm"]["Version"],
item["meta"]["revision-id"],
)
except (ValueError, KeyError, IndexError, TypeError) as exc:
raise DatasetComparisonError(
"EXPERT_SOURCE_INVALID", "Persisted CMR source is invalid."
) from exc
identity = verified.identity
if exact != (
identity.concept_id,
identity.native_id,
identity.short_name,
identity.version,
identity.cmr_revision_id,
):
raise DatasetComparisonError("EXPERT_SOURCE_IDENTITY_MISMATCH", "CMR identity differs.")
# Authorization occurs only after every non-expert prerequisite above succeeds.
try:
envelope = unsealer.unseal(sealed_bytes)
except Exception as exc:
raise DatasetComparisonError(
"EXPERT_UNSEAL_FAILED", "Sealed expert artifact is invalid."
) from exc
if (
envelope.identity != identity
or envelope.cmr_source_sha256 != verified.cmr_source_sha256
or envelope.sealed_payload_sha256 != verified.sealed_payload_sha256
or tuple(item["umm"].get("ScienceKeywords", ())) != envelope.sealed_payload.science_keywords
):
raise DatasetComparisonError(
"EXPERT_SOURCE_BINDING_INVALID", "Unsealed expert source differs."
)
return blind, envelope
def _relationships(blind_uuid, experts, vocabulary):
by_uuid: dict[str, list[int]] = defaultdict(list)
for item in experts:
by_uuid[item.UUID].append(item.source_position)
relationships = []
for expert_uuid in sorted(by_uuid):
positions = tuple(by_uuid[expert_uuid])
if blind_uuid == expert_uuid:
kind = HierarchyRelationshipKind.EXACT
distance = 0
path = vocabulary.get(blind_uuid).path_components
elif vocabulary.is_descendant(blind_uuid, expert_uuid):
kind = HierarchyRelationshipKind.BLIND_DESCENDANT
distance = _distance(vocabulary, expert_uuid, blind_uuid)
path = vocabulary.get(blind_uuid).path_components
elif vocabulary.is_descendant(expert_uuid, blind_uuid):
kind = HierarchyRelationshipKind.EXPERT_DESCENDANT
distance = _distance(vocabulary, blind_uuid, expert_uuid)
path = vocabulary.get(expert_uuid).path_components
else:
kind = HierarchyRelationshipKind.INDEPENDENT
distance = None
path = ()
relationships.append(
ComparisonRelationship(
relationship=kind,
expert_UUID=expert_uuid,
expert_source_positions=positions,
distance=distance,
authoritative_path=path,
)
)
return tuple(relationships)
def _item(index, blind, experts, unresolved, vocabulary):
relationships = list(_relationships(blind.UUID, experts, vocabulary))
if unresolved:
relationships.extend(
ComparisonRelationship(
relationship=HierarchyRelationshipKind.UNRESOLVED,
expert_source_positions=(item.source_position,),
)
for item in unresolved
)
category = ComparisonCategory.REVIEW_REQUIRED
primary = HierarchyRelationshipKind.UNRESOLVED
proposal = ProposalType.NONE
review = True
reasons = tuple(sorted({item.resolution_status.value for item in unresolved}))
rationale = "Unresolved expert assignments prevent safe deterministic categorization."
else:
kinds = {item.relationship for item in relationships}
if HierarchyRelationshipKind.EXACT in kinds:
category = ComparisonCategory.ALREADY_ASSIGNED
primary = HierarchyRelationshipKind.EXACT
rationale = "Blind terminal UUID exactly matches an existing expert assignment."
elif HierarchyRelationshipKind.EXPERT_DESCENDANT in kinds:
category = ComparisonCategory.REDUNDANT_RESULT
primary = HierarchyRelationshipKind.EXPERT_DESCENDANT
rationale = "A more specific expert descendant already exists; no change is proposed."
elif HierarchyRelationshipKind.BLIND_DESCENDANT in kinds:
category = ComparisonCategory.PROPOSED_REFINEMENT
primary = HierarchyRelationshipKind.BLIND_DESCENDANT
rationale = (
"Blind assignment is a strict descendant of an existing broader expert assignment."
)
else:
category = ComparisonCategory.PROPOSED_ADDITION
primary = HierarchyRelationshipKind.INDEPENDENT
rationale = "Blind assignment is independent of every resolved expert assignment."
proposal = {
ComparisonCategory.PROPOSED_ADDITION: ProposalType.ADDITION,
ComparisonCategory.PROPOSED_REFINEMENT: ProposalType.REFINEMENT,
}.get(category, ProposalType.NONE)
review = category in {
ComparisonCategory.PROPOSED_ADDITION,
ComparisonCategory.PROPOSED_REFINEMENT,
}
reasons = (category.value,) if review else ()
related = tuple(
sorted({position for item in relationships for position in item.expert_source_positions})
)
material = {
"blind_index": index,
"blind_uuid": blind.UUID,
"category": category.value,
"relationships": [item.model_dump(mode="json") for item in relationships],
}
return DatasetComparisonItem(
comparison_item_id=_hash(material),
blind_index=index,
blind_classification=blind,
category=category.value,
primary_relationship=primary,
relationships=tuple(relationships),
related_expert_source_positions=related,
proposal_type=proposal,
rationale=rationale,
review_required=review,
review_reason_codes=reasons,
)
def _comparison_hash_material(result: DatasetComparisonResult) -> dict:
return result.model_dump(mode="json", exclude={"comparison_sha256", "comparison_artifact"})
def compare_verified_dataset(
*,
store: DatasetArtifactStore,
session_id: str,
verified: VerifiedPersistedBlindResultReference,
vocabulary: VocabularyIndex,
unsealer: ExpertUnsealer | None = None,
logger: logging.Logger | None = None,
) -> PersistedComparisonOutcome:
"""Resolve, compare, persist, and verify advisory records with zero model calls."""
blind, envelope = _verified_inputs(
store=store,
session_id=session_id,
verified=verified,
vocabulary=vocabulary,
unsealer=unsealer or TypedExpertUnsealer(),
)
resolved, unresolved, duplicates = resolve_expert_assignments(
envelope.sealed_payload.science_keywords, vocabulary
)
items = tuple(
_item(index, item, resolved, unresolved, vocabulary)
for index, item in enumerate(blind.classifications)
)
counts = Counter(item.category for item in items)
category_counts = ComparisonCategoryCounts(
**{name: counts[name] for name in ComparisonCategoryCounts.model_fields}
)
related_uuids = {
relationship.expert_UUID
for item in items
for relationship in item.relationships
if relationship.expert_UUID
and relationship.relationship is not HierarchyRelationshipKind.INDEPENDENT
}
positions: dict[str, list[int]] = defaultdict(list)
paths = {}
for item in resolved:
positions[item.UUID].append(item.source_position)
paths[item.UUID] = item.canonical_path
expert_only = tuple(
ExpertOnlyAssignment(
UUID=uuid, source_positions=tuple(positions[uuid]), canonical_path=paths[uuid]
)
for uuid in sorted(positions)
if uuid not in related_uuids
)
policy_identity = {
"comparison": COMPARISON_POLICY,
"resolution": RESOLUTION_POLICY,
"normalization": NORMALIZATION_POLICY,
"review": REVIEW_POLICY,
"schema": "dataset-comparison-v2",
"vocabulary": vocabulary.vocabulary_version,
}
component_hashes = ComparisonComponentHashes(
blind_input_sha256=verified.blind_result_sha256,
expert_input_sha256=verified.sealed_payload_sha256,
resolved_experts_sha256=_hash(
{"resolved": resolved, "unresolved": unresolved, "duplicates": duplicates}
),
comparison_items_sha256=_hash(items),
policy_identity_sha256=_hash(policy_identity),
)
comparison_id = _hash(
{
"blind": verified.blind_result_sha256,
"expert": verified.sealed_payload_sha256,
"vocabulary": vocabulary.vocabulary_version,
**policy_identity,
}
)
review_reasons = tuple(
sorted(
{reason for item in items for reason in item.review_reason_codes}
| {item.resolution_status.value for item in unresolved}
)
)
draft = DatasetComparisonResult(
comparison_id=comparison_id,
comparison_sha256="0" * 64,
comparison_policy_version=COMPARISON_POLICY,
expert_resolution_policy_version=RESOLUTION_POLICY,
normalization_policy_version=NORMALIZATION_POLICY,
review_policy_version=REVIEW_POLICY,
run_id=verified.run.run_id,
identity=verified.identity,
verified_blind_result=verified,
blind_result_sha256=verified.blind_result_sha256,
evidence_packet_sha256=verified.evidence_packet_sha256,
expert_source_artifact=verified.cmr_source_artifact,
expert_source_sha256=verified.cmr_source_sha256,
sealed_expert_artifact=verified.sealed_expert_artifact,
sealed_payload_sha256=verified.sealed_payload_sha256,
vocabulary_hash=vocabulary.vocabulary_version,
comparison_status=(
ComparisonRunStatus.BLIND_NOT_CLASSIFIED
if blind.classification_outcome is DatasetClassificationOutcome.NOT_CLASSIFIED
else ComparisonRunStatus.BLIND_CLASSIFIED
),
resolved_expert_assignments=resolved,
unresolved_expert_assignments=unresolved,
expert_duplicate_groups=duplicates,
comparison_items=items,
expert_only_assignments=expert_only,
category_counts=category_counts,
review_status=(
DatasetReviewStatus.PENDING if review_reasons else DatasetReviewStatus.NOT_REQUIRED
),
review_required_reasons=review_reasons,
component_hashes=component_hashes,
)
comparison = draft.model_copy(
update={"comparison_sha256": _hash(_comparison_hash_material(draft))}
)
reference = store.persist_json(
run=verified.run,
session_id=session_id,
relative_name="comparison.json",
value=comparison,
artifact_type=ArtifactType.COMPARISON,
schema_name="dataset-comparison-v2",
upstream_sha256=(verified.blind_result_sha256, verified.sealed_payload_sha256),
validate=lambda data: DatasetComparisonResult.model_validate_json(data),
)
persisted = DatasetComparisonResult.model_validate_json(
store.load_verified(verified.run, session_id, reference)
)
if persisted.comparison_sha256 != _hash(_comparison_hash_material(persisted)):
raise DatasetComparisonError(
"COMPARISON_HASH_INVALID", "Persisted comparison hash is invalid."
)
item_review_records = tuple(
HumanReviewRecord(
review_record_id=_hash(
{"comparison": comparison.comparison_sha256, "item": item.comparison_item_id}
),
comparison_sha256=comparison.comparison_sha256,
comparison_item_id=item.comparison_item_id,
identity=comparison.identity,
blind_classification=item.blind_classification,
related_expert_source_positions=item.related_expert_source_positions,
comparison_category=item.category,
proposal_type=item.proposal_type,
relationships=item.relationships,
review_reason_codes=item.review_reason_codes,
comparison_policy_version=COMPARISON_POLICY,
review_policy_version=REVIEW_POLICY,
blind_result_sha256=verified.blind_result_sha256,
sealed_payload_sha256=verified.sealed_payload_sha256,
)
for item in items
if item.review_required
)
unresolved_review_records = tuple(
HumanReviewRecord(
review_record_id=_hash(
{
"comparison": comparison.comparison_sha256,
"unresolved_source_position": item.source_position,
}
),
comparison_sha256=comparison.comparison_sha256,
comparison_item_id=_hash({"unresolved_source_position": item.source_position}),
identity=comparison.identity,
related_expert_source_positions=(item.source_position,),
comparison_category=ComparisonCategory.REVIEW_REQUIRED,
proposal_type=ProposalType.NONE,
relationships=(
ComparisonRelationship(
relationship=HierarchyRelationshipKind.UNRESOLVED,
expert_source_positions=(item.source_position,),
),
),
review_reason_codes=(item.resolution_status.value,),
comparison_policy_version=COMPARISON_POLICY,
review_policy_version=REVIEW_POLICY,
blind_result_sha256=verified.blind_result_sha256,
sealed_payload_sha256=verified.sealed_payload_sha256,
)
for item in unresolved
)
review_records = (*item_review_records, *unresolved_review_records)
review_refs = tuple(
store.persist_json(
run=verified.run,
session_id=session_id,
relative_name=f"reviews/{record.review_record_id}.json",
value=record,
artifact_type=ArtifactType.HUMAN_REVIEW,
schema_name="dataset-human-review-v1",
upstream_sha256=(comparison.comparison_sha256,),
validate=lambda data: HumanReviewRecord.model_validate_json(data),
)
for record in review_records
)
log_event(
logger or get_logger("gcmd_classifier.datasets.comparison"),
"dataset_comparison_completed",
run_id=verified.run.run_id,
concept_id=verified.identity.concept_id,
blind_result_sha256=verified.blind_result_sha256,
expert_source_sha256=verified.cmr_source_sha256,
hierarchy_hash=vocabulary.vocabulary_version,
comparison_sha256=comparison.comparison_sha256,
resolved_count=len(resolved),
unresolved_count=len(unresolved),
duplicate_group_count=len(duplicates),
category_counts=category_counts.model_dump(),
review_record_count=len(review_records),
comparison_reference=reference.relative_path,
)
return PersistedComparisonOutcome(
comparison=comparison,
comparison_artifact=reference,
review_records=review_records,
review_artifacts=review_refs,
)