diff --git a/plugins/openflywheel/.codex-plugin/plugin.json b/plugins/openflywheel/.codex-plugin/plugin.json index 3a64262..7d58950 100644 --- a/plugins/openflywheel/.codex-plugin/plugin.json +++ b/plugins/openflywheel/.codex-plugin/plugin.json @@ -1,7 +1,7 @@ { "name": "openflywheel", - "version": "0.6.0", - "description": "Initialize ITSM-bench harness workspaces, query Langfuse trajectories, and mine compact failure diagnoses and exact patterns.", + "version": "0.7.0", + "description": "Initialize ITSM-bench workspaces, mine exact failure patterns, and record evidence-bound curations.", "author": { "name": "OpenFlyWheel" }, @@ -10,8 +10,8 @@ "skills": "./skills/", "interface": { "displayName": "OpenFlyWheel", - "shortDescription": "Diagnose and mine ITSM failure patterns", - "longDescription": "Initialize an ITSM-bench agent-harness optimization workspace, inspect bounded Langfuse evidence, record authoritative outcomes and compact diagnoses, and mine exact recurring patterns without copying trace payloads.", + "shortDescription": "Mine and curate ITSM failure patterns", + "longDescription": "Initialize an ITSM-bench agent-harness optimization workspace, inspect bounded Langfuse evidence, record authoritative outcomes and compact diagnoses, mine exact recurring patterns, and persist evidence-bound cross-task curations without copying trace payloads.", "developerName": "OpenFlyWheel", "category": "Productivity", "capabilities": ["Read", "Write"], diff --git a/plugins/openflywheel/.mcp.json b/plugins/openflywheel/.mcp.json index 2280c7c..c0fcfae 100644 --- a/plugins/openflywheel/.mcp.json +++ b/plugins/openflywheel/.mcp.json @@ -4,7 +4,7 @@ "command": "uvx", "args": [ "--from", - "git+https://github.com/divo12/OpenFlyWheel.git@ab0ef62cbe1e6cddf0bfd8ec61374d10120c61aa", + "git+https://github.com/divo12/OpenFlyWheel.git@726594a56c6683d1d32207a462fe524b202f4b8b", "--with", "mcp>=1.13,<2", "openflywheel-mcp" diff --git a/plugins/openflywheel/program_templates/itsm.md b/plugins/openflywheel/program_templates/itsm.md index 9174db3..b7f328b 100644 --- a/plugins/openflywheel/program_templates/itsm.md +++ b/plugins/openflywheel/program_templates/itsm.md @@ -49,6 +49,14 @@ After recording the diagnoses for one bounded run, use `$failure-pattern-miner` type and exact normalized root cause; results are not semantic clusters. Keep inconclusive diagnoses separate, and reread each supporting diagnosis before forming a shared hypothesis. +After every failed outcome has a retained artifact, use `$failure-curator` for one bounded +cross-task debugger pass. Use only the current run's retained `record_failure` receipt IDs. +Do not glob `.workspace/failures/`, which may contain earlier runs. Submit the complete set of at +most 50 IDs to `record_failure_curation`, assigning each exactly once to a repeated +evidence-bound group or a deferred entry. Retain the `.workspace/failure-curations/` artifact +before forming a harness hypothesis. If the run has more than 50 failure artifacts, report the +curation bound and stop before forming a harness hypothesis. + An intermediate tool error is evidence, not an outcome failure, when the agent recovered and the verifier passed. A technically clean trajectory is still a failure when the ITSM verifier shows that the required environment state was not achieved. diff --git a/plugins/openflywheel/skills/failure-curator/SKILL.md b/plugins/openflywheel/skills/failure-curator/SKILL.md new file mode 100644 index 0000000..0f1e348 --- /dev/null +++ b/plugins/openflywheel/skills/failure-curator/SKILL.md @@ -0,0 +1,33 @@ +--- +name: failure-curator +description: Group a completed bounded set of recorded failure diagnoses into evidence-bound cross-task patterns before forming a harness hypothesis. Use after failure-miner has recorded every verifier-backed failure; do not query traces, edit the harness, or evaluate candidates. +--- + +# Failure Curator + +Turn only the retained `record_failure` receipt IDs produced for the current completed baseline +or candidate run into a compact cross-task debugger overview. Do not glob or enumerate +`.workspace/failures/`; it may also contain artifacts from earlier runs. Langfuse remains the +source of trace content; read only the current run's recorded diagnosis artifacts in this phase. + +## Curate + +1. Read every retained failure artifact for the run, up to the tool's 50-artifact bound. If + the run exceeds that bound, stop and report that complete curation is unavailable. +2. Defer every inconclusive diagnosis with the specific evidence or recurrence gap. A + supported diagnosis may also be deferred when it has no matching failure from another task. +3. Group supported diagnoses only when at least two distinct task IDs share the same causal + mechanism and top-level `issue_type`. Shared entities, words, tools, or symptoms alone do + not establish a pattern. +4. Give each group a stable lowercase `pattern_key`, concise title, causal mechanism, general + prevention mechanism, and exactly one most relevant harness `target_component`. Keep the + prevention mechanism structural rather than task-specific; do not describe a file edit yet. +5. Assign every source artifact exactly once, either to one group or one deferred entry. + +Call `record_failure_curation` once with the prepared worktree as `workspace_root`, the full +source artifact ID set, the proposed groups, and deferred entries. Retain the returned +`.workspace/failure-curations/.json` path for the hypothesis phase. + +Do not call trace tools, copy trace content, invent missing causal evidence, combine different +failure types, form a repair hypothesis, edit the harness, run a benchmark, or promote a change +while following this skill. diff --git a/pyproject.toml b/pyproject.toml index dd61077..3721a01 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "openflywheel" -version = "0.6.0" +version = "0.7.0" description = "A governed self-improving agent harness" requires-python = ">=3.11" dependencies = [ diff --git a/src/ofw/__init__.py b/src/ofw/__init__.py index d34c096..86a08ae 100644 --- a/src/ofw/__init__.py +++ b/src/ofw/__init__.py @@ -27,16 +27,23 @@ WorkspaceFile, ) from ofw.evaluation import ( + DeferredFailure, + FailureCuration, + FailureCurationErrorCode, + FailureCurationFailure, FailureDiagnosis, FailureDiagnosisError, FailureErrorCode, FailureEvidenceStatus, + FailureGroup, + FailureGroupMember, FailurePatternMiningError, FailurePatternMiningErrorCode, FailurePatternMiningObservation, FailurePatternMiningStatus, FailurePatternOrdering, FailurePatternSummary, + FailureSource, FailureType, LangfuseOutcomeStore, MineFailurePatternsInput, @@ -109,19 +116,26 @@ def editable(self, path: Path) -> EditableFile: "CollectionErrorCode", "CommandLoop", "CommandVerifier", + "DeferredFailure", "E2BSandbox", "EditableFile", "EvidenceReference", + "FailureCuration", + "FailureCurationErrorCode", + "FailureCurationFailure", "FailureDiagnosis", "FailureDiagnosisError", "FailureErrorCode", "FailureEvidenceStatus", + "FailureGroup", + "FailureGroupMember", "FailurePatternMiningError", "FailurePatternMiningErrorCode", "FailurePatternMiningObservation", "FailurePatternMiningStatus", "FailurePatternOrdering", "FailurePatternSummary", + "FailureSource", "FailureType", "GitCommit", "Harness", diff --git a/src/ofw/evaluation/__init__.py b/src/ofw/evaluation/__init__.py index acf1962..ce95c89 100644 --- a/src/ofw/evaluation/__init__.py +++ b/src/ofw/evaluation/__init__.py @@ -7,6 +7,15 @@ FailureEvidenceStatus, FailureType, ) +from ofw.evaluation.failure_curation import ( + DeferredFailure, + FailureCuration, + FailureCurationErrorCode, + FailureCurationFailure, + FailureGroup, + FailureGroupMember, + FailureSource, +) from ofw.evaluation.failure_patterns import ( FailurePatternMiningError, FailurePatternMiningErrorCode, @@ -31,16 +40,23 @@ ) __all__ = [ + "DeferredFailure", + "FailureCuration", + "FailureCurationErrorCode", + "FailureCurationFailure", "FailureDiagnosis", "FailureDiagnosisError", "FailureErrorCode", "FailureEvidenceStatus", + "FailureGroup", + "FailureGroupMember", "FailurePatternMiningError", "FailurePatternMiningErrorCode", "FailurePatternMiningObservation", "FailurePatternMiningStatus", "FailurePatternOrdering", "FailurePatternSummary", + "FailureSource", "FailureType", "LangfuseOutcomeStore", "MineFailurePatternsInput", diff --git a/src/ofw/evaluation/failure_curation.py b/src/ofw/evaluation/failure_curation.py new file mode 100644 index 0000000..cbc4e7d --- /dev/null +++ b/src/ofw/evaluation/failure_curation.py @@ -0,0 +1,551 @@ +"""Typed, evidence-bound grouping of recorded failure diagnoses.""" + +from __future__ import annotations + +from dataclasses import dataclass +from enum import StrEnum +from pathlib import Path +from typing import Annotated, Literal, Protocol +from uuid import NAMESPACE_URL, uuid5 + +from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator + +from ofw.contracts import ComponentKind, Sha256Digest +from ofw.evaluation.failure import FailureEvidenceStatus, FailureType +from ofw.evaluation.outcome import TaskId +from ofw.observability.langfuse.domain import ObservationId, ScoreId, TraceId + +_ARTIFACT_ID_PATTERN = r"^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$" +_PATTERN_KEY_PATTERN = r"[a-z0-9]+(?:-[a-z0-9]+)*" +_TEXT_PATTERN = r"[^\x00]+" + +ArtifactIdentifier = Annotated[str, Field(pattern=_ARTIFACT_ID_PATTERN)] +DigestValue = Annotated[str, Field(pattern=r"sha256:[0-9a-f]{64}")] +PatternKey = Annotated[ + str, + Field(min_length=1, max_length=80, pattern=_PATTERN_KEY_PATTERN), +] +TitleText = Annotated[str, Field(min_length=1, max_length=160, pattern=_TEXT_PATTERN)] +ExplanationText = Annotated[str, Field(min_length=1, max_length=1000, pattern=_TEXT_PATTERN)] +SourceArtifactIds = Annotated[ + tuple[ArtifactIdentifier, ...], + Field(min_length=1, max_length=50), +] +GroupArtifactIds = Annotated[ + tuple[ArtifactIdentifier, ...], + Field(min_length=2, max_length=20), +] + + +class StrictModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True, strict=True) + + +class FailureGroupInput(StrictModel): + pattern_key: PatternKey + title: TitleText + mechanism: ExplanationText + prevention: ExplanationText + target_component: ComponentKind + failure_artifact_ids: GroupArtifactIds + + @field_validator("title", "mechanism", "prevention") + @classmethod + def normalize_text(cls, value: str) -> str: + normalized = value.strip() + if not normalized: + raise ValueError("text must not be blank") + return normalized + + @field_validator("failure_artifact_ids") + @classmethod + def validate_unique_artifacts(cls, values: tuple[str, ...]) -> tuple[str, ...]: + if len(set(values)) != len(values): + raise ValueError("group failure artifacts must be unique") + return values + + +class DeferredFailureInput(StrictModel): + failure_artifact_id: ArtifactIdentifier + reason: ExplanationText + + @field_validator("reason") + @classmethod + def normalize_reason(cls, value: str) -> str: + normalized = value.strip() + if not normalized: + raise ValueError("reason must not be blank") + return normalized + + +class RecordFailureCurationInput(StrictModel): + workspace_root: Path = Field(strict=False) + source_artifact_ids: SourceArtifactIds + groups: tuple[FailureGroupInput, ...] = Field(max_length=25) + deferred: tuple[DeferredFailureInput, ...] = Field(max_length=50) + + @field_validator("workspace_root") + @classmethod + def validate_workspace_root(cls, value: Path) -> Path: + if not value.is_absolute(): + raise ValueError("workspace_root must be absolute") + return value + + @field_validator("source_artifact_ids") + @classmethod + def validate_unique_sources(cls, values: tuple[str, ...]) -> tuple[str, ...]: + if len(set(values)) != len(values): + raise ValueError("source failure artifacts must be unique") + return values + + @model_validator(mode="after") + def validate_partition(self) -> RecordFailureCurationInput: + _require_unique_pattern_keys(self.groups) + _require_complete_partition( + self.source_artifact_ids, + _assigned_artifact_ids(self.groups, self.deferred), + ) + return self + + +@dataclass(frozen=True, slots=True) +class FailureSource: + artifact_id: str + artifact_digest: Sha256Digest + trace_id: TraceId + task_id: TaskId + outcome_score_id: ScoreId + evidence_status: FailureEvidenceStatus + issue_type: FailureType | None + critical_observation_id: ObservationId | None + + +@dataclass(frozen=True, slots=True) +class FailureGroupMember: + artifact_id: str + artifact_digest: Sha256Digest + trace_id: TraceId + task_id: TaskId + outcome_score_id: ScoreId + critical_observation_id: ObservationId + + +@dataclass(frozen=True, slots=True) +class FailureGroup: + id: str + pattern_key: str + title: str + mechanism: str + prevention: str + target_component: ComponentKind + issue_type: FailureType + members: tuple[FailureGroupMember, ...] + + +@dataclass(frozen=True, slots=True) +class DeferredFailure: + source: FailureSource + reason: str + + +@dataclass(frozen=True, slots=True) +class FailureCuration: + id: str + source_artifact_ids: tuple[str, ...] + groups: tuple[FailureGroup, ...] + deferred: tuple[DeferredFailure, ...] + + +@dataclass(frozen=True, slots=True) +class FailureCurationReceipt: + curation_id: str + relative_path: Path + + +class FailureCurationErrorCode(StrEnum): + INVALID_WORKSPACE = "invalid_workspace" + SOURCE_NOT_FOUND = "source_not_found" + SOURCE_INVALID = "source_invalid" + UNSUPPORTED_SOURCE = "unsupported_source" + MIXED_ISSUE_TYPES = "mixed_issue_types" + INSUFFICIENT_RECURRENCE = "insufficient_recurrence" + ARTIFACT_CONFLICT = "artifact_conflict" + ARTIFACT_TOO_LARGE = "artifact_too_large" + WRITE_FAILED = "write_failed" + + +class FailureCurationFailure(Exception): + __slots__ = ("code", "subject") + + def __init__(self, code: FailureCurationErrorCode, subject: str) -> None: + self.code = code + self.subject = subject + super().__init__(f"{code.value}: {subject}") + + +class FailureCurationStatus(StrEnum): + SUCCESS = "success" + + +class FailureCurationObservation(StrictModel): + status: FailureCurationStatus + summary: str = Field(min_length=1, max_length=256) + next_actions: tuple[str, ...] = Field(max_length=2) + artifacts: tuple[str, ...] = Field(min_length=2, max_length=2) + curation_id: ArtifactIdentifier + relative_path: Path + source_failure_count: int = Field(strict=True, ge=1, le=50) + group_count: int = Field(strict=True, ge=0, le=25) + deferred_count: int = Field(strict=True, ge=0, le=50) + + +class FailureGroupMemberArtifact(StrictModel): + artifact_id: ArtifactIdentifier + artifact_digest: DigestValue + trace_id: str = Field(min_length=1, max_length=256) + task_id: str = Field(min_length=1, max_length=256) + outcome_score_id: str = Field(min_length=1, max_length=256) + critical_observation_id: str = Field(min_length=1, max_length=256) + + @classmethod + def from_member(cls, member: FailureGroupMember) -> FailureGroupMemberArtifact: + return cls( + artifact_id=member.artifact_id, + artifact_digest=member.artifact_digest.value, + trace_id=member.trace_id.value, + task_id=member.task_id.value, + outcome_score_id=member.outcome_score_id.value, + critical_observation_id=member.critical_observation_id.value, + ) + + +class FailureGroupArtifact(StrictModel): + group_id: ArtifactIdentifier + pattern_key: PatternKey + title: TitleText + mechanism: ExplanationText + prevention: ExplanationText + target_component: ComponentKind + issue_type: FailureType + members: tuple[FailureGroupMemberArtifact, ...] = Field(min_length=2, max_length=20) + + @classmethod + def from_group(cls, group: FailureGroup) -> FailureGroupArtifact: + return cls( + group_id=group.id, + pattern_key=group.pattern_key, + title=group.title, + mechanism=group.mechanism, + prevention=group.prevention, + target_component=group.target_component, + issue_type=group.issue_type, + members=tuple(FailureGroupMemberArtifact.from_member(item) for item in group.members), + ) + + +class FailureSourceArtifact(StrictModel): + artifact_id: ArtifactIdentifier + artifact_digest: DigestValue + trace_id: str = Field(min_length=1, max_length=256) + task_id: str = Field(min_length=1, max_length=256) + outcome_score_id: str = Field(min_length=1, max_length=256) + evidence_status: FailureEvidenceStatus + issue_type: FailureType | None + critical_observation_id: str | None = Field(default=None, min_length=1, max_length=256) + + @classmethod + def from_source(cls, source: FailureSource) -> FailureSourceArtifact: + observation_id = source.critical_observation_id + return cls( + artifact_id=source.artifact_id, + artifact_digest=source.artifact_digest.value, + trace_id=source.trace_id.value, + task_id=source.task_id.value, + outcome_score_id=source.outcome_score_id.value, + evidence_status=source.evidence_status, + issue_type=source.issue_type, + critical_observation_id=None if observation_id is None else observation_id.value, + ) + + +class DeferredFailureArtifact(StrictModel): + source: FailureSourceArtifact + reason: ExplanationText + + @classmethod + def from_deferred(cls, deferred: DeferredFailure) -> DeferredFailureArtifact: + return cls( + source=FailureSourceArtifact.from_source(deferred.source), + reason=deferred.reason, + ) + + +class FailureCurationArtifact(StrictModel): + schema_version: Literal[1] = 1 + curation_id: ArtifactIdentifier + source_artifact_ids: SourceArtifactIds + groups: tuple[FailureGroupArtifact, ...] = Field(max_length=25) + deferred: tuple[DeferredFailureArtifact, ...] = Field(max_length=50) + + @classmethod + def from_curation(cls, curation: FailureCuration) -> FailureCurationArtifact: + return cls( + curation_id=curation.id, + source_artifact_ids=curation.source_artifact_ids, + groups=tuple(FailureGroupArtifact.from_group(group) for group in curation.groups), + deferred=tuple( + DeferredFailureArtifact.from_deferred(item) for item in curation.deferred + ), + ) + + +class FailureCurationGateway(Protocol): + def load(self, root: Path, artifact_ids: tuple[str, ...]) -> tuple[FailureSource, ...]: ... + + def store(self, root: Path, curation: FailureCuration) -> FailureCurationReceipt: ... + + +@dataclass(frozen=True, slots=True) +class FailureCurationService: + gateway: FailureCurationGateway + + def record(self, request: RecordFailureCurationInput) -> FailureCurationObservation: + curation = self.build(request) + receipt = self.gateway.store(request.workspace_root, curation) + return FailureCurationObservation( + status=FailureCurationStatus.SUCCESS, + summary=_summary(curation), + next_actions=_next_actions(curation), + artifacts=(str(receipt.relative_path), receipt.curation_id), + curation_id=receipt.curation_id, + relative_path=receipt.relative_path, + source_failure_count=len(curation.source_artifact_ids), + group_count=len(curation.groups), + deferred_count=len(curation.deferred), + ) + + def build(self, request: RecordFailureCurationInput) -> FailureCuration: + sources = self.gateway.load(request.workspace_root, request.source_artifact_ids) + _validate_loaded_sources(request.source_artifact_ids, sources) + groups = tuple(sorted((_group(item, sources) for item in request.groups), key=_group_key)) + deferred = tuple( + sorted( + ( + DeferredFailure( + _find_source(item.failure_artifact_id, sources), + item.reason, + ) + for item in request.deferred + ), + key=_deferred_key, + ) + ) + source_ids = tuple(sorted(request.source_artifact_ids)) + return FailureCuration( + id=_curation_id(source_ids, groups, deferred), + source_artifact_ids=source_ids, + groups=groups, + deferred=deferred, + ) + + +def _validate_loaded_sources( + requested: tuple[str, ...], + sources: tuple[FailureSource, ...], +) -> None: + loaded = tuple(source.artifact_id for source in sources) + if len(set(loaded)) != len(loaded): + raise FailureCurationFailure(FailureCurationErrorCode.SOURCE_INVALID, "source_artifacts") + _reject_unrequested_sources(requested, loaded) + _reject_missing_sources(requested, loaded) + + +def _group(item: FailureGroupInput, sources: tuple[FailureSource, ...]) -> FailureGroup: + selected = tuple( + _find_source(artifact_id, sources) for artifact_id in item.failure_artifact_ids + ) + _require_supported_sources(selected) + issue_type = _shared_issue_type(selected, item.pattern_key) + _require_distinct_tasks(selected, item.pattern_key) + members = tuple(sorted((_member(source) for source in selected), key=_member_key)) + return FailureGroup( + id=_group_id(item, issue_type, members), + pattern_key=item.pattern_key, + title=item.title, + mechanism=item.mechanism, + prevention=item.prevention, + target_component=item.target_component, + issue_type=issue_type, + members=members, + ) + + +def _require_supported_sources(selected: tuple[FailureSource, ...]) -> None: + unsupported = next((source for source in selected if not _is_supported(source)), None) + if unsupported is not None: + raise FailureCurationFailure( + FailureCurationErrorCode.UNSUPPORTED_SOURCE, + unsupported.artifact_id, + ) + + +def _shared_issue_type( + selected: tuple[FailureSource, ...], + pattern_key: str, +) -> FailureType: + issue_type = selected[0].issue_type + if issue_type is None: + raise FailureCurationFailure( + FailureCurationErrorCode.UNSUPPORTED_SOURCE, + selected[0].artifact_id, + ) + if any(source.issue_type is not issue_type for source in selected[1:]): + raise FailureCurationFailure( + FailureCurationErrorCode.MIXED_ISSUE_TYPES, + pattern_key, + ) + return issue_type + + +def _require_distinct_tasks(selected: tuple[FailureSource, ...], pattern_key: str) -> None: + if len({source.task_id for source in selected}) < 2: + raise FailureCurationFailure( + FailureCurationErrorCode.INSUFFICIENT_RECURRENCE, + pattern_key, + ) + + +def _reject_unrequested_sources(requested: tuple[str, ...], loaded: tuple[str, ...]) -> None: + unexpected = next((artifact_id for artifact_id in loaded if artifact_id not in requested), None) + if unexpected is not None: + raise FailureCurationFailure(FailureCurationErrorCode.SOURCE_INVALID, unexpected) + + +def _reject_missing_sources(requested: tuple[str, ...], loaded: tuple[str, ...]) -> None: + missing = next((artifact_id for artifact_id in requested if artifact_id not in loaded), None) + if missing is not None: + raise FailureCurationFailure(FailureCurationErrorCode.SOURCE_NOT_FOUND, missing) + + +def _require_unique_pattern_keys(groups: tuple[FailureGroupInput, ...]) -> None: + pattern_keys = tuple(group.pattern_key for group in groups) + if len(set(pattern_keys)) != len(pattern_keys): + raise ValueError("failure pattern keys must be unique") + + +def _assigned_artifact_ids( + groups: tuple[FailureGroupInput, ...], + deferred: tuple[DeferredFailureInput, ...], +) -> tuple[str, ...]: + grouped = tuple(artifact_id for group in groups for artifact_id in group.failure_artifact_ids) + return (*grouped, *(item.failure_artifact_id for item in deferred)) + + +def _require_complete_partition(sources: tuple[str, ...], assigned: tuple[str, ...]) -> None: + if len(set(assigned)) != len(assigned): + raise ValueError("each source failure must be assigned once") + if set(assigned) != set(sources): + raise ValueError("curation must partition every source failure") + + +def _find_source(artifact_id: str, sources: tuple[FailureSource, ...]) -> FailureSource: + # ponytail: linear lookup is bounded at 50; add an index only if that limit grows. + source = next((item for item in sources if item.artifact_id == artifact_id), None) + if source is None: + raise FailureCurationFailure(FailureCurationErrorCode.SOURCE_NOT_FOUND, artifact_id) + return source + + +def _is_supported(source: FailureSource) -> bool: + return ( + source.evidence_status is FailureEvidenceStatus.SUPPORTED + and source.issue_type is not None + and source.critical_observation_id is not None + ) + + +def _member(source: FailureSource) -> FailureGroupMember: + observation_id = source.critical_observation_id + if observation_id is None: + raise FailureCurationFailure( + FailureCurationErrorCode.UNSUPPORTED_SOURCE, + source.artifact_id, + ) + return FailureGroupMember( + artifact_id=source.artifact_id, + artifact_digest=source.artifact_digest, + trace_id=source.trace_id, + task_id=source.task_id, + outcome_score_id=source.outcome_score_id, + critical_observation_id=observation_id, + ) + + +def _group_id( + item: FailureGroupInput, + issue_type: FailureType, + members: tuple[FailureGroupMember, ...], +) -> str: + identity = "\0".join( + ( + "ofw.failure-group", + item.pattern_key, + item.title, + item.mechanism, + item.prevention, + item.target_component.value, + issue_type.value, + *(f"{member.artifact_id}:{member.artifact_digest.value}" for member in members), + ) + ) + return str(uuid5(NAMESPACE_URL, identity)) + + +def _curation_id( + source_ids: tuple[str, ...], + groups: tuple[FailureGroup, ...], + deferred: tuple[DeferredFailure, ...], +) -> str: + identity = "\0".join( + ( + "ofw.failure-curation", + *source_ids, + *(group.id for group in groups), + *( + f"{item.source.artifact_id}:{item.source.artifact_digest.value}:{item.reason}" + for item in deferred + ), + ) + ) + return str(uuid5(NAMESPACE_URL, identity)) + + +def _member_key(member: FailureGroupMember) -> tuple[str, str, str]: + return member.task_id.value, member.trace_id.value, member.artifact_id + + +def _group_key(group: FailureGroup) -> tuple[str, str]: + return group.pattern_key, group.id + + +def _deferred_key(item: DeferredFailure) -> str: + return item.source.artifact_id + + +def _summary(curation: FailureCuration) -> str: + group_word = "group" if len(curation.groups) == 1 else "groups" + deferred_word = "failure" if len(curation.deferred) == 1 else "failures" + return ( + f"Stored {_quantity(len(curation.groups))} evidence-bound failure {group_word} " + f"and {_quantity(len(curation.deferred))} deferred {deferred_word}." + ) + + +def _next_actions(curation: FailureCuration) -> tuple[str, ...]: + if curation.groups: + return ("Use one recorded group to form the next harness hypothesis.",) + return ("Do not form a harness hypothesis until a repeated supported pattern exists.",) + + +def _quantity(value: int) -> str: + return "one" if value == 1 else str(value) diff --git a/src/ofw/evaluation/failure_workspace.py b/src/ofw/evaluation/failure_workspace.py index b2e0d3f..5fe87f8 100644 --- a/src/ofw/evaluation/failure_workspace.py +++ b/src/ofw/evaluation/failure_workspace.py @@ -2,6 +2,7 @@ from __future__ import annotations +import hashlib import os import re import stat @@ -14,20 +15,34 @@ from typing import Annotated, Literal, Never, Protocol from uuid import NAMESPACE_URL, uuid4, uuid5 -from pydantic import BaseModel, ConfigDict, Field, ValidationError, field_validator +from pydantic import BaseModel, ConfigDict, Field, ValidationError, field_validator, model_validator +from ofw.contracts import Sha256Digest from ofw.evaluation.failure import ( FailureDiagnosis, FailureDiagnosisError, FailureEvidenceStatus, FailureType, ) +from ofw.evaluation.failure_curation import ( + FailureCuration, + FailureCurationArtifact, + FailureCurationErrorCode, + FailureCurationFailure, + FailureCurationReceipt, + FailureSource, +) from ofw.evaluation.failure_patterns import ( FailureDiagnosisRecord, FailurePatternMiningError, FailurePatternMiningErrorCode, ) -from ofw.evaluation.outcome import OutcomeEvaluation, OutcomeEvaluationError, TaskId, VerifierId +from ofw.evaluation.outcome import ( + OutcomeEvaluation, + OutcomeEvaluationError, + TaskId, + VerifierId, +) from ofw.observability.langfuse.domain import ObservationId, ScoreId, TraceId from ofw.runtime import EvidenceReference, VerifierVerdict @@ -37,6 +52,7 @@ _ARTIFACT_LIMIT_BYTES = 64 * 1024 _WORKSPACE_DIRECTORY = ".workspace" _FAILURE_DIRECTORY = "failures" +_CURATION_DIRECTORY = "failure-curations" _IGNORE_CONTENT = "*\n" _WORKSPACE_MARKERS = ("PROGRAM.md", "experiment_config.yaml") _DIRECTORY_FLAGS = os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW @@ -164,6 +180,14 @@ class FailureArtifact(StrictModel): counterfactual_action: DiagnosisText | None inconclusive_reason: DiagnosisText | None + @model_validator(mode="after") + def validate_domain_contract(self) -> FailureArtifact: + try: + self.to_diagnosis() + except (FailureDiagnosisError, OutcomeEvaluationError) as error: + raise ValueError("invalid failure artifact") from error + return self + @classmethod def from_diagnosis( cls, @@ -230,7 +254,7 @@ class FailureArtifactReceipt: class _DirectoryChainIdentity: root: tuple[int, int] workspace: tuple[int, int] - failures: tuple[int, int] + artifacts: tuple[int, int] class FailureWorkspaceErrorCode(StrEnum): @@ -280,7 +304,15 @@ class FileFailureWorkspace: def store(self, root: Path, diagnosis: FailureDiagnosis) -> FailureArtifactReceipt: artifact_id = _artifact_id(diagnosis) try: - return self._store(root, diagnosis, artifact_id) + artifact = FailureArtifact.from_diagnosis(artifact_id, diagnosis) + content = (artifact.model_dump_json(indent=2) + "\n").encode("utf-8") + relative_path = _store_artifact( + root, + _FAILURE_DIRECTORY, + artifact_id, + content, + ) + return FailureArtifactReceipt(artifact_id, relative_path) except (OSError, RuntimeError, UnicodeError): raise FailureWorkspaceFailure( FailureWorkspaceErrorCode.WRITE_FAILED, @@ -304,46 +336,95 @@ def read( "failure_artifacts", ) from None - def _store( + def _read( self, root: Path, - diagnosis: FailureDiagnosis, - artifact_id: str, - ) -> FailureArtifactReceipt: + artifact_ids: tuple[str, ...], + ) -> tuple[FailureDiagnosisRecord, ...]: prepared_root = _prepared_root(root) - workspace, failures = _workspace_paths(prepared_root) - artifact = FailureArtifact.from_diagnosis(artifact_id, diagnosis) - content = (artifact.model_dump_json(indent=2) + "\n").encode("utf-8") - _validate_artifact_size(content) - directory_identity = _prepare_workspace_directories( + workspace, failures = _workspace_artifact_paths(prepared_root, _FAILURE_DIRECTORY) + directory_identity = _existing_workspace_directories( prepared_root, workspace, failures, ) - path = failures / f"{artifact_id}.json" - receipt = FailureArtifactReceipt(artifact_id, path.relative_to(prepared_root)) - with _failure_directory_handle( + with _artifact_directory_handle( prepared_root, + _FAILURE_DIRECTORY, directory_identity, ) as directory: - _publish_or_validate(directory, path.name, content, artifact_id) - return receipt + return tuple(_read_artifact(directory, artifact_id) for artifact_id in artifact_ids) - def _read( +class FileFailureCurationWorkspace: + """Read compact diagnoses and store one bounded cross-failure curation.""" + + def load(self, root: Path, artifact_ids: tuple[str, ...]) -> tuple[FailureSource, ...]: + try: + return self._load(root, artifact_ids) + except FailureCurationFailure: + raise + except FailureWorkspaceFailure as error: + code = ( + FailureCurationErrorCode.SOURCE_INVALID + if error.code is FailureWorkspaceErrorCode.ARTIFACT_TOO_LARGE + else FailureCurationErrorCode.INVALID_WORKSPACE + ) + raise FailureCurationFailure(code, error.subject) from None + except (OSError, RuntimeError, UnicodeError): + raise FailureCurationFailure( + FailureCurationErrorCode.SOURCE_INVALID, + "failure_artifacts", + ) from None + + def store(self, root: Path, curation: FailureCuration) -> FailureCurationReceipt: + try: + artifact = FailureCurationArtifact.from_curation(curation) + content = (artifact.model_dump_json(indent=2) + "\n").encode("utf-8") + relative_path = _store_artifact( + root, + _CURATION_DIRECTORY, + curation.id, + content, + ) + return FailureCurationReceipt(curation.id, relative_path) + except FailureWorkspaceFailure as error: + raise FailureCurationFailure(_curation_workspace_error(error), error.subject) from None + except (OSError, RuntimeError, UnicodeError): + raise FailureCurationFailure( + FailureCurationErrorCode.WRITE_FAILED, + curation.id, + ) from None + + def _load( self, root: Path, artifact_ids: tuple[str, ...], - ) -> tuple[FailureDiagnosisRecord, ...]: + ) -> tuple[FailureSource, ...]: prepared_root = _prepared_root(root) - workspace, failures = _workspace_paths(prepared_root) - directory_identity = _existing_workspace_directories( - prepared_root, - workspace, - failures, - ) - with _failure_directory_handle(prepared_root, directory_identity) as directory: - return tuple(_read_artifact(directory, artifact_id) for artifact_id in artifact_ids) - + workspace, failures = _workspace_artifact_paths(prepared_root, _FAILURE_DIRECTORY) + try: + identity = _existing_directory_identity(prepared_root, workspace, failures) + with _artifact_directory_handle( + prepared_root, + _FAILURE_DIRECTORY, + identity, + ) as directory: + return tuple( + _read_failure_source(directory, artifact_id) for artifact_id in artifact_ids + ) + except FileNotFoundError: + missing = next( + ( + artifact_id + for artifact_id in artifact_ids + if not (failures / f"{artifact_id}.json").exists() + ), + artifact_ids[0], + ) + raise FailureCurationFailure( + FailureCurationErrorCode.SOURCE_NOT_FOUND, + missing, + ) from None def _observation_id(value: str | None) -> ObservationId | None: return None if value is None else ObservationId(value) @@ -360,6 +441,48 @@ def _required_score(outcome: OutcomeEvaluation) -> float: return score +def _read_failure_source(directory: int, artifact_id: str) -> FailureSource: + try: + content = _read_existing(directory, f"{artifact_id}.json", artifact_id) + artifact = FailureArtifact.model_validate_json(content) + except FileNotFoundError: + raise FailureCurationFailure( + FailureCurationErrorCode.SOURCE_NOT_FOUND, + artifact_id, + ) from None + except (FailureWorkspaceFailure, ValidationError): + raise FailureCurationFailure( + FailureCurationErrorCode.SOURCE_INVALID, + artifact_id, + ) from None + if artifact.artifact_id != artifact_id: + raise FailureCurationFailure( + FailureCurationErrorCode.SOURCE_INVALID, + artifact_id, + ) + diagnosis = artifact.to_diagnosis() + return FailureSource( + artifact_id=artifact.artifact_id, + artifact_digest=Sha256Digest(f"sha256:{hashlib.sha256(content).hexdigest()}"), + trace_id=diagnosis.outcome.trace_id, + task_id=diagnosis.outcome.task_id, + outcome_score_id=diagnosis.outcome_score_id, + evidence_status=diagnosis.evidence_status, + issue_type=diagnosis.issue_type, + critical_observation_id=diagnosis.critical_observation_id, + ) + + +def _curation_workspace_error(error: FailureWorkspaceFailure) -> FailureCurationErrorCode: + if error.code is FailureWorkspaceErrorCode.INVALID_WORKSPACE: + return FailureCurationErrorCode.INVALID_WORKSPACE + if error.code is FailureWorkspaceErrorCode.ARTIFACT_CONFLICT: + return FailureCurationErrorCode.ARTIFACT_CONFLICT + if error.code is FailureWorkspaceErrorCode.ARTIFACT_TOO_LARGE: + return FailureCurationErrorCode.ARTIFACT_TOO_LARGE + return FailureCurationErrorCode.WRITE_FAILED + + def _artifact_id(diagnosis: FailureDiagnosis) -> str: identity = "\0".join( ( @@ -389,15 +512,26 @@ def _is_prepared_root(root: Path) -> bool: return all((root / name).is_file() for name in _WORKSPACE_MARKERS) -def _workspace_paths(root: Path) -> tuple[Path, Path]: +def _workspace_artifact_paths(root: Path, directory_name: str) -> tuple[Path, Path]: workspace = root / _WORKSPACE_DIRECTORY - failures = workspace / _FAILURE_DIRECTORY + failures = workspace / directory_name _require_contained(root, workspace.resolve(strict=False)) _require_contained(root, failures.resolve(strict=False)) _require_directory_if_present(workspace) return workspace, failures +def _store_artifact(root: Path, directory_name: str, artifact_id: str, content: bytes) -> Path: + prepared_root = _prepared_root(root) + workspace, artifacts = _workspace_artifact_paths(prepared_root, directory_name) + _validate_artifact_size(content) + identity = _prepare_workspace_directories(prepared_root, workspace, artifacts) + path = artifacts / f"{artifact_id}.json" + with _artifact_directory_handle(prepared_root, directory_name, identity) as directory: + _publish_or_validate(directory, path.name, content, artifact_id) + return path.relative_to(prepared_root) + + def _require_directory_if_present(path: Path) -> None: if path.exists() and not path.is_dir(): _invalid_workspace(_WORKSPACE_DIRECTORY) @@ -439,7 +573,21 @@ def _prepare_workspace_directories( return _DirectoryChainIdentity( root=_path_identity(root), workspace=_path_identity(workspace), - failures=_path_identity(failures), + artifacts=_path_identity(failures), + ) + + +def _existing_directory_identity( + root: Path, + workspace: Path, + artifacts: Path, +) -> _DirectoryChainIdentity: + _require_contained(root, workspace.resolve(strict=True)) + _require_contained(root, artifacts.resolve(strict=True)) + return _DirectoryChainIdentity( + root=_path_identity(root), + workspace=_path_identity(workspace), + artifacts=_path_identity(artifacts), ) @@ -453,11 +601,7 @@ def _existing_workspace_directories( FailurePatternMiningErrorCode.ARTIFACT_NOT_FOUND, _FAILURE_DIRECTORY, ) - return _DirectoryChainIdentity( - root=_path_identity(root), - workspace=_path_identity(workspace), - failures=_path_identity(failures), - ) + return _existing_directory_identity(root, workspace, failures) def _write_ignore_file(directory: int) -> None: @@ -613,28 +757,31 @@ def _child_directory_handle(parent: int, name: str) -> Iterator[int]: @contextmanager -def _failure_directory_handle( +def _artifact_directory_handle( root: Path, + directory_name: str, expected: _DirectoryChainIdentity, ) -> Iterator[int]: with ( _directory_handle(root) as root_directory, _child_directory_handle(root_directory, _WORKSPACE_DIRECTORY) as workspace, - _child_directory_handle(workspace, _FAILURE_DIRECTORY) as failures, + _child_directory_handle(workspace, directory_name) as artifacts, ): _require_directory_chain( root, root_directory, workspace, - failures, + directory_name, + artifacts, expected, ) - yield failures + yield artifacts _require_directory_chain( root, root_directory, workspace, - failures, + directory_name, + artifacts, expected, ) @@ -643,12 +790,13 @@ def _require_directory_chain( root: Path, root_directory: int, workspace: int, - failures: int, + directory_name: str, + artifacts: int, expected: _DirectoryChainIdentity, ) -> None: _require_directory_identity(root_directory, root, expected.root) _require_child_identity(root_directory, _WORKSPACE_DIRECTORY, workspace, expected.workspace) - _require_child_identity(workspace, _FAILURE_DIRECTORY, failures, expected.failures) + _require_child_identity(workspace, directory_name, artifacts, expected.artifacts) def _require_directory_identity( diff --git a/src/ofw/mcp.py b/src/ofw/mcp.py index 572caa9..0c20bd8 100644 --- a/src/ofw/mcp.py +++ b/src/ofw/mcp.py @@ -14,6 +14,11 @@ from mcp.types import ToolAnnotations from pydantic import BaseModel, Field +from ofw.evaluation.failure_curation import ( + FailureCurationObservation, + FailureCurationService, + RecordFailureCurationInput, +) from ofw.evaluation.failure_patterns import ( FailurePatternMiningObservation, FailurePatternMiningService, @@ -22,6 +27,7 @@ from ofw.evaluation.failure_workspace import ( FailureRecordObservation, FailureWorkspaceService, + FileFailureCurationWorkspace, FileFailureWorkspace, RecordFailureInput, ) @@ -73,8 +79,9 @@ name="openflywheel", instructions=( "Prepare isolated ITSM harness workspaces, read bounded Langfuse trace evidence, and " - "record authoritative outcomes plus compact failure diagnoses and exact patterns. " - "Never infer outcomes, mutate traces, or copy trace payloads into local storage." + "record authoritative outcomes, compact failure diagnoses, exact patterns, and " + "evidence-bound curations. Never infer outcomes, mutate traces, or copy trace payloads " + "into local storage." ), log_level="DEBUG", ) @@ -139,6 +146,10 @@ def _failure_pattern_service() -> FailurePatternMiningService: return FailurePatternMiningService(FileFailureWorkspace()) +def _curation_service() -> FailureCurationService: + return FailureCurationService(FileFailureCurationWorkspace()) + + def _program_template(name: str) -> str: content = files("ofw.preparation.templates").joinpath(name).read_bytes() if len(content) > _PROGRAM_TEMPLATE_LIMIT_BYTES: @@ -276,6 +287,14 @@ def mine_failure_patterns( return _failure_pattern_service().mine(request) +@server.tool(annotations=record_write, structured_output=True) +def record_failure_curation( + request: RecordFailureCurationInput, +) -> FailureCurationObservation: + """Store one bounded cross-failure curation under a prepared local workspace.""" + return _curation_service().record(request) + + def main() -> None: """Run the OpenFlywheel MCP server over stdio.""" server.run(transport="stdio") diff --git a/src/ofw/preparation/templates/itsm.md b/src/ofw/preparation/templates/itsm.md index 9174db3..b7f328b 100644 --- a/src/ofw/preparation/templates/itsm.md +++ b/src/ofw/preparation/templates/itsm.md @@ -49,6 +49,14 @@ After recording the diagnoses for one bounded run, use `$failure-pattern-miner` type and exact normalized root cause; results are not semantic clusters. Keep inconclusive diagnoses separate, and reread each supporting diagnosis before forming a shared hypothesis. +After every failed outcome has a retained artifact, use `$failure-curator` for one bounded +cross-task debugger pass. Use only the current run's retained `record_failure` receipt IDs. +Do not glob `.workspace/failures/`, which may contain earlier runs. Submit the complete set of at +most 50 IDs to `record_failure_curation`, assigning each exactly once to a repeated +evidence-bound group or a deferred entry. Retain the `.workspace/failure-curations/` artifact +before forming a harness hypothesis. If the run has more than 50 failure artifacts, report the +curation bound and stop before forming a harness hypothesis. + An intermediate tool error is evidence, not an outcome failure, when the agent recovered and the verifier passed. A technically clean trajectory is still a failure when the ITSM verifier shows that the required environment state was not achieved. diff --git a/tests/test_failure_curation.py b/tests/test_failure_curation.py new file mode 100644 index 0000000..7aa75d8 --- /dev/null +++ b/tests/test_failure_curation.py @@ -0,0 +1,372 @@ +"""Evidence-bound cross-failure curation tests.""" + +from __future__ import annotations + +from pathlib import Path + +import pytest +from pydantic import ValidationError + +from ofw.contracts import ComponentKind, Sha256Digest +from ofw.evaluation.failure import FailureEvidenceStatus, FailureType +from ofw.evaluation.failure_curation import ( + DeferredFailureInput, + FailureCuration, + FailureCurationErrorCode, + FailureCurationFailure, + FailureCurationObservation, + FailureCurationReceipt, + FailureCurationService, + FailureCurationStatus, + FailureGroupInput, + FailureSource, + RecordFailureCurationInput, +) +from ofw.evaluation.outcome import TaskId +from ofw.observability.langfuse.domain import ObservationId, ScoreId, TraceId + +_FIRST_ARTIFACT_ID = "00000000-0000-0000-0000-000000000001" +_SECOND_ARTIFACT_ID = "00000000-0000-0000-0000-000000000002" +_THIRD_ARTIFACT_ID = "00000000-0000-0000-0000-000000000003" + + +class _FakeCurationGateway: + def __init__(self, sources: tuple[FailureSource, ...]) -> None: + self.sources = sources + self.loaded: list[tuple[Path, tuple[str, ...]]] = [] + self.stored: list[tuple[Path, FailureCuration]] = [] + + def load(self, root: Path, artifact_ids: tuple[str, ...]) -> tuple[FailureSource, ...]: + self.loaded.append((root, artifact_ids)) + return self.sources + + def store(self, root: Path, curation: FailureCuration) -> FailureCurationReceipt: + self.stored.append((root, curation)) + return FailureCurationReceipt( + curation_id=curation.id, + relative_path=Path(f".workspace/failure-curations/{curation.id}.json"), + ) + + +def _source( + artifact_id: str, + task_id: str, + *, + trace_id: str | None = None, + issue_type: FailureType | None = FailureType.CONTROL_FLOW_FAILURE, + evidence_status: FailureEvidenceStatus = FailureEvidenceStatus.SUPPORTED, + digest_character: str | None = None, +) -> FailureSource: + suffix = artifact_id[-1] + return FailureSource( + artifact_id=artifact_id, + artifact_digest=Sha256Digest(f"sha256:{(digest_character or suffix) * 64}"), + trace_id=TraceId(trace_id or f"trace-{suffix}"), + task_id=TaskId(task_id), + outcome_score_id=ScoreId(f"score-{suffix}"), + evidence_status=evidence_status, + issue_type=issue_type, + critical_observation_id=( + ObservationId(f"observation-{suffix}") + if evidence_status is FailureEvidenceStatus.SUPPORTED + else None + ), + ) + + +def _group(*artifact_ids: str) -> FailureGroupInput: + return FailureGroupInput( + pattern_key="finalizes-before-verification", + title="Finalizes before verifying state", + mechanism="The control loop treats a successful mutation as task completion.", + prevention="Require a state read after mutation and before finalization.", + target_component=ComponentKind.PROMPT, + failure_artifact_ids=artifact_ids, + ) + + +def _request(root: Path) -> RecordFailureCurationInput: + return RecordFailureCurationInput( + workspace_root=root, + source_artifact_ids=( + _THIRD_ARTIFACT_ID, + _SECOND_ARTIFACT_ID, + _FIRST_ARTIFACT_ID, + ), + groups=(_group(_SECOND_ARTIFACT_ID, _FIRST_ARTIFACT_ID),), + deferred=( + DeferredFailureInput( + failure_artifact_id=_THIRD_ARTIFACT_ID, + reason="No second task supports this mechanism yet.", + ), + ), + ) + + +def test_curates_repeated_failures_into_one_deterministic_actionable_group( + tmp_path: Path, +) -> None: + sources = ( + _source(_THIRD_ARTIFACT_ID, "task-3"), + _source(_SECOND_ARTIFACT_ID, "task-2"), + _source(_FIRST_ARTIFACT_ID, "task-1"), + ) + gateway = _FakeCurationGateway(sources) + service = FailureCurationService(gateway) + request = _request(tmp_path) + + observation = service.record(request) + + curation = gateway.stored[0][1] + group = curation.groups[0] + expected = FailureCurationObservation( + status=FailureCurationStatus.SUCCESS, + summary="Stored one evidence-bound failure group and one deferred failure.", + next_actions=("Use one recorded group to form the next harness hypothesis.",), + artifacts=(str(observation.relative_path), observation.curation_id), + curation_id=observation.curation_id, + relative_path=observation.relative_path, + source_failure_count=3, + group_count=1, + deferred_count=1, + ) + assert ( + gateway.loaded, + observation, + tuple(member.task_id.value for member in group.members), + group.issue_type, + group.target_component, + curation.deferred[0].source.artifact_id, + ) == ( + [(tmp_path, request.source_artifact_ids)], + expected, + ("task-1", "task-2"), + FailureType.CONTROL_FLOW_FAILURE, + ComponentKind.PROMPT, + _THIRD_ARTIFACT_ID, + ) + + +def test_curation_identity_and_order_do_not_depend_on_request_order(tmp_path: Path) -> None: + sources = ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source(_SECOND_ARTIFACT_ID, "task-2"), + _source(_THIRD_ARTIFACT_ID, "task-3"), + ) + service = FailureCurationService(_FakeCurationGateway(sources)) + first = service.build(_request(tmp_path)) + reordered = RecordFailureCurationInput( + workspace_root=tmp_path, + source_artifact_ids=( + _FIRST_ARTIFACT_ID, + _SECOND_ARTIFACT_ID, + _THIRD_ARTIFACT_ID, + ), + groups=(_group(_FIRST_ARTIFACT_ID, _SECOND_ARTIFACT_ID),), + deferred=( + DeferredFailureInput( + failure_artifact_id=_THIRD_ARTIFACT_ID, + reason="No second task supports this mechanism yet.", + ), + ), + ) + + assert service.build(reordered) == first + + +def test_curation_identity_is_bound_to_source_artifact_content(tmp_path: Path) -> None: + original_sources = ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source(_SECOND_ARTIFACT_ID, "task-2"), + _source(_THIRD_ARTIFACT_ID, "task-3"), + ) + changed_sources = ( + _source(_FIRST_ARTIFACT_ID, "task-1", digest_character="a"), + original_sources[1], + original_sources[2], + ) + request = _request(tmp_path) + + original = FailureCurationService(_FakeCurationGateway(original_sources)).build(request) + changed = FailureCurationService(_FakeCurationGateway(changed_sources)).build(request) + + assert (changed.id, changed.groups[0].id) != (original.id, original.groups[0].id) + + +def test_all_deferred_curation_blocks_a_harness_hypothesis(tmp_path: Path) -> None: + source = _source( + _FIRST_ARTIFACT_ID, + "task-1", + issue_type=None, + evidence_status=FailureEvidenceStatus.INCONCLUSIVE, + ) + service = FailureCurationService(_FakeCurationGateway((source,))) + request = RecordFailureCurationInput( + workspace_root=tmp_path, + source_artifact_ids=(_FIRST_ARTIFACT_ID,), + groups=(), + deferred=( + DeferredFailureInput( + failure_artifact_id=_FIRST_ARTIFACT_ID, + reason="The trace lacks the terminal state observation.", + ), + ), + ) + + observation = service.record(request) + + assert observation.next_actions == ( + "Do not form a harness hypothesis until a repeated supported pattern exists.", + ) + + +@pytest.mark.parametrize( + ("sources", "code"), + [ + ( + ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source( + _SECOND_ARTIFACT_ID, + "task-2", + issue_type=FailureType.TOOL_INTERACTION_FAILURE, + ), + _source(_THIRD_ARTIFACT_ID, "task-3"), + ), + FailureCurationErrorCode.MIXED_ISSUE_TYPES, + ), + ( + ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source( + _SECOND_ARTIFACT_ID, + "task-2", + issue_type=None, + evidence_status=FailureEvidenceStatus.INCONCLUSIVE, + ), + _source(_THIRD_ARTIFACT_ID, "task-3"), + ), + FailureCurationErrorCode.UNSUPPORTED_SOURCE, + ), + ( + ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source(_SECOND_ARTIFACT_ID, "task-1", trace_id="trace-other"), + _source(_THIRD_ARTIFACT_ID, "task-3"), + ), + FailureCurationErrorCode.INSUFFICIENT_RECURRENCE, + ), + ], +) +def test_group_members_must_be_supported_same_type_failures_from_distinct_tasks( + tmp_path: Path, + sources: tuple[FailureSource, ...], + code: FailureCurationErrorCode, +) -> None: + service = FailureCurationService(_FakeCurationGateway(sources)) + + with pytest.raises(FailureCurationFailure) as raised: + service.build(_request(tmp_path)) + + assert raised.value.code is code + + +def test_gateway_must_return_every_requested_source(tmp_path: Path) -> None: + gateway = _FakeCurationGateway( + ( + _source(_FIRST_ARTIFACT_ID, "task-1"), + _source(_SECOND_ARTIFACT_ID, "task-2"), + ) + ) + + with pytest.raises(FailureCurationFailure) as raised: + FailureCurationService(gateway).build(_request(tmp_path)) + + assert raised.value.code is FailureCurationErrorCode.SOURCE_NOT_FOUND + assert raised.value.subject == _THIRD_ARTIFACT_ID + + +def test_curation_input_rejects_relative_workspace() -> None: + with pytest.raises(ValidationError): + RecordFailureCurationInput( + workspace_root=Path("relative"), + source_artifact_ids=(_FIRST_ARTIFACT_ID,), + groups=(), + deferred=( + DeferredFailureInput( + failure_artifact_id=_FIRST_ARTIFACT_ID, + reason="No repeated mechanism.", + ), + ), + ) + + +def test_curation_input_rejects_an_empty_source_set(tmp_path: Path) -> None: + with pytest.raises(ValidationError): + RecordFailureCurationInput( + workspace_root=tmp_path, + source_artifact_ids=(), + groups=(), + deferred=(), + ) + + +def test_curation_input_rejects_an_artifact_path(tmp_path: Path) -> None: + escaped = f"../{_FIRST_ARTIFACT_ID}" + + with pytest.raises(ValidationError): + RecordFailureCurationInput( + workspace_root=tmp_path, + source_artifact_ids=(escaped,), + groups=(), + deferred=( + DeferredFailureInput( + failure_artifact_id=escaped, + reason="No repeated mechanism.", + ), + ), + ) + + +@pytest.mark.parametrize( + ("sources", "groups", "deferred"), + [ + ( + (_FIRST_ARTIFACT_ID, _FIRST_ARTIFACT_ID), + (), + (), + ), + ( + (_FIRST_ARTIFACT_ID, _SECOND_ARTIFACT_ID), + (_group(_FIRST_ARTIFACT_ID, _SECOND_ARTIFACT_ID),), + (DeferredFailureInput(failure_artifact_id=_FIRST_ARTIFACT_ID, reason="duplicate"),), + ), + ( + (_FIRST_ARTIFACT_ID, _SECOND_ARTIFACT_ID, _THIRD_ARTIFACT_ID), + (_group(_FIRST_ARTIFACT_ID, _SECOND_ARTIFACT_ID),), + (), + ), + ], +) +def test_curation_input_requires_a_unique_complete_partition( + tmp_path: Path, + sources: tuple[str, ...], + groups: tuple[FailureGroupInput, ...], + deferred: tuple[DeferredFailureInput, ...], +) -> None: + with pytest.raises(ValidationError): + RecordFailureCurationInput( + workspace_root=tmp_path, + source_artifact_ids=sources, + groups=groups, + deferred=deferred, + ) + + +def test_curation_input_forbids_unknown_fields(tmp_path: Path) -> None: + request = _request(tmp_path) + + with pytest.raises(ValidationError): + RecordFailureCurationInput.model_validate_json( + request.model_dump_json().removesuffix("}") + ',"unexpected":"field"}' + ) diff --git a/tests/test_failure_workspace.py b/tests/test_failure_workspace.py index 23c2241..cc2d700 100644 --- a/tests/test_failure_workspace.py +++ b/tests/test_failure_workspace.py @@ -9,12 +9,26 @@ from datetime import UTC, datetime from pathlib import Path from threading import Barrier +from uuid import UUID import pytest from pydantic import ValidationError import ofw.evaluation.failure_workspace as failure_workspace_module +from ofw.contracts import ComponentKind, Sha256Digest from ofw.evaluation.failure import FailureEvidenceStatus, FailureType +from ofw.evaluation.failure_curation import ( + DeferredFailureInput, + FailureCuration, + FailureCurationArtifact, + FailureCurationErrorCode, + FailureCurationFailure, + FailureCurationService, + FailureGroup, + FailureGroupInput, + FailureGroupMember, + RecordFailureCurationInput, +) from ofw.evaluation.failure_workspace import ( FailedOutcomeInput, FailureArtifact, @@ -23,13 +37,17 @@ FailureWorkspaceErrorCode, FailureWorkspaceFailure, FailureWorkspaceService, + FileFailureCurationWorkspace, FileFailureWorkspace, RecordFailureInput, ) +from ofw.evaluation.outcome import TaskId +from ofw.observability.langfuse.domain import ObservationId, ScoreId, TraceId _EVALUATED_AT = datetime(2026, 8, 28, 6, 0, tzinfo=UTC) _ARTIFACT_LIMIT_BYTES = 64 * 1024 RecordResult = FailureRecordObservation | FailureWorkspaceErrorCode +RecordedFailures = tuple[tuple[str, str, str], tuple[Path, Path, Path]] def _git(root: Path, *arguments: str) -> str: @@ -57,31 +75,116 @@ def _prepared_workspace(tmp_path: Path) -> Path: def _request( root: Path, root_cause: str = "The agent finalized before reading state.", + *, + trace_id: str = "trace-1", + task_id: str = "task-1", + outcome_score_id: str = "outcome-score-1", + critical_observation_id: str = "observation-7", ) -> RecordFailureInput: - critical = "observation-7" return RecordFailureInput( workspace_root=root, outcome=FailedOutcomeInput( - trace_id="trace-1", - task_id="task-1", + trace_id=trace_id, + task_id=task_id, verifier_id="itsm-bench@v1", evaluated_at=_EVALUATED_AT, score=0.0, evidence=("harbor://trial-1/verifier/result",), - outcome_score_id="outcome-score-1", + outcome_score_id=outcome_score_id, ), evidence_status=FailureEvidenceStatus.SUPPORTED, issue_type=FailureType.CONTROL_FLOW_FAILURE, expected_outcome="Incident INC-123 is closed.", actual_outcome="Incident INC-123 remains open.", - critical_observation_id=critical, - evidence_observation_ids=(critical, "observation-9"), + critical_observation_id=critical_observation_id, + evidence_observation_ids=(critical_observation_id, "observation-9"), root_cause=root_cause, counterfactual_action="Read the incident state before finalizing.", inconclusive_reason=None, ) +def _curation_request(root: Path, artifact_ids: tuple[str, str, str]) -> RecordFailureCurationInput: + return RecordFailureCurationInput( + workspace_root=root, + source_artifact_ids=artifact_ids, + groups=( + FailureGroupInput( + pattern_key="finalizes-before-verification", + title="Finalizes before verifying state", + mechanism="The control loop treats a successful mutation as completion.", + prevention="Require a state read after mutation and before finalization.", + target_component=ComponentKind.PROMPT, + failure_artifact_ids=(artifact_ids[0], artifact_ids[1]), + ), + ), + deferred=( + DeferredFailureInput( + failure_artifact_id=artifact_ids[2], + reason="No second task supports this mechanism yet.", + ), + ), + ) + + +def _record_failures(root: Path) -> RecordedFailures: + service = FailureWorkspaceService(FileFailureWorkspace()) + receipts = tuple( + service.record( + _request( + root, + trace_id=f"trace-{index}", + task_id=f"task-{index}", + outcome_score_id=f"score-{index}", + critical_observation_id=f"observation-{index}", + ) + ) + for index in range(1, 4) + ) + return ( + (receipts[0].artifact_id, receipts[1].artifact_id, receipts[2].artifact_id), + ( + root / receipts[0].relative_path, + root / receipts[1].relative_path, + root / receipts[2].relative_path, + ), + ) + + +def _group_member(index: int) -> FailureGroupMember: + return FailureGroupMember( + artifact_id=str(UUID(int=index)), + artifact_digest=Sha256Digest(f"sha256:{index:064x}"), + trace_id=TraceId(f"trace-{index}"), + task_id=TaskId(f"task-{index}"), + outcome_score_id=ScoreId(f"score-{index}"), + critical_observation_id=ObservationId(f"observation-{index}"), + ) + + +def _oversized_curation() -> FailureCuration: + members = tuple(_group_member(index) for index in range(1, 51)) + groups = tuple( + FailureGroup( + id=str(UUID(int=100 + index)), + pattern_key=f"pattern-{index}", + title="T" * 160, + mechanism="M" * 1000, + prevention="P" * 1000, + target_component=ComponentKind.PROMPT, + issue_type=FailureType.CONTROL_FLOW_FAILURE, + members=(members[index * 2], members[index * 2 + 1]), + ) + for index in range(25) + ) + return FailureCuration( + id=str(UUID(int=999)), + source_artifact_ids=tuple(member.artifact_id for member in members), + groups=groups, + deferred=(), + ) + + def _expected_artifact(artifact_id: str) -> FailureArtifact: return FailureArtifact( artifact_id=artifact_id, @@ -368,6 +471,73 @@ def prepare_then_swap( assert not tuple((root / ".workspace").rglob("*.json")) +def test_curates_recorded_failures_without_copying_trace_content(tmp_path: Path) -> None: + root = _prepared_workspace(tmp_path) + artifact_ids, _ = _record_failures(root) + service = FailureCurationService(FileFailureCurationWorkspace()) + + observation = service.record(_curation_request(root, artifact_ids)) + + artifact_path = root / observation.relative_path + artifact = FailureCurationArtifact.model_validate_json(artifact_path.read_text()) + assert ( + tuple(member.task_id for member in artifact.groups[0].members), + artifact.deferred[0].source.artifact_id, + artifact_path.read_text().find("The agent finalized before reading state."), + service.record(_curation_request(root, artifact_ids)), + len(tuple(artifact_path.parent.glob("*.json"))), + _git(root, "status", "--short"), + ) == (("task-1", "task-2"), artifact_ids[2], -1, observation, 1, "") + + +def test_curation_rejects_a_missing_or_tampered_failure_artifact(tmp_path: Path) -> None: + root = _prepared_workspace(tmp_path) + artifact_ids, artifact_paths = _record_failures(root) + missing_path = artifact_paths[0] + missing_path.unlink() + service = FailureCurationService(FileFailureCurationWorkspace()) + + with pytest.raises(FailureCurationFailure) as missing: + service.record(_curation_request(root, artifact_ids)) + + assert missing.value.code is FailureCurationErrorCode.SOURCE_NOT_FOUND + missing_path.write_text("{}", encoding="utf-8") + + with pytest.raises(FailureCurationFailure) as invalid: + service.record(_curation_request(root, artifact_ids)) + + assert invalid.value.code is FailureCurationErrorCode.SOURCE_INVALID + assert not (root / ".workspace/failure-curations").exists() + + +def test_curation_does_not_follow_a_failure_artifact_symlink(tmp_path: Path) -> None: + root = _prepared_workspace(tmp_path) + artifact_ids, artifact_paths = _record_failures(root) + source_path = artifact_paths[0] + outside = tmp_path / "outside.json" + source_path.replace(outside) + source_path.symlink_to(outside) + + with pytest.raises(FailureCurationFailure) as raised: + FailureCurationService(FileFailureCurationWorkspace()).record( + _curation_request(root, artifact_ids) + ) + + assert raised.value.code is FailureCurationErrorCode.SOURCE_INVALID + assert outside.read_bytes() + assert not (root / ".workspace/failure-curations").exists() + + +def test_oversized_curation_fails_before_workspace_mutation(tmp_path: Path) -> None: + root = _prepared_workspace(tmp_path) + + with pytest.raises(FailureCurationFailure) as raised: + FileFailureCurationWorkspace().store(root, _oversized_curation()) + + assert raised.value.code is FailureCurationErrorCode.ARTIFACT_TOO_LARGE + assert not (root / ".workspace/failure-curations").exists() + + def test_record_input_rejects_relative_workspace_and_extra_fields(tmp_path: Path) -> None: request = _request(_prepared_workspace(tmp_path)) diff --git a/tests/test_openflywheel_mcp.py b/tests/test_openflywheel_mcp.py index f2f2dd1..0492ec4 100644 --- a/tests/test_openflywheel_mcp.py +++ b/tests/test_openflywheel_mcp.py @@ -12,7 +12,15 @@ from mcp.server.fastmcp import FastMCP from mcp.types import Tool +from ofw.contracts import ComponentKind from ofw.evaluation.failure import FailureEvidenceStatus, FailureType +from ofw.evaluation.failure_curation import ( + FailureCurationObservation, + FailureCurationService, + FailureCurationStatus, + FailureGroupInput, + RecordFailureCurationInput, +) from ofw.evaluation.failure_patterns import ( FailurePatternMiningObservation, FailurePatternMiningStatus, @@ -59,6 +67,8 @@ from ofw.runtime import EvidenceReference, VerifierVerdict _FAILURE_ARTIFACT_ID = "00000000-0000-0000-0000-000000000001" +_SECOND_FAILURE_ARTIFACT_ID = "00000000-0000-0000-0000-000000000002" +_CURATION_ID = "00000000-0000-0000-0000-000000000003" class OpenFlywheelMcpModule(Protocol): @@ -69,6 +79,8 @@ def _preparation_service(self) -> WorkspacePreparationService: ... def _failure_service(self) -> FailureWorkspaceService: ... + def _curation_service(self) -> FailureCurationService: ... + def _program_template(self, name: str) -> str: ... def prepare_workspace( @@ -120,6 +132,11 @@ def mine_failure_patterns( request: MineFailurePatternsInput, ) -> FailurePatternMiningObservation: ... + def record_failure_curation( + self, + request: RecordFailureCurationInput, + ) -> FailureCurationObservation: ... + class _FakeOutcomeStore: def __init__(self) -> None: @@ -160,6 +177,16 @@ def record(self, request: RecordFailureInput) -> FailureRecordObservation: return self.observation +class _FakeCurationService: + def __init__(self, observation: FailureCurationObservation) -> None: + self.observation = observation + self.requests: list[RecordFailureCurationInput] = [] + + def record(self, request: RecordFailureCurationInput) -> FailureCurationObservation: + self.requests.append(request) + return self.observation + + def _module() -> OpenFlywheelMcpModule: return cast(OpenFlywheelMcpModule, importlib.import_module("ofw.mcp")) @@ -241,6 +268,39 @@ def _inconclusive_failure_request(root: Path) -> RecordFailureInput: ) +def _curation_request(root: Path) -> RecordFailureCurationInput: + return RecordFailureCurationInput( + workspace_root=root, + source_artifact_ids=(_FAILURE_ARTIFACT_ID, _SECOND_FAILURE_ARTIFACT_ID), + groups=( + FailureGroupInput( + pattern_key="finalizes-before-verification", + title="Finalizes before verifying state", + mechanism="A successful mutation is treated as task completion.", + prevention="Require a state read before finalization.", + target_component=ComponentKind.PROMPT, + failure_artifact_ids=(_FAILURE_ARTIFACT_ID, _SECOND_FAILURE_ARTIFACT_ID), + ), + ), + deferred=(), + ) + + +def _curation_observation() -> FailureCurationObservation: + relative_path = Path(f".workspace/failure-curations/{_CURATION_ID}.json") + return FailureCurationObservation( + status=FailureCurationStatus.SUCCESS, + summary="Stored one evidence-bound failure group and 0 deferred failures.", + next_actions=("Use one recorded group to form the next harness hypothesis.",), + artifacts=(str(relative_path), _CURATION_ID), + curation_id=_CURATION_ID, + relative_path=relative_path, + source_failure_count=2, + group_count=1, + deferred_count=0, + ) + + def test_mcp_exposes_scoped_read_and_recording_tools() -> None: tools = asyncio.run(_server().list_tools()) @@ -253,6 +313,7 @@ def test_mcp_exposes_scoped_read_and_recording_tools() -> None: "record_outcome", "record_failure", "mine_failure_patterns", + "record_failure_curation", ] assert tuple(map(_annotation_flags, tools)) == ( (False, False, True), @@ -263,6 +324,7 @@ def test_mcp_exposes_scoped_read_and_recording_tools() -> None: (False, False, True), (False, False, True), (True, False, True), + (False, False, True), ) @@ -504,3 +566,17 @@ def test_mine_failure_patterns_passes_one_bounded_object_to_the_service( assert result.source_artifact_count == 1 assert result.patterns == () assert result.inconclusive_artifact_ids == (recorded.artifact_id,) + + +def test_record_failure_curation_passes_one_strict_object_to_the_service( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + module = _module() + expected = _curation_observation() + service = _FakeCurationService(expected) + request = _curation_request(tmp_path) + monkeypatch.setattr(module, "_curation_service", lambda: service) + + assert module.record_failure_curation(request) == expected + assert service.requests == [request] diff --git a/tests/test_plugin_packaging.py b/tests/test_plugin_packaging.py index 9653e81..44808ae 100644 --- a/tests/test_plugin_packaging.py +++ b/tests/test_plugin_packaging.py @@ -28,7 +28,7 @@ def test_openflywheel_mcp_uses_pinned_portable_runtime() -> None: server = manifest.mcpServers["openflywheel"] assert ( - "git+https://github.com/divo12/OpenFlyWheel.git@ab0ef62cbe1e6cddf0bfd8ec61374d10120c61aa" + "git+https://github.com/divo12/OpenFlyWheel.git@726594a56c6683d1d32207a462fe524b202f4b8b" in server.args ) assert "openflywheel-mcp" in server.args diff --git a/tests/test_program_templates.py b/tests/test_program_templates.py index 0557d16..d7fe654 100644 --- a/tests/test_program_templates.py +++ b/tests/test_program_templates.py @@ -20,6 +20,11 @@ def test_packaged_program_template_matches_plugin_asset(name: str) -> None: "$failure-miner", "record_failure", ".workspace/failures/", + "$failure-curator", + "record_failure_curation", + ".workspace/failure-curations/", + "Do not glob", + "stop before forming a harness hypothesis", "Do not copy Langfuse trace payloads", ), ) diff --git a/tests/test_typing.py b/tests/test_typing.py index 2f8004e..8f80686 100644 --- a/tests/test_typing.py +++ b/tests/test_typing.py @@ -50,6 +50,20 @@ def test_namespace_exports_failure_pattern_contract() -> None: assert "MineFailurePatternsInput" in package.__all__ +def test_namespace_exports_failure_curation_contract() -> None: + expected = { + "DeferredFailure", + "FailureCuration", + "FailureCurationErrorCode", + "FailureCurationFailure", + "FailureGroup", + "FailureGroupMember", + "FailureSource", + } + + assert expected <= set(package.__all__) + + def test_namespace_exports_workspace_preparation_contract() -> None: assert "PreparationErrorCode" in package.__all__ assert "PreparationPhase" in package.__all__