Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
201 changes: 201 additions & 0 deletions docs/research/2026-08-22-failure-mining-evals-tuning-ab-research.md

Large diffs are not rendered by default.

8 changes: 8 additions & 0 deletions src/ofw/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,11 +53,15 @@
)
from ofw.diagnosis import (
ClusterId,
ClusterReview,
ClusterReviewDecision,
ClusterReviewerId,
ClusterRevisionRef,
ClusterState,
DiagnosisError,
DiagnosisErrorCode,
DiagnosisResult,
DiagnosisReview,
DiagnosisRun,
EvidenceAnchor,
EvidenceAnchorKind,
Expand Down Expand Up @@ -304,6 +308,9 @@ def promote(
"ClusterPartitionRule",
"ClusterFamilyId",
"ClusterId",
"ClusterReview",
"ClusterReviewDecision",
"ClusterReviewerId",
"ClusterState",
"ChangePrediction",
"ConsentStatus",
Expand All @@ -317,6 +324,7 @@ def promote(
"ContentCaptureMode",
"DockerCompose",
"DiagnosisResult",
"DiagnosisReview",
"DiagnosisRun",
"DataLicense",
"DiagnosisError",
Expand Down
129 changes: 126 additions & 3 deletions src/ofw/diagnosis.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
import hashlib
import math
import sys
from dataclasses import dataclass
from dataclasses import dataclass, replace
from datetime import datetime
from enum import IntEnum, StrEnum
from pathlib import Path
Expand Down Expand Up @@ -59,6 +59,12 @@ class ClusterState(StrEnum):
TARGETED = "targeted"
RESOLVED = "resolved"
REOPENED = "reopened"
REJECTED = "rejected"


class ClusterReviewDecision(StrEnum):
CONFIRM = "confirm"
REJECT = "reject"


class DiagnosisStatus(StrEnum):
Expand All @@ -70,6 +76,7 @@ class DiagnosisErrorCode(StrEnum):
STALE_HARNESS = "stale_harness"
REVISION_MISMATCH = "revision_mismatch"
ARTIFACT_INVALID = "artifact_invalid"
REVIEW_INVALID = "review_invalid"


class DiagnosisError(Exception):
Expand Down Expand Up @@ -200,6 +207,35 @@ def __str__(self) -> str:
return self.value


@dataclass(frozen=True, slots=True)
class ClusterReviewerId:
value: str

def __post_init__(self) -> None:
if not self.value.strip() or "\0" in self.value:
raise DiagnosisError(DiagnosisErrorCode.REVIEW_INVALID, self.value)

def __str__(self) -> str:
return self.value


@dataclass(frozen=True, slots=True)
class ClusterReview:
cluster_id: ClusterId
cluster_revision: int
cluster_content_digest: Sha256Digest
reviewer_id: ClusterReviewerId
decision: ClusterReviewDecision
reviewed_at: datetime

def __post_init__(self) -> None:
if self.cluster_revision < 1 or self.reviewed_at.utcoffset() is None:
raise DiagnosisError(
DiagnosisErrorCode.REVIEW_INVALID,
self.cluster_id.value,
)


@dataclass(frozen=True, slots=True)
class ClusterRevisionRef:
id: ClusterId
Expand Down Expand Up @@ -244,6 +280,7 @@ class DiagnosisResult:
diagnoses: tuple[TraceDiagnosis, ...]
clusters: tuple[FailureCluster, ...]
root: Path
reviews: tuple[ClusterReview, ...] = ()

@property
def abstained_count(self) -> int:
Expand All @@ -263,6 +300,34 @@ def to_json(self) -> str:
tuple[TraceDiagnosis, ...]
)
_RESULT_ADAPTER: TypeAdapter[DiagnosisResult] = TypeAdapter(DiagnosisResult)
_REVIEWS_ADAPTER: TypeAdapter[tuple[ClusterReview, ...]] = TypeAdapter(tuple[ClusterReview, ...])


@dataclass(frozen=True, slots=True)
class DiagnosisReview:
source: DiagnosisResult
reviews: tuple[ClusterReview, ...]

def run(self) -> DiagnosisResult:
_validate_reviews(self.source, self.reviews)
combined = (*self.source.reviews, *self.reviews)
review_digest = digest_bytes(_REVIEWS_ADAPTER.dump_json(combined))
run_id = DiagnosisRunId(
"diagnosis_"
+ hashlib.sha256(
f"{self.source.id}\0{review_digest}\0{int(DiagnosisSchemaVersion.V1)}".encode()
).hexdigest()
)
result = replace(
self.source,
id=run_id,
clusters=tuple(
_reviewed_cluster(cluster, self.reviews) for cluster in self.source.clusters
),
reviews=combined,
)
write_artifact(result.manifest_path, f"{result.to_json()}\n".encode())
return result


@dataclass(frozen=True, slots=True)
Expand Down Expand Up @@ -313,6 +378,7 @@ def run(self) -> DiagnosisResult:
diagnoses,
clusters,
revision.root,
() if self.previous is None else self.previous.reviews,
)
write_artifact(result.manifest_path, f"{result.to_json()}\n".encode())
return result
Expand All @@ -332,6 +398,53 @@ def _diagnose(
return diagnosis


def _validate_reviews(
source: DiagnosisResult,
reviews: tuple[ClusterReview, ...],
) -> None:
keys = tuple((review.cluster_id, review.cluster_revision) for review in reviews)
previous_keys = tuple((review.cluster_id, review.cluster_revision) for review in source.reviews)
if not reviews or len(set(keys)) != len(keys) or any(key in previous_keys for key in keys):
raise DiagnosisError(DiagnosisErrorCode.REVIEW_INVALID, str(source.id))
for review in reviews:
cluster = next(
(item for item in source.clusters if item.id == review.cluster_id),
None,
)
if (
cluster is None
or cluster.revision != review.cluster_revision
or cluster.content_digest != review.cluster_content_digest
or cluster.state not in (ClusterState.PROPOSED, ClusterState.REOPENED)
):
raise DiagnosisError(
DiagnosisErrorCode.REVIEW_INVALID,
review.cluster_id.value,
)


def _reviewed_cluster(
cluster: FailureCluster,
reviews: tuple[ClusterReview, ...],
) -> FailureCluster:
review = next(
(
item
for item in reviews
if item.cluster_id == cluster.id and item.cluster_revision == cluster.revision
),
None,
)
if review is None:
return cluster
state = (
ClusterState.CONFIRMED
if review.decision is ClusterReviewDecision.CONFIRM
else ClusterState.REJECTED
)
return replace(cluster, state=state)


def _clusters(
diagnoses: tuple[TraceDiagnosis, ...],
diagnoser_digest: Sha256Digest,
Expand Down Expand Up @@ -401,7 +514,17 @@ def _cluster_revision(
state = (
ClusterState.PROPOSED
if previous is None
else (ClusterState.REOPENED if previous.state is ClusterState.RESOLVED else previous.state)
else (
ClusterState.REOPENED
if previous.state
in (
ClusterState.CONFIRMED,
ClusterState.TARGETED,
ClusterState.RESOLVED,
ClusterState.REJECTED,
)
else previous.state
)
)
return FailureCluster(
_cluster_id(mechanism),
Expand All @@ -423,7 +546,7 @@ def _cluster_revision(


def _resolved_cluster(previous: FailureCluster) -> FailureCluster:
if previous.state is ClusterState.RESOLVED:
if previous.state in (ClusterState.RESOLVED, ClusterState.REJECTED):
return previous
return FailureCluster(
previous.id,
Expand Down
22 changes: 19 additions & 3 deletions src/ofw/exports.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
from pydantic import TypeAdapter

from ofw.contracts import ComponentKind, HarnessRevision, HarnessRevisionId, Sha256Digest
from ofw.diagnosis import DiagnosisResult, FailureCluster, read_snapshot
from ofw.diagnosis import ClusterState, DiagnosisResult, FailureCluster, read_snapshot
from ofw.mine import (
MineResult,
SnapshotObservation,
Expand Down Expand Up @@ -385,15 +385,32 @@ def _eligible(self) -> tuple[TraceAdmission, ...]:
def _ledger_entry(self, admission: TraceAdmission) -> LedgerEntry:
snapshot = read_snapshot(admission, self.mine)
family_id = _trace_family(snapshot)
cluster = _cluster_for_trace(self.diagnosis, admission.trace_id)
if admission.partition is TracePartition.VERIFIED_FAILURE and (
cluster is None or cluster.state not in (ClusterState.CONFIRMED, ClusterState.TARGETED)
):
return LedgerEntry(
admission.trace_id,
family_id,
None if cluster is None else ClusterFamilyId(cluster.id.value),
ExportPartition.REVIEW,
_snapshot_reference(admission),
)
previous = self._previous_partition(family_id, admission.trace_id)
if (
previous is ExportPartition.REVIEW
and admission.partition is TracePartition.VERIFIED_FAILURE
and cluster is not None
and cluster.state in (ClusterState.CONFIRMED, ClusterState.TARGETED)
):
previous = None
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,
Expand All @@ -410,7 +427,6 @@ def _ledger_entry(self, admission: TraceAdmission) -> LedgerEntry:
ExportPartition.TRAINING,
_snapshot_reference(admission),
)
cluster = _cluster_for_trace(self.diagnosis, admission.trace_id)
if cluster is None:
return LedgerEntry(
admission.trace_id,
Expand Down
Loading