diff --git a/src/ofw/__init__.py b/src/ofw/__init__.py index 1a515e3..e42c4ed 100644 --- a/src/ofw/__init__.py +++ b/src/ofw/__init__.py @@ -41,6 +41,19 @@ Severity, TraceDiagnosis, ) +from ofw.exports import ( + ClusterFamilyId, + ClusterPartitionRule, + ConsentStatus, + DataLicense, + ExportBundle, + ExportPartition, + ExportPolicy, + LeakageError, + LeakageErrorCode, + MineExports, + PrivacyTransform, +) from ofw.harness import EditableFile, Harness, Subagent, Tool, editable from ofw.mine import ( Mine, @@ -125,6 +138,7 @@ class _OfwNamespace: MiningPolicy = MiningPolicy ScoreName = ScoreName PythonDiagnoser = PythonDiagnoser + ExportPolicy = ExportPolicy def editable(self, path: Path) -> EditableFile: return editable(path) @@ -174,7 +188,10 @@ def read_snapshot_content( "AssetAccess", "CanaryCase", "CaseId", + "ClusterPartitionRule", + "ClusterFamilyId", "ClusterState", + "ConsentStatus", "ClusterRevisionRef", "ComponentKind", "CollectionError", @@ -186,11 +203,15 @@ def read_snapshot_content( "DockerCompose", "DiagnosisResult", "DiagnosisRun", + "DataLicense", "DiagnosisError", "DiagnosisErrorCode", "EditableFile", "EvidenceAnchor", "EvidenceAnchorKind", + "ExportBundle", + "ExportPartition", + "ExportPolicy", "FailureCluster", "GitCommit", "Harness", @@ -205,12 +226,15 @@ def read_snapshot_content( "LangfuseOtelSpanAttributes", "LangfuseProject", "LangfuseSpan", + "LeakageError", + "LeakageErrorCode", "LocalProcess", "ModelFingerprint", "ModuleName", "Mine", "MineError", "MineErrorCode", + "MineExports", "MiningPolicy", "MechanismKey", "ObservationContent", @@ -223,6 +247,7 @@ def read_snapshot_content( "RepositorySnapshot", "ProcessCommand", "ProcessLimits", + "PrivacyTransform", "PythonEntrypoint", "PythonLoop", "PythonDiagnoser", diff --git a/src/ofw/diagnosis.py b/src/ofw/diagnosis.py index 60ecb76..255ea58 100644 --- a/src/ofw/diagnosis.py +++ b/src/ofw/diagnosis.py @@ -287,7 +287,7 @@ def run(self) -> DiagnosisResult: prepared = environment.prepare(revision, CanaryCase(CaseId("diagnosis"), "")) try: diagnoses = tuple( - self._diagnose(_read_snapshot(admission, self.mine), prepared) + self._diagnose(read_snapshot(admission, self.mine), prepared) for admission in failures if admission.snapshot_path is not None ) @@ -490,7 +490,7 @@ def _anchor_exists(anchor: EvidenceAnchor, snapshot: TraceSnapshot) -> bool: return any(score.id.value == anchor.id for score in snapshot.scores) -def _read_snapshot(admission: TraceAdmission, mine: MineResult) -> TraceSnapshot: +def read_snapshot(admission: TraceAdmission, mine: MineResult) -> TraceSnapshot: path = admission.snapshot_path digest = admission.snapshot_digest if path is None or digest is None: diff --git a/src/ofw/exports.py b/src/ofw/exports.py new file mode 100644 index 0000000..464bd6f --- /dev/null +++ b/src/ofw/exports.py @@ -0,0 +1,615 @@ +"""Authoritative family ledger and privacy-safe Mine export manifests.""" + +from __future__ import annotations + +import hashlib +import math +from dataclasses import dataclass +from enum import IntEnum, StrEnum +from pathlib import Path + +from pydantic import TypeAdapter + +from ofw.contracts import ComponentKind, HarnessRevision, HarnessRevisionId, Sha256Digest +from ofw.diagnosis import DiagnosisResult, FailureCluster, read_snapshot +from ofw.mine import ( + MineResult, + SnapshotObservation, + TraceAdmission, + TracePartition, + TraceSnapshot, + write_artifact, +) +from ofw.observability.langfuse.domain import ScoreId, TraceId + + +class ExportSchemaVersion(IntEnum): + V1 = 1 + + +class ExportPartition(StrEnum): + TRAINING = "training" + MEMORY = "memory" + FRONTIER = "frontier" + REGRESSION = "regression" + SELECTION = "selection" + ADMISSION = "admission" + REVIEW = "review" + + +class DatasetSplit(StrEnum): + TRAIN = "train" + VALIDATION = "validation" + + +class ConsentStatus(StrEnum): + APPROVED = "approved" + DENIED = "denied" + + +class PrivacyTransform(StrEnum): + METADATA_ONLY = "metadata_only" + + +class LeakageErrorCode(StrEnum): + FAMILY_CONFLICT = "family_conflict" + REVISION_MISMATCH = "revision_mismatch" + INVALID_POLICY = "invalid_policy" + + +class LeakageError(Exception): + __slots__ = ("code", "subject") + + def __init__(self, code: LeakageErrorCode, subject: str) -> None: + self.code = code + self.subject = subject + super().__init__(f"{code.value}: {subject}") + + +@dataclass(frozen=True, slots=True) +class DataLicense: + value: str + + def __post_init__(self) -> None: + if not self.value: + raise LeakageError(LeakageErrorCode.INVALID_POLICY, "license is required") + + +@dataclass(frozen=True, slots=True) +class ClusterPartitionRule: + cluster_family_id: ClusterFamilyId + partition: ExportPartition + + +@dataclass(frozen=True, slots=True) +class ExportPolicy: + failure_partitions: tuple[ExportPartition, ...] + validation_fraction: float + license: DataLicense + consent: ConsentStatus + cluster_rules: tuple[ClusterPartitionRule, ...] = () + + def __post_init__(self) -> None: + allowed = ( + ExportPartition.MEMORY, + ExportPartition.FRONTIER, + ExportPartition.REGRESSION, + ExportPartition.SELECTION, + ExportPartition.ADMISSION, + ) + if ( + not self.failure_partitions + or any(partition not in allowed for partition in self.failure_partitions) + or any(rule.partition not in allowed for rule in self.cluster_rules) + or self.consent is not ConsentStatus.APPROVED + or not math.isfinite(self.validation_fraction) + or self.validation_fraction < 0 + or self.validation_fraction > 1 + or len({rule.cluster_family_id for rule in self.cluster_rules}) + != len(self.cluster_rules) + ): + raise LeakageError(LeakageErrorCode.INVALID_POLICY, "invalid export policy") + + @property + def digest(self) -> Sha256Digest: + return _digest_text( + "\0".join( + ( + *(partition.value for partition in self.failure_partitions), + str(self.validation_fraction), + self.license.value, + self.consent.value, + *( + f"{rule.cluster_family_id.value}:{rule.partition.value}" + for rule in self.cluster_rules + ), + str(int(ExportSchemaVersion.V1)), + ) + ) + ) + + +@dataclass(frozen=True, slots=True) +class TraceFamilyId: + value: str + + +@dataclass(frozen=True, slots=True) +class ClusterFamilyId: + value: str + + +@dataclass(frozen=True, slots=True) +class LedgerEntry: + trace_id: TraceId + trace_family_id: TraceFamilyId + cluster_family_id: ClusterFamilyId | None + partition: ExportPartition + snapshot: SnapshotReference + + +@dataclass(frozen=True, slots=True) +class PartitionLedger: + entries: tuple[LedgerEntry, ...] + + def validate(self) -> bool: + return all( + len( + { + candidate.partition + for candidate in self.entries + if candidate.trace_family_id == entry.trace_family_id + } + ) + == 1 + for entry in self.entries + ) and all( + len( + { + candidate.partition + for candidate in self.entries + if entry.cluster_family_id is not None + and candidate.cluster_family_id == entry.cluster_family_id + } + ) + == 1 + for entry in self.entries + if entry.cluster_family_id is not None + ) + + +@dataclass(frozen=True, slots=True) +class SnapshotReference: + path: Path + digest: Sha256Digest + + +@dataclass(frozen=True, slots=True) +class GoodTraceExample: + trace_id: TraceId + family_id: TraceFamilyId + snapshot: SnapshotReference + split: DatasetSplit + + +@dataclass(frozen=True, slots=True) +class GoodTraceDataset: + id: str + revision_id: HarnessRevisionId + license: DataLicense + consent: ConsentStatus + privacy_transform: PrivacyTransform + examples: tuple[GoodTraceExample, ...] + path: Path + + +@dataclass(frozen=True, slots=True) +class EvalCase: + id: str + trace_id: TraceId + family_id: TraceFamilyId + cluster_family_id: ClusterFamilyId + partition: ExportPartition + snapshot: SnapshotReference + verifier_score_ids: tuple[ScoreId, ...] + deterministic: bool = True + repeats: int = 1 + + +@dataclass(frozen=True, slots=True) +class EvalSuite: + id: str + revision_id: HarnessRevisionId + cases: tuple[EvalCase, ...] + path: Path + + +@dataclass(frozen=True, slots=True) +class MemoryCandidate: + cluster_family_id: ClusterFamilyId + mechanism: str + proposal: str + components: tuple[ComponentKind, ...] + source_trace_ids: tuple[TraceId, ...] + + +@dataclass(frozen=True, slots=True) +class MemoryPatchSet: + id: str + revision_id: HarnessRevisionId + candidates: tuple[MemoryCandidate, ...] + path: Path + + +@dataclass(frozen=True, slots=True) +class Benchmark: + id: str + revision_id: HarnessRevisionId + developer_suite_id: str + selection_suite_id: str + admission_suite_id: str + execution_digest: Sha256Digest | None + lifecycle_digest: Sha256Digest | None + path: Path + + +@dataclass(frozen=True, slots=True) +class ExportBundle: + id: str + previous_id: str | None + revision_id: HarnessRevisionId + ledger: PartitionLedger + good_traces: GoodTraceDataset + developer_evals: EvalSuite + selection_holdout: EvalSuite + admission_holdout: EvalSuite + memory: MemoryPatchSet + benchmark: Benchmark + root: Path + + @property + def manifest_path(self) -> Path: + return self.root / ".ofw" / "mine" / "exports" / self.id / "manifest.json" + + def to_json(self) -> str: + return _BUNDLE_ADAPTER.dump_json(self).decode() + + +_LEDGER_ADAPTER: TypeAdapter[PartitionLedger] = TypeAdapter(PartitionLedger) +_GOOD_ADAPTER: TypeAdapter[GoodTraceDataset] = TypeAdapter(GoodTraceDataset) +_SUITE_ADAPTER: TypeAdapter[EvalSuite] = TypeAdapter(EvalSuite) +_MEMORY_ADAPTER: TypeAdapter[MemoryPatchSet] = TypeAdapter(MemoryPatchSet) +_BENCHMARK_ADAPTER: TypeAdapter[Benchmark] = TypeAdapter(Benchmark) +_BUNDLE_ADAPTER: TypeAdapter[ExportBundle] = TypeAdapter(ExportBundle) + + +@dataclass(frozen=True, slots=True) +class MineExports: + revision: HarnessRevision + mine: MineResult + diagnosis: DiagnosisResult + policy: ExportPolicy + previous: ExportBundle | None = None + + def run(self) -> ExportBundle: + self._validate_lineage() + export_id = ( + "exports_" + + hashlib.sha256( + f"{self.mine.id}\0{self.diagnosis.id}\0{self.policy.digest}".encode() + if self.previous is None + else f"{self.mine.id}\0{self.diagnosis.id}\0{self.policy.digest}\0" + f"{self.previous.id}".encode() + ).hexdigest() + ) + root = self.revision.root / ".ofw" / "mine" / "exports" / export_id + current_entries = tuple(self._ledger_entry(admission) for admission in self._eligible()) + prior_entries = () if self.previous is None else self.previous.ledger.entries + entries = ( + tuple( + entry + for entry in prior_entries + if all(current.trace_id != entry.trace_id for current in current_entries) + ) + + current_entries + ) + ledger = PartitionLedger(entries) + if not ledger.validate(): + raise LeakageError(LeakageErrorCode.FAMILY_CONFLICT, export_id) + good = self._good_dataset(export_id, root, entries) + developer = self._suite( + export_id, + root / "developer.json", + entries, + (ExportPartition.FRONTIER, ExportPartition.REGRESSION), + ) + selection = self._suite( + export_id, + root / "selection.json", + entries, + (ExportPartition.SELECTION,), + ) + admission = self._suite( + export_id, + root / "admission.json", + entries, + (ExportPartition.ADMISSION,), + ) + memory = self._memory(export_id, root / "memory.json", entries) + runtime = self.revision.runtime + benchmark = Benchmark( + f"benchmark_{export_id[8:]}", + self.revision.id, + developer.id, + selection.id, + admission.id, + None if runtime is None else runtime.execution, + None if runtime is None else runtime.lifecycle, + root / "benchmark.json", + ) + bundle = ExportBundle( + export_id, + None if self.previous is None else self.previous.id, + self.revision.id, + ledger, + good, + developer, + selection, + admission, + memory, + benchmark, + self.revision.root, + ) + self._write(bundle) + return bundle + + def _validate_lineage(self) -> None: + if ( + self.revision.id != self.mine.revision_id + or self.revision.id != self.diagnosis.revision_id + or self.mine.id != self.diagnosis.mine_id + ): + raise LeakageError(LeakageErrorCode.REVISION_MISMATCH, str(self.revision.id)) + + def _eligible(self) -> tuple[TraceAdmission, ...]: + return tuple( + admission + for admission in self.mine.admissions + if admission.partition + in (TracePartition.VERIFIED_GOOD, TracePartition.VERIFIED_FAILURE) + and admission.snapshot_path is not None + and admission.snapshot_digest is not None + ) + + def _ledger_entry(self, admission: TraceAdmission) -> LedgerEntry: + snapshot = read_snapshot(admission, self.mine) + family_id = _trace_family(snapshot) + previous = self._previous_partition(family_id, admission.trace_id) + if previous is not None: + is_good = admission.partition is TracePartition.VERIFIED_GOOD + if is_good != (previous is ExportPartition.TRAINING): + raise LeakageError( + LeakageErrorCode.FAMILY_CONFLICT, + family_id.value, + ) + cluster = _cluster_for_trace(self.diagnosis, admission.trace_id) + cluster_family = None if cluster is None else ClusterFamilyId(cluster.id.value) + return LedgerEntry( + admission.trace_id, + family_id, + cluster_family, + previous, + _snapshot_reference(admission), + ) + if admission.partition is TracePartition.VERIFIED_GOOD: + return LedgerEntry( + admission.trace_id, + family_id, + None, + ExportPartition.TRAINING, + _snapshot_reference(admission), + ) + cluster = _cluster_for_trace(self.diagnosis, admission.trace_id) + if cluster is None: + return LedgerEntry( + admission.trace_id, + family_id, + None, + ExportPartition.REVIEW, + _snapshot_reference(admission), + ) + cluster_family = ClusterFamilyId(cluster.id.value) + partition = _failure_partition(cluster_family, self.policy) + return LedgerEntry( + admission.trace_id, + family_id, + cluster_family, + partition, + _snapshot_reference(admission), + ) + + def _previous_partition( + self, + family_id: TraceFamilyId, + trace_id: TraceId, + ) -> ExportPartition | None: + if self.previous is None: + return None + cluster = _cluster_for_trace(self.diagnosis, trace_id) + cluster_family = None if cluster is None else ClusterFamilyId(cluster.id.value) + return next( + ( + entry.partition + for entry in self.previous.ledger.entries + if entry.trace_family_id == family_id + or (cluster_family is not None and entry.cluster_family_id == cluster_family) + ), + None, + ) + + def _good_dataset( + self, + export_id: str, + root: Path, + entries: tuple[LedgerEntry, ...], + ) -> GoodTraceDataset: + examples = tuple( + GoodTraceExample( + admission.trace_id, + entry.trace_family_id, + _snapshot_reference(admission), + _split(entry.trace_family_id, self.policy.validation_fraction), + ) + for admission in self.mine.admissions + for entry in entries + if admission.trace_id == entry.trace_id + and admission.partition is TracePartition.VERIFIED_GOOD + and entry.partition is ExportPartition.TRAINING + ) + return GoodTraceDataset( + f"good_{export_id[8:]}", + self.revision.id, + self.policy.license, + self.policy.consent, + PrivacyTransform.METADATA_ONLY, + examples, + root / "good.json", + ) + + def _suite( + self, + export_id: str, + path: Path, + entries: tuple[LedgerEntry, ...], + partitions: tuple[ExportPartition, ...], + ) -> EvalSuite: + cases = tuple( + _eval_case(admission, entry) + for admission in self.mine.admissions + for entry in entries + if admission.trace_id == entry.trace_id + and admission.partition is TracePartition.VERIFIED_FAILURE + and entry.partition in partitions + ) + return EvalSuite( + f"suite_{path.stem}_{export_id[8:]}", + self.revision.id, + cases, + path, + ) + + def _memory( + self, + export_id: str, + path: Path, + entries: tuple[LedgerEntry, ...], + ) -> MemoryPatchSet: + candidates = tuple( + MemoryCandidate( + ClusterFamilyId(cluster.id.value), + cluster.mechanism.value, + cluster.description, + cluster.components, + cluster.source_trace_ids, + ) + for cluster in self.diagnosis.clusters + if any( + entry.cluster_family_id == ClusterFamilyId(cluster.id.value) + and entry.partition is ExportPartition.MEMORY + for entry in entries + ) + ) + return MemoryPatchSet( + f"memory_{export_id[8:]}", + self.revision.id, + candidates, + path, + ) + + def _write(self, bundle: ExportBundle) -> None: + root = bundle.manifest_path.parent + write_artifact(root / "ledger.json", _LEDGER_ADAPTER.dump_json(bundle.ledger) + b"\n") + write_artifact(bundle.good_traces.path, _GOOD_ADAPTER.dump_json(bundle.good_traces) + b"\n") + for suite in ( + bundle.developer_evals, + bundle.selection_holdout, + bundle.admission_holdout, + ): + write_artifact(suite.path, _SUITE_ADAPTER.dump_json(suite) + b"\n") + write_artifact(bundle.memory.path, _MEMORY_ADAPTER.dump_json(bundle.memory) + b"\n") + write_artifact( + bundle.benchmark.path, _BENCHMARK_ADAPTER.dump_json(bundle.benchmark) + b"\n" + ) + write_artifact(bundle.manifest_path, f"{bundle.to_json()}\n".encode()) + + +def _trace_family(snapshot: TraceSnapshot) -> TraceFamilyId: + payload = "\0".join( + f"{observation.type.value}:{observation.name or ''}:{observation.is_root}" + f":{observation.level}:{observation.parent_observation_id is not None}" + for observation in sorted(snapshot.observations, key=_observation_family_key) + ) + return TraceFamilyId(f"family_{hashlib.sha256(payload.encode()).hexdigest()}") + + +def _observation_family_key(observation: SnapshotObservation) -> tuple[str, str, str, str]: + return ( + observation.type.value, + observation.name or "", + str(observation.is_root), + str(observation.parent_observation_id is not None), + ) + + +def _cluster_for_trace(diagnosis: DiagnosisResult, trace_id: TraceId) -> FailureCluster | None: + return next( + (cluster for cluster in diagnosis.clusters if trace_id in cluster.source_trace_ids), + None, + ) + + +def _failure_partition( + cluster: ClusterFamilyId, + policy: ExportPolicy, +) -> ExportPartition: + explicit = next( + (rule.partition for rule in policy.cluster_rules if rule.cluster_family_id == cluster), + None, + ) + if explicit is not None: + return explicit + index = int(hashlib.sha256(cluster.value.encode()).hexdigest()[:8], 16) + return policy.failure_partitions[index % len(policy.failure_partitions)] + + +def _snapshot_reference(admission: TraceAdmission) -> SnapshotReference: + if admission.snapshot_path is None or admission.snapshot_digest is None: + raise LeakageError(LeakageErrorCode.FAMILY_CONFLICT, admission.trace_id.value) + return SnapshotReference(admission.snapshot_path, admission.snapshot_digest) + + +def _split(family: TraceFamilyId, validation_fraction: float) -> DatasetSplit: + bucket = int(hashlib.sha256(family.value.encode()).hexdigest()[:8], 16) / 0xFFFFFFFF + return DatasetSplit.VALIDATION if bucket < validation_fraction else DatasetSplit.TRAIN + + +def _eval_case(admission: TraceAdmission, entry: LedgerEntry) -> EvalCase: + cluster = entry.cluster_family_id + if cluster is None: + raise LeakageError(LeakageErrorCode.FAMILY_CONFLICT, admission.trace_id.value) + case_digest = hashlib.sha256( + (admission.trace_id.value + entry.partition.value).encode() + ).hexdigest() + return EvalCase( + f"eval_{case_digest}", + admission.trace_id, + entry.trace_family_id, + cluster, + entry.partition, + _snapshot_reference(admission), + admission.evidence_score_ids, + ) + + +def _digest_text(value: str) -> Sha256Digest: + return Sha256Digest(f"sha256:{hashlib.sha256(value.encode()).hexdigest()}") diff --git a/src/ofw/mine.py b/src/ofw/mine.py index 07bbbc5..be30686 100644 --- a/src/ofw/mine.py +++ b/src/ofw/mine.py @@ -332,7 +332,7 @@ def _admit( payload = _SNAPSHOT_ADAPTER.dump_json(snapshot) digest = digest_bytes(payload) path = revision.root / ".ofw" / "mine" / str(run_id) / "traces" / f"{digest.value[7:]}.json" - write_artifact(path, payload + b"\n") + write_artifact(path, payload) return TraceAdmission(trace.id, partition, reason, evidence, digest, path) def _snapshot_observation( diff --git a/tests/test_exports.py b/tests/test_exports.py new file mode 100644 index 0000000..a86b25e --- /dev/null +++ b/tests/test_exports.py @@ -0,0 +1,392 @@ +"""Authoritative no-leak ledger and Mine export behavior.""" + +from __future__ import annotations + +import hashlib +import subprocess +from dataclasses import replace +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest +from pydantic import TypeAdapter + +from ofw import ( + ClusterFamilyId, + ClusterPartitionRule, + ConsentStatus, + DataLicense, + ExportPartition, + ExportPolicy, + Harness, + LeakageError, + LeakageErrorCode, + MineExports, +) +from ofw.contracts import ComponentKind, HarnessRevision, Sha256Digest +from ofw.diagnosis import ( + ClusterId, + ClusterState, + DiagnosisResult, + DiagnosisRunId, + DiagnosisSchemaVersion, + FailureCluster, + MechanismKey, + Severity, +) +from ofw.mine import ( + AdmissionReason, + MineResult, + MineRunId, + MineSchemaVersion, + SnapshotObservation, + SnapshotTrace, + TraceAdmission, + TracePartition, + TraceSnapshot, +) +from ofw.observability.langfuse.contracts import TraceWindow +from ofw.observability.langfuse.domain import ( + AttributionLevel, + ObservationId, + ObservationType, + TraceId, +) + +_SNAPSHOT_ADAPTER: TypeAdapter[TraceSnapshot] = TypeAdapter(TraceSnapshot) + + +def _run_git(root: Path, *arguments: str) -> None: + subprocess.run( + ("git", "-C", str(root), *arguments), + check=True, + capture_output=True, + text=True, + ) + + +def _revision(tmp_path: Path) -> HarnessRevision: + root = tmp_path / "export-agent" + root.mkdir() + (root / "prompt.md").write_text("Be accurate.\n", encoding="utf-8") + _run_git(root, "init", "-q") + _run_git(root, "config", "user.email", "fixture@example.test") + _run_git(root, "config", "user.name", "FixtureCo") + _run_git(root, "add", ".") + _run_git(root, "commit", "-qm", "fixture baseline") + harness = Harness("export-agent", root=root) + harness.connect_prompt(Path("prompt.md")) + return harness.process() + + +def _snapshot( + revision: HarnessRevision, + mine_id: MineRunId, + trace: str, + name: str, +) -> tuple[Path, Sha256Digest]: + observation_id = ObservationId(f"observation-{trace}") + snapshot = TraceSnapshot( + MineSchemaVersion.V1, + revision.id, + Sha256Digest("sha256:collection"), + SnapshotTrace( + TraceId(trace), + (observation_id,), + (observation_id,), + (), + AttributionLevel.EXACT, + (), + Sha256Digest(f"sha256:trace-{trace}"), + ), + ( + SnapshotObservation( + observation_id, + TraceId(trace), + datetime(2026, 8, 22, tzinfo=UTC), + datetime(2026, 8, 22, 0, 1, tzinfo=UTC), + None, + ObservationType.AGENT, + True, + name, + None, + None, + Sha256Digest(f"sha256:observation-{trace}"), + ), + ), + (), + ) + payload = _SNAPSHOT_ADAPTER.dump_json(snapshot) + digest = Sha256Digest(f"sha256:{hashlib.sha256(payload).hexdigest()}") + path = revision.root / ".ofw" / "mine" / str(mine_id) / "traces" / f"{digest.value[7:]}.json" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_bytes(payload) + return path, digest + + +def _admission( + revision: HarnessRevision, + mine_id: MineRunId, + trace: str, + name: str, + partition: TracePartition, +) -> TraceAdmission: + path, digest = _snapshot(revision, mine_id, trace, name) + return TraceAdmission( + TraceId(trace), + partition, + ( + AdmissionReason.VERIFIED_PASS + if partition is TracePartition.VERIFIED_GOOD + else AdmissionReason.VERIFIED_FAIL + ), + (), + digest, + path, + ) + + +def _inputs( + tmp_path: Path, *, conflict: bool = False +) -> tuple[HarnessRevision, MineResult, DiagnosisResult]: + revision = _revision(tmp_path) + mine_id = MineRunId("mine_exports") + good_name = "shared-topology" + failure_name = good_name if conflict else "tool-topology" + admissions = ( + _admission(revision, mine_id, "good-one", good_name, TracePartition.VERIFIED_GOOD), + _admission(revision, mine_id, "good-two", good_name, TracePartition.VERIFIED_GOOD), + _admission( + revision, mine_id, "tool-failure", failure_name, TracePartition.VERIFIED_FAILURE + ), + _admission( + revision, + mine_id, + "prompt-failure", + "prompt-topology", + TracePartition.VERIFIED_FAILURE, + ), + ) + start = datetime(2026, 8, 22, tzinfo=UTC) + mine = MineResult( + MineSchemaVersion.V1, + mine_id, + revision.id, + TraceWindow(start, start + timedelta(hours=1)), + Sha256Digest("sha256:collection"), + Sha256Digest("sha256:policy"), + admissions, + revision.root, + ) + clusters = ( + _cluster("prompt-gap", TraceId("prompt-failure"), ComponentKind.PROMPT), + _cluster("tool-schema", TraceId("tool-failure"), ComponentKind.TOOL), + ) + diagnosis = DiagnosisResult( + DiagnosisSchemaVersion.V1, + DiagnosisRunId("diagnosis_exports"), + mine_id, + revision.id, + Sha256Digest("sha256:diagnoser"), + mine.window.end, + (), + clusters, + revision.root, + ) + return revision, mine, diagnosis + + +def _cluster(mechanism: str, trace: TraceId, component: ComponentKind) -> FailureCluster: + cluster_id = ClusterId(f"cluster-{mechanism}") + return FailureCluster( + cluster_id, + 1, + Sha256Digest(f"sha256:{mechanism}"), + MechanismKey(mechanism), + mechanism, + "fixture cluster", + (trace,), + (), + (component,), + 1, + Severity.HIGH, + 0.9, + 0.0, + ClusterState.CONFIRMED, + ) + + +def _policy() -> ExportPolicy: + return ExportPolicy( + failure_partitions=(ExportPartition.FRONTIER, ExportPartition.ADMISSION), + validation_fraction=0.2, + license=DataLicense("fixture-approved"), + consent=ConsentStatus.APPROVED, + cluster_rules=( + ClusterPartitionRule( + ClusterFamilyId("cluster-prompt-gap"), + ExportPartition.FRONTIER, + ), + ClusterPartitionRule( + ClusterFamilyId("cluster-tool-schema"), + ExportPartition.ADMISSION, + ), + ), + ) + + +def test_family_ledger_prevents_cross_partition_leakage(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + + bundle = MineExports(revision, mine, diagnosis, _policy()).run() + + assert len(bundle.good_traces.examples) == 2 + assert bundle.good_traces.examples[0].family_id == bundle.good_traces.examples[1].family_id + assert bundle.good_traces.examples[0].split == bundle.good_traces.examples[1].split + training_families = tuple(example.family_id for example in bundle.good_traces.examples) + eval_families = tuple(case.family_id for case in bundle.developer_evals.cases) + holdout_families = tuple(case.family_id for case in bundle.admission_holdout.cases) + assert all(family not in eval_families for family in training_families) + assert all(family not in holdout_families for family in training_families) + assert bundle.ledger.validate() + + +def test_holdout_artifacts_are_separate_from_developer_suite(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + + bundle = MineExports(revision, mine, diagnosis, _policy()).run() + + assert bundle.developer_evals.path != bundle.admission_holdout.path + assert bundle.admission_holdout.cases + assert all( + case.partition is not ExportPartition.ADMISSION for case in bundle.developer_evals.cases + ) + + +def test_export_bundle_is_idempotent_and_content_addressed(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + exports = MineExports(revision, mine, diagnosis, _policy()) + + first = exports.run() + second = exports.run() + + assert first == second + assert first.manifest_path.read_text(encoding="utf-8") == f"{first.to_json()}\n" + + +def test_same_family_requested_for_training_and_eval_fails_closed(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path, conflict=True) + + with pytest.raises(LeakageError) as raised: + MineExports(revision, mine, diagnosis, _policy()).run() + + assert raised.value.code is LeakageErrorCode.FAMILY_CONFLICT + + +def test_memory_partition_produces_proposal_without_mutating_harness(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + policy = ExportPolicy( + failure_partitions=(ExportPartition.MEMORY, ExportPartition.REGRESSION), + validation_fraction=0.2, + license=DataLicense("fixture-approved"), + consent=ConsentStatus.APPROVED, + cluster_rules=( + ClusterPartitionRule( + ClusterFamilyId("cluster-prompt-gap"), + ExportPartition.MEMORY, + ), + ClusterPartitionRule( + ClusterFamilyId("cluster-tool-schema"), + ExportPartition.REGRESSION, + ), + ), + ) + prompt_before = (revision.root / "prompt.md").read_bytes() + + bundle = MineExports(revision, mine, diagnosis, policy).run() + + assert len(bundle.memory.candidates) == 1 + assert (revision.root / "prompt.md").read_bytes() == prompt_before + + +def test_new_cluster_does_not_reassign_existing_cluster_families(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + first = MineExports(revision, mine, diagnosis, _policy()).run() + new_admission = _admission( + revision, + mine.id, + "new-failure", + "new-topology", + TracePartition.VERIFIED_FAILURE, + ) + expanded_mine = replace(mine, admissions=(*mine.admissions, new_admission)) + expanded_diagnosis = replace( + diagnosis, + clusters=( + *diagnosis.clusters, + _cluster("aaa-new-mechanism", TraceId("new-failure"), ComponentKind.TOOL), + ), + ) + + second = MineExports(revision, expanded_mine, expanded_diagnosis, _policy()).run() + + for entry in first.ledger.entries: + if entry.cluster_family_id is None: + continue + matching = next( + candidate + for candidate in second.ledger.entries + if candidate.cluster_family_id == entry.cluster_family_id + ) + assert matching.partition is entry.partition + + +def test_reordered_same_topology_cannot_cross_training_and_eval(tmp_path: Path) -> None: + revision, mine, diagnosis = _inputs(tmp_path) + good_admission = mine.admissions[0] + failed_admission = mine.admissions[2] + assert good_admission.snapshot_path is not None + assert failed_admission.snapshot_path is not None + good_snapshot = _SNAPSHOT_ADAPTER.validate_json(good_admission.snapshot_path.read_bytes()) + failed_snapshot = _SNAPSHOT_ADAPTER.validate_json(failed_admission.snapshot_path.read_bytes()) + good_root = replace(good_snapshot.observations[0], name="shared-root") + good_child = replace( + good_root, + id=ObservationId("good-child"), + parent_observation_id=good_root.id, + is_root=False, + name="shared-child", + ) + failed_root = replace(failed_snapshot.observations[0], name="shared-root") + failed_child = replace( + failed_root, + id=ObservationId("failed-child"), + parent_observation_id=failed_root.id, + is_root=False, + name="shared-child", + ) + changed_good = replace(good_snapshot, observations=(good_root, good_child)) + changed_failed = replace(failed_snapshot, observations=(failed_child, failed_root)) + good_admission = _replace_snapshot(good_admission, changed_good) + failed_admission = _replace_snapshot(failed_admission, changed_failed) + changed_mine = replace( + mine, + admissions=(good_admission, mine.admissions[1], failed_admission, mine.admissions[3]), + ) + + with pytest.raises(LeakageError) as raised: + MineExports(revision, changed_mine, diagnosis, _policy()).run() + + assert raised.value.code is LeakageErrorCode.FAMILY_CONFLICT + + +def _replace_snapshot( + admission: TraceAdmission, + snapshot: TraceSnapshot, +) -> TraceAdmission: + payload = _SNAPSHOT_ADAPTER.dump_json(snapshot) + digest = Sha256Digest(f"sha256:{hashlib.sha256(payload).hexdigest()}") + assert admission.snapshot_path is not None + path = admission.snapshot_path.with_name(f"{digest.value[7:]}.json") + path.write_bytes(payload) + return replace(admission, snapshot_digest=digest, snapshot_path=path) diff --git a/tests/test_mine.py b/tests/test_mine.py index f6f5c4a..7b1cdce 100644 --- a/tests/test_mine.py +++ b/tests/test_mine.py @@ -22,6 +22,7 @@ read_snapshot_content, ) from ofw.contracts import HarnessRevision, Sha256Digest +from ofw.diagnosis import read_snapshot from ofw.mine import TraceSnapshot from ofw.observability.langfuse.contracts import ( LangfuseConnectionId, @@ -327,6 +328,7 @@ def test_mine_is_content_addressed_and_idempotent(tmp_path: Path) -> None: snapshot = good.snapshot_path.read_text(encoding="utf-8") assert "token" not in snapshot assert "reviewed" not in snapshot + assert read_snapshot(good, first).trace.id == TraceId("good") def test_foreign_score_subject_cannot_label_trace(tmp_path: Path) -> None: