diff --git a/src/ofw/__init__.py b/src/ofw/__init__.py index c59e32c..fa338c7 100644 --- a/src/ofw/__init__.py +++ b/src/ofw/__init__.py @@ -13,20 +13,54 @@ ) from ofw.contracts import ( - AssetAccess, - ComponentKind, GitCommit, - HarnessAsset, - HarnessComponent, HarnessErrorCode, HarnessRevision, HarnessRevisionId, HarnessValidationError, RepositorySnapshot, Sha256Digest, - WorkspaceFile, ) -from ofw.harness import EditableFile, Harness, Subagent, Tool, editable +from ofw.mine import ( + AdaptationRequest, + AdaptationResult, + BehaviorObservation, + CompletionCheck, + CompletionStatus, + Confidence, + ConstraintKind, + EnvironmentCheckId, + EnvironmentCheckRequest, + EnvironmentSource, + EnvironmentSourceId, + EnvironmentSourceKind, + EnvironmentVerification, + EvidenceKind, + EvidenceRecordId, + EvidenceReference, + FailureBehavior, + FailureBehaviorKind, + FailureMiningResult, + FailureMiningRun, + FailurePhase, + FailureSource, + FailureSourceId, + FailureSourceKind, + Mine, + MiningContext, + MiningInvalidReason, + MiningNomination, + MiningTask, + MiningVerdict, + RecoveryStatus, + RequiredOutcome, + TaskConstraint, + TaskId, + ToolAccess, + ToolCapability, + ToolName, + TraceMiningCase, +) from ofw.observability.langfuse import ( CollectionError, CollectionErrorCode, @@ -50,6 +84,7 @@ read_trace_observations, search_observation_content, ) +from ofw.repository import process_repository from ofw.runtime import ( CanaryCase, CaseId, @@ -79,9 +114,6 @@ class _OfwNamespace: CanaryCase = CanaryCase CaseId = CaseId - def editable(self, path: Path) -> EditableFile: - return editable(path) - def collect( self, revision: HarnessRevision, @@ -113,25 +145,42 @@ def read_observation_content( ) -> ObservationContent: return read_observation_content(collection, reference) - ofw = _OfwNamespace() __all__ = [ - "AssetAccess", + "AdaptationRequest", + "AdaptationResult", + "BehaviorObservation", "CanaryCase", "CaseId", - "ComponentKind", "CollectionError", "CollectionErrorCode", "CollectionResult", "CommandLoop", "CommandVerifier", + "CompletionCheck", + "CompletionStatus", + "Confidence", + "ConstraintKind", "E2BSandbox", - "EditableFile", + "EnvironmentCheckId", + "EnvironmentCheckRequest", + "EnvironmentSource", + "EnvironmentSourceId", + "EnvironmentSourceKind", + "EnvironmentVerification", + "EvidenceKind", + "EvidenceRecordId", + "EvidenceReference", + "FailureMiningResult", + "FailureMiningRun", + "FailureBehavior", + "FailureBehaviorKind", + "FailurePhase", + "FailureSource", + "FailureSourceId", + "FailureSourceKind", "GitCommit", - "Harness", - "HarnessAsset", - "HarnessComponent", "HarnessErrorCode", "HarnessRevision", "HarnessRevisionId", @@ -141,6 +190,12 @@ def read_observation_content( "LangfuseProject", "LangfuseSpan", "ModelFingerprint", + "Mine", + "MiningContext", + "MiningInvalidReason", + "MiningNomination", + "MiningTask", + "MiningVerdict", "ObservationContent", "ObservationContentField", "ObservationContentHit", @@ -148,24 +203,29 @@ def read_observation_content( "ObservationContentQuery", "ObservationContentReference", "RepositorySnapshot", + "RecoveryStatus", + "RequiredOutcome", "ProcessCommand", "ProcessLimits", "RunErrorCode", "RunResult", "RunStatus", "Sha256Digest", - "Subagent", - "Tool", + "TaskConstraint", + "TaskId", + "ToolAccess", + "ToolCapability", + "ToolName", "TraceWindow", + "TraceMiningCase", "VerifierResult", "VerifierVerdict", - "WorkspaceFile", "collect", - "editable", "get_client", "is_default_export_span", "observe", "ofw", + "process_repository", "propagate_attributes", "read_observation_content", "read_trace_observations", diff --git a/src/ofw/contracts.py b/src/ofw/contracts.py index 28f9741..ccb0832 100644 --- a/src/ofw/contracts.py +++ b/src/ofw/contracts.py @@ -1,4 +1,4 @@ -"""Immutable, language-neutral harness component contracts.""" +"""Immutable, language-neutral repository revision contracts.""" from __future__ import annotations @@ -14,43 +14,14 @@ class HarnessSchemaVersion(IntEnum): V1 = 1 -class ComponentKind(StrEnum): - PROMPT = "prompt" - TOOL = "tool" - SKILL = "skill" - SUBAGENT = "subagent" - MIDDLEWARE = "middleware" - - -class AssetAccess(StrEnum): - FROZEN = "frozen" - FIT_EDITABLE = "fit_editable" - - class HarnessErrorCode(StrEnum): INVALID_NAME = "invalid_name" INVALID_SOURCE = "invalid_source" ROOT_NOT_FOUND = "root_not_found" ROOT_NOT_DIRECTORY = "root_not_directory" - PROMPT_REQUIRED = "prompt_required" - MISSING_ASSET = "missing_asset" - PATH_OUTSIDE_ROOT = "path_outside_root" - NOT_A_FILE = "not_a_file" - DUPLICATE_ASSET = "duplicate_asset" - CONFLICTING_ACCESS = "conflicting_access" - COMPONENT_OVERLAP = "component_overlap" - INVALID_TOOL_NAME = "invalid_tool_name" - DUPLICATE_TOOL = "duplicate_tool" - INVALID_SUBAGENT_NAME = "invalid_subagent_name" - DUPLICATE_SUBAGENT = "duplicate_subagent" GIT_REPOSITORY_REQUIRED = "git_repository_required" GIT_COMMAND_FAILED = "git_command_failed" MANIFEST_WRITE_FAILED = "manifest_write_failed" - SENSITIVE_ASSET = "sensitive_asset" - RUNTIME_INCOMPLETE = "runtime_incomplete" - CANARY_FAILED = "canary_failed" - DUPLICATE_VERIFIER = "duplicate_verifier" - RUNTIME_INVALID = "runtime_invalid" class HarnessValidationError(Exception): @@ -88,26 +59,6 @@ def __str__(self) -> str: return self.value -@dataclass(frozen=True, slots=True) -class WorkspaceFile: - relative_path: Path - - -@dataclass(frozen=True, slots=True) -class HarnessAsset: - name: str | None - access: AssetAccess - source: WorkspaceFile - digest: Sha256Digest - - -@dataclass(frozen=True, slots=True) -class HarnessComponent: - kind: ComponentKind - assets: tuple[HarnessAsset, ...] - digest: Sha256Digest - - @dataclass(frozen=True, slots=True) class RepositorySnapshot: commit: GitCommit @@ -115,31 +66,12 @@ class RepositorySnapshot: dirty_digest: Sha256Digest | None -@dataclass(frozen=True, slots=True) -class RuntimeConfiguration: - execution: Sha256Digest - lifecycle: Sha256Digest - verifiers: tuple[Sha256Digest, ...] - - def canonical_json(self) -> str: - verifiers = ",".join(_quote(str(verifier)) for verifier in self.verifiers) - return ( - "{" - f'"execution":{_quote(str(self.execution))},' - f'"lifecycle":{_quote(str(self.lifecycle))},' - f'"verifiers":[{verifiers}]' - "}" - ) - - @dataclass(frozen=True, slots=True) class HarnessRevisionContent: schema_version: HarnessSchemaVersion harness_name: str repository: RepositorySnapshot - components: tuple[HarnessComponent, ...] observability: LangfuseConnectionManifest | None - runtime: RuntimeConfiguration | None def canonical_json(self) -> str: return _render_content(self) @@ -152,47 +84,13 @@ class HarnessRevision: harness_name: str root: Path repository: RepositorySnapshot - components: tuple[HarnessComponent, ...] observability: LangfuseConnectionManifest | None - runtime: RuntimeConfiguration | None - canary_digest: Sha256Digest | None @property def manifest_path(self) -> Path: return self.root / ".ofw" / "revisions" / str(self.id) / "manifest.json" - @property - def canary_path(self) -> Path: - return self.manifest_path.with_name("canary.json") - - @property - def assets(self) -> tuple[HarnessAsset, ...]: - return tuple(asset for component in self.components for asset in component.assets) - - @property - def editable_files(self) -> tuple[Path, ...]: - return tuple( - asset.source.relative_path - for asset in self.assets - if asset.access is AssetAccess.FIT_EDITABLE - ) - - @property - def frozen_files(self) -> tuple[Path, ...]: - return tuple( - asset.source.relative_path - for asset in self.assets - if asset.access is AssetAccess.FROZEN - ) - - def component(self, kind: ComponentKind) -> HarnessComponent | None: - for component in self.components: - if component.kind is kind: - return component - return None - def to_json(self) -> str: - components = ",".join(_render_component(component) for component in self.components) return ( "{" f'"schema_version":{int(self.schema_version)},' @@ -200,24 +98,18 @@ def to_json(self) -> str: f'"harness_name":{_quote(self.harness_name)},' f'"root":{_quote(self.root.as_posix())},' f'"repository":{_render_repository(self.repository)},' - f'"observability":{_render_observability(self.observability)},' - f'"runtime":{_render_runtime(self.runtime)},' - f'"canary_digest":{_render_digest(self.canary_digest)},' - f'"components":[{components}]' + f'"observability":{_render_observability(self.observability)}' "}" ) def _render_content(content: HarnessRevisionContent) -> str: - components = ",".join(_render_component(component) for component in content.components) return ( "{" f'"schema_version":{int(content.schema_version)},' f'"harness_name":{_quote(content.harness_name)},' f'"repository":{_render_repository(content.repository)},' - f'"observability":{_render_observability(content.observability)},' - f'"runtime":{_render_runtime(content.runtime)},' - f'"components":[{components}]' + f'"observability":{_render_observability(content.observability)}' "}" ) @@ -240,36 +132,5 @@ def _render_observability(connection: LangfuseConnectionManifest | None) -> str: return "null" if connection is None else connection.to_json() -def _render_runtime(runtime: RuntimeConfiguration | None) -> str: - return "null" if runtime is None else runtime.canonical_json() - - -def _render_digest(digest: Sha256Digest | None) -> str: - return "null" if digest is None else _quote(str(digest)) - - -def _render_component(component: HarnessComponent) -> str: - assets = ",".join(_render_asset(asset) for asset in component.assets) - return ( - "{" - f'"kind":{_quote(component.kind.value)},' - f'"digest":{_quote(str(component.digest))},' - f'"assets":[{assets}]' - "}" - ) - - -def _render_asset(asset: HarnessAsset) -> str: - name = "null" if asset.name is None else _quote(asset.name) - return ( - "{" - f'"name":{name},' - f'"access":{_quote(asset.access.value)},' - f'"source":{{"relative_path":{_quote(asset.source.relative_path.as_posix())}}},' - f'"digest":{_quote(str(asset.digest))}' - "}" - ) - - def _quote(value: str) -> str: return json.dumps(value, ensure_ascii=False, separators=(",", ":")) diff --git a/src/ofw/harness.py b/src/ofw/harness.py deleted file mode 100644 index 831da84..0000000 --- a/src/ofw/harness.py +++ /dev/null @@ -1,498 +0,0 @@ -"""Compile a file-level harness workspace into an immutable revision.""" - -from __future__ import annotations - -import hashlib -import logging -import os -import re -import subprocess # nosec B404 -import tempfile -from dataclasses import dataclass, field -from pathlib import Path - -from ofw.contracts import ( - AssetAccess, - ComponentKind, - GitCommit, - HarnessAsset, - HarnessComponent, - HarnessErrorCode, - HarnessRevision, - HarnessRevisionContent, - HarnessRevisionId, - HarnessSchemaVersion, - HarnessValidationError, - RepositorySnapshot, - RuntimeConfiguration, - Sha256Digest, - WorkspaceFile, -) -from ofw.observability.langfuse.contracts import LangfuseProject -from ofw.runtime import ( - CanaryCase, - CanaryReport, - ExecutionEnvironment, - LifecycleAdapter, - VerifierAdapter, - run_canary, - runtime_configuration, -) - -logger = logging.getLogger(__name__) - -_NAME_PATTERN = re.compile(r"[a-z0-9]+(?:-[a-z0-9]+)*") -_NAMED_SOURCE_PATTERN = re.compile(r"[a-z][a-z0-9_-]*") - - -@dataclass(frozen=True, slots=True) -class EditableFile: - path: Path - - -@dataclass(frozen=True, slots=True) -class Tool: - name: str - source: Path | EditableFile - - def __post_init__(self) -> None: - _validate_named_source(self.name, self.source, HarnessErrorCode.INVALID_TOOL_NAME) - - -@dataclass(frozen=True, slots=True) -class Subagent: - name: str - source: Path | EditableFile - - def __post_init__(self) -> None: - _validate_named_source(self.name, self.source, HarnessErrorCode.INVALID_SUBAGENT_NAME) - - -@dataclass(frozen=True, slots=True) -class _FileRegistration: - component: ComponentKind - name: str | None - path: Path - access: AssetAccess - - -@dataclass(frozen=True, slots=True) -class _CompiledAsset: - component: ComponentKind - asset: HarnessAsset - - -def editable(path: Path) -> EditableFile: - """Grant Fit authority to edit one workspace file.""" - if not isinstance(path, Path): - raise HarnessValidationError(HarnessErrorCode.INVALID_SOURCE, repr(path)) - return EditableFile(path=path) - - -@dataclass(slots=True) -class Harness: - """Mutable component registry; ``process`` returns an immutable revision.""" - - name: str - root: Path - _files: list[_FileRegistration] = field(default_factory=list, init=False, repr=False) - _observability: LangfuseProject | None = field(default=None, init=False, repr=False) - _execution: ExecutionEnvironment | None = field(default=None, init=False, repr=False) - _lifecycle: LifecycleAdapter | None = field(default=None, init=False, repr=False) - _verifiers: list[VerifierAdapter] = field(default_factory=list, init=False, repr=False) - - def __post_init__(self) -> None: - if _NAME_PATTERN.fullmatch(self.name) is None: - raise HarnessValidationError(HarnessErrorCode.INVALID_NAME, self.name) - if not isinstance(self.root, Path): - raise HarnessValidationError(HarnessErrorCode.INVALID_SOURCE, repr(self.root)) - - def connect_prompt(self, *sources: Path | EditableFile) -> Harness: - self._register_files(ComponentKind.PROMPT, sources) - return self - - def connect_tools(self, *tools: Tool) -> Harness: - for tool in tools: - if any( - registration.component is ComponentKind.TOOL and registration.name == tool.name - for registration in self._files - ): - raise HarnessValidationError(HarnessErrorCode.DUPLICATE_TOOL, tool.name) - self._files.append(_registration(ComponentKind.TOOL, tool.source, tool.name)) - return self - - def connect_skills(self, *sources: Path | EditableFile) -> Harness: - self._register_files(ComponentKind.SKILL, sources) - return self - - def connect_subagents(self, *subagents: Subagent) -> Harness: - for subagent in subagents: - if any( - registration.component is ComponentKind.SUBAGENT - and registration.name == subagent.name - for registration in self._files - ): - raise HarnessValidationError( - HarnessErrorCode.DUPLICATE_SUBAGENT, - subagent.name, - ) - self._files.append( - _registration(ComponentKind.SUBAGENT, subagent.source, subagent.name) - ) - return self - - def connect_middleware(self, *sources: Path | EditableFile) -> Harness: - self._register_files(ComponentKind.MIDDLEWARE, sources) - return self - - def connect_observability(self, project: LangfuseProject) -> Harness: - self._observability = project - return self - - def connect_execute(self, environment: ExecutionEnvironment) -> Harness: - self._execution = environment - return self - - def connect_lifecycle(self, lifecycle: LifecycleAdapter) -> Harness: - self._lifecycle = lifecycle - return self - - def connect_verifiers(self, *verifiers: VerifierAdapter) -> Harness: - for verifier in verifiers: - if any(existing.name == verifier.name for existing in self._verifiers): - raise HarnessValidationError(HarnessErrorCode.DUPLICATE_VERIFIER, verifier.name) - self._verifiers.append(verifier) - return self - - def _register_files( - self, - component: ComponentKind, - sources: tuple[Path | EditableFile, ...], - ) -> None: - for source in sources: - self._files.append(_registration(component, source, None)) - - def process(self, *, canary: CanaryCase | None = None) -> HarnessRevision: - logger.debug("Compiling harness revision: %s", self.name) - root = _resolve_root(self.root) - if not _has_component(self._files, ComponentKind.PROMPT): - raise HarnessValidationError(HarnessErrorCode.PROMPT_REQUIRED, self.name) - - components = _compile_components(root, self._files) - repository = _snapshot_repository(root) - runtime = self._runtime(root) - content = HarnessRevisionContent( - schema_version=HarnessSchemaVersion.V1, - harness_name=self.name, - repository=repository, - components=components, - observability=(None if self._observability is None else self._observability.manifest()), - runtime=runtime, - ) - revision = _revision_from_content(content, root) - report: CanaryReport | None = None - if canary is not None: - if self._execution is None or self._lifecycle is None or not self._verifiers: - raise HarnessValidationError(HarnessErrorCode.RUNTIME_INCOMPLETE, self.name) - report = run_canary( - revision, - canary, - self._execution, - self._lifecycle, - tuple(self._verifiers), - ) - if not report.passed: - _write_canary(revision, report) - raise HarnessValidationError(HarnessErrorCode.CANARY_FAILED, canary.id.value) - revision = _revision_from_content(content, root, report.digest) - _write_manifest(revision) - if report is not None: - _write_canary(revision, report) - logger.debug("Compiled harness revision %s", revision.id) - return revision - - def _runtime(self, root: Path) -> RuntimeConfiguration | None: - connections = ( - self._execution is not None, - self._lifecycle is not None, - bool(self._verifiers), - ) - if not any(connections): - return None - if not all(connections) or self._execution is None or self._lifecycle is None: - raise HarnessValidationError(HarnessErrorCode.RUNTIME_INCOMPLETE, self.name) - try: - return runtime_configuration( - root, - self._execution, - self._lifecycle, - tuple(self._verifiers), - ) - except ValueError as error: - raise HarnessValidationError(HarnessErrorCode.RUNTIME_INVALID, self.name) from error - - -def _has_component(registrations: list[_FileRegistration], kind: ComponentKind) -> bool: - return any(registration.component is kind for registration in registrations) - - -def _revision_from_content( - content: HarnessRevisionContent, - root: Path, - canary_digest: Sha256Digest | None = None, -) -> HarnessRevision: - content_digest = _digest_text(content.canonical_json()) - return HarnessRevision( - schema_version=content.schema_version, - id=HarnessRevisionId(f"ofw_{content_digest.value[7:]}"), - harness_name=content.harness_name, - root=root, - repository=content.repository, - components=content.components, - observability=content.observability, - runtime=content.runtime, - canary_digest=canary_digest, - ) - - -def _resolve_root(root: Path) -> Path: - try: - resolved = root.expanduser().resolve(strict=True) - except FileNotFoundError as error: - raise HarnessValidationError(HarnessErrorCode.ROOT_NOT_FOUND, str(root)) from error - if not resolved.is_dir(): - raise HarnessValidationError(HarnessErrorCode.ROOT_NOT_DIRECTORY, str(resolved)) - return resolved - - -def _compile_components( - root: Path, - registrations: list[_FileRegistration], -) -> tuple[HarnessComponent, ...]: - compiled: list[_CompiledAsset] = [] - for registration in registrations: - resolved, relative = _resolve_file(root, registration.path) - compiled.append( - _CompiledAsset( - component=registration.component, - asset=HarnessAsset( - name=registration.name, - access=registration.access, - source=WorkspaceFile(relative_path=relative), - digest=_digest_file(resolved), - ), - ) - ) - compiled.sort(key=_compiled_asset_sort_key) - _validate_component_boundaries(compiled) - - components: list[HarnessComponent] = [] - for kind in ComponentKind: - assets = tuple(item.asset for item in compiled if item.component is kind) - if not assets: - continue - components.append( - HarnessComponent( - kind=kind, - assets=assets, - digest=_component_digest(kind, assets), - ) - ) - return tuple(components) - - -def _validate_component_boundaries(compiled: list[_CompiledAsset]) -> None: - for index, item in enumerate(compiled): - for existing in compiled[:index]: - if existing.asset.source.relative_path != item.asset.source.relative_path: - continue - if existing.component is not item.component: - raise HarnessValidationError( - HarnessErrorCode.COMPONENT_OVERLAP, - item.asset.source.relative_path.as_posix(), - ) - if ( - item.component in (ComponentKind.TOOL, ComponentKind.SUBAGENT) - and existing.asset.name != item.asset.name - ): - continue - code = ( - HarnessErrorCode.DUPLICATE_ASSET - if existing.asset.access is item.asset.access - else HarnessErrorCode.CONFLICTING_ACCESS - ) - raise HarnessValidationError(code, item.asset.source.relative_path.as_posix()) - - -def _component_digest( - kind: ComponentKind, - assets: tuple[HarnessAsset, ...], -) -> Sha256Digest: - fields = [kind.value] - for asset in assets: - fields.extend( - ( - asset.name or "", - asset.access.value, - asset.source.relative_path.as_posix(), - str(asset.digest), - ) - ) - return _digest_text("\0".join(fields)) - - -def _compiled_asset_sort_key(item: _CompiledAsset) -> tuple[str, str, str]: - return ( - item.component.value, - item.asset.source.relative_path.as_posix(), - item.asset.name or "", - ) - - -def _registration( - component: ComponentKind, - source: Path | EditableFile, - name: str | None, -) -> _FileRegistration: - if isinstance(source, EditableFile): - return _FileRegistration(component, name, source.path, AssetAccess.FIT_EDITABLE) - if isinstance(source, Path): - return _FileRegistration(component, name, source, AssetAccess.FROZEN) - raise HarnessValidationError(HarnessErrorCode.INVALID_SOURCE, repr(source)) - - -def _validate_named_source( - name: str, - source: Path | EditableFile, - error_code: HarnessErrorCode, -) -> None: - if _NAMED_SOURCE_PATTERN.fullmatch(name) is None: - raise HarnessValidationError(error_code, name) - if not isinstance(source, (Path, EditableFile)): - raise HarnessValidationError(HarnessErrorCode.INVALID_SOURCE, repr(source)) - - -def _resolve_file(root: Path, source: Path) -> tuple[Path, Path]: - if _is_sensitive_path(source): - raise HarnessValidationError(HarnessErrorCode.SENSITIVE_ASSET, str(source)) - candidate = source if source.is_absolute() else root / source - try: - resolved = candidate.resolve(strict=True) - except FileNotFoundError as error: - raise HarnessValidationError(HarnessErrorCode.MISSING_ASSET, str(source)) from error - try: - relative = resolved.relative_to(root) - except ValueError as error: - raise HarnessValidationError( - HarnessErrorCode.PATH_OUTSIDE_ROOT, - str(source), - ) from error - if not resolved.is_file(): - raise HarnessValidationError(HarnessErrorCode.NOT_A_FILE, str(source)) - return resolved, Path(relative.as_posix()) - - -def _is_sensitive_path(path: Path) -> bool: - return any( - part == ".env" or (part.startswith(".env.") and part != ".env.example") - for part in path.parts - ) - - -def _snapshot_repository(root: Path) -> RepositorySnapshot: - top_level = _run_git(root, "rev-parse", "--show-toplevel", repository_probe=True) - try: - git_root = Path(top_level.decode().strip()).resolve(strict=True) - except (UnicodeDecodeError, FileNotFoundError) as error: - raise HarnessValidationError( - HarnessErrorCode.GIT_REPOSITORY_REQUIRED, - str(root), - ) from error - if git_root != root: - raise HarnessValidationError(HarnessErrorCode.GIT_REPOSITORY_REQUIRED, str(root)) - - commit_bytes = _run_git(root, "rev-parse", "HEAD") - try: - commit = GitCommit(commit_bytes.decode().strip()) - except UnicodeDecodeError as error: - raise HarnessValidationError( - HarnessErrorCode.GIT_COMMAND_FAILED, "rev-parse HEAD" - ) from error - diff = _run_git(root, "diff", "--binary", "--no-ext-diff", "HEAD", "--") - return RepositorySnapshot( - commit=commit, - is_dirty=bool(diff), - dirty_digest=_digest_bytes(diff) if diff else None, - ) - - -def _run_git( - root: Path, - *arguments: str, - repository_probe: bool = False, -) -> bytes: - try: - # Every argument is selected internally; the validated root is one argv value. - result: subprocess.CompletedProcess[bytes] = subprocess.run( # nosec B603 - ("git", "-C", str(root), *arguments), - check=False, - capture_output=True, - ) - except OSError as error: - raise HarnessValidationError(HarnessErrorCode.GIT_COMMAND_FAILED, arguments[0]) from error - if result.returncode != 0: - code = ( - HarnessErrorCode.GIT_REPOSITORY_REQUIRED - if repository_probe - else HarnessErrorCode.GIT_COMMAND_FAILED - ) - raise HarnessValidationError(code, arguments[0]) - return result.stdout - - -def _digest_file(path: Path) -> Sha256Digest: - try: - return _digest_bytes(path.read_bytes()) - except OSError as error: - raise HarnessValidationError(HarnessErrorCode.MISSING_ASSET, str(path)) from error - - -def _digest_text(value: str) -> Sha256Digest: - return _digest_bytes(value.encode()) - - -def _digest_bytes(value: bytes) -> Sha256Digest: - return Sha256Digest(f"sha256:{hashlib.sha256(value).hexdigest()}") - - -def _write_manifest(revision: HarnessRevision) -> None: - _write_revision_file(revision.manifest_path, f"{revision.to_json()}\n") - - -def _write_canary(revision: HarnessRevision, report: CanaryReport) -> None: - _write_revision_file(revision.canary_path, f"{report.to_json()}\n") - - -def _write_revision_file(path: Path, payload: str) -> None: - try: - path.parent.mkdir(parents=True, exist_ok=True) - descriptor, temporary_name = tempfile.mkstemp( - dir=path.parent, - prefix=f".{path.stem}-", - suffix=".json", - text=True, - ) - temporary_path = Path(temporary_name) - try: - with os.fdopen(descriptor, "w", encoding="utf-8") as stream: - stream.write(payload) - stream.flush() - os.fsync(stream.fileno()) - temporary_path.replace(path) - finally: - temporary_path.unlink(missing_ok=True) - except OSError as error: - raise HarnessValidationError( - HarnessErrorCode.MANIFEST_WRITE_FAILED, - str(path), - ) from error diff --git a/src/ofw/mine.py b/src/ofw/mine.py new file mode 100644 index 0000000..ef25f33 --- /dev/null +++ b/src/ofw/mine.py @@ -0,0 +1,998 @@ +"""Evidence-backed failure mining over full Langfuse trajectories.""" + +from __future__ import annotations + +import math +from dataclasses import dataclass, field +from datetime import datetime +from enum import StrEnum +from typing import Protocol + +from ofw.contracts import HarnessRevision, HarnessRevisionId, Sha256Digest +from ofw.observability.langfuse.contracts import CollectionError +from ofw.observability.langfuse.domain import ( + AttributionLevel, + CollectionResult, + ObservationContent, + ObservationContentField, + ObservationContentHit, + ObservationContentMatch, + ObservationContentQuery, + ObservationId, + ObservationRecord, + TraceId, + TraceRecord, +) +from ofw.observability.langfuse.store import CollectionStore + + +class FailureSourceKind(StrEnum): + HUMAN_FEEDBACK = "human_feedback" + USER_CORRECTION = "user_correction" + TRUSTED_SCORE = "trusted_score" + DOWNSTREAM_FAILURE = "downstream_failure" + INCIDENT = "incident" + ROLLBACK = "rollback" + REOPENED_WORK = "reopened_work" + ENVIRONMENT_MISMATCH = "environment_mismatch" + AGENT_ERROR = "agent_error" + + +class EnvironmentSourceKind(StrEnum): + RECORDED_STATE = "recorded_state" + AUDIT_LOG = "audit_log" + PRODUCTION_API = "production_api" + DETERMINISTIC_CHECK = "deterministic_check" + + +class EvidenceKind(StrEnum): + TRAJECTORY = "trajectory" + ENVIRONMENT = "environment" + PRODUCTION_SIGNAL = "production_signal" + + +class CompletionStatus(StrEnum): + COMPLETED = "completed" + NOT_COMPLETED = "not_completed" + UNKNOWN = "unknown" + + +class MiningVerdict(StrEnum): + CONFIRMED_FAILURE = "confirmed_failure" + NO_FAILURE = "no_failure" + AMBIGUOUS = "ambiguous" + INVALID = "invalid" + + +class MiningInvalidReason(StrEnum): + REVISION_MISMATCH = "revision_mismatch" + TRACE_NOT_FOUND = "trace_not_found" + CORRUPT_TRACE = "corrupt_trace" + JUDGE_OUTPUT = "judge_output" + + +class ToolStatus(StrEnum): + OK = "ok" + NOT_FOUND = "not_found" + UNAVAILABLE = "unavailable" + BLOCKED = "blocked" + ERROR = "error" + + +class ToolAction(StrEnum): + SEARCH_TRAJECTORY = "search_trajectory" + SEARCH_PRIOR_TRAJECTORIES = "search_prior_trajectories" + READ_TRAJECTORY = "read_trajectory" + VERIFY_ENVIRONMENT = "verify_environment" + ADAPT = "adapt" + RETURN_VERDICT = "return_verdict" + + +class ConstraintKind(StrEnum): + TIME = "time" + RESOURCE = "resource" + NETWORK = "network" + ACCESS = "access" + POLICY = "policy" + + +class ToolAccess(StrEnum): + READ_ONLY = "read_only" + MUTATING = "mutating" + + +class FailureBehaviorKind(StrEnum): + OUTCOME_MISMATCH = "outcome_mismatch" + FALSE_COMPLETION = "false_completion" + REQUIRED_ACTION_OMITTED = "required_action_omitted" + FORBIDDEN_STATE_CHANGE = "forbidden_state_change" + UNRECOVERED_ACTION_FAILURE = "unrecovered_action_failure" + NO_PROGRESS_LOOP = "no_progress_loop" + ABANDONED_BEFORE_COMPLETION = "abandoned_before_completion" + + +class FailurePhase(StrEnum): + ACTION = "action" + RECOVERY = "recovery" + COMPLETION = "completion" + VERIFICATION = "verification" + + +class RecoveryStatus(StrEnum): + RECOVERED = "recovered" + NOT_RECOVERED = "not_recovered" + UNKNOWN = "unknown" + + +@dataclass(frozen=True, slots=True) +class FailureSourceId: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class EnvironmentSourceId: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class EnvironmentCheckId: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class TaskId: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class ToolName: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class EvidenceRecordId: + value: str + + def __post_init__(self) -> None: + _require_identifier(self.value) + + +@dataclass(frozen=True, slots=True) +class Confidence: + value: float + + def __post_init__(self) -> None: + if not math.isfinite(self.value) or not 0.0 <= self.value <= 1.0: + raise ValueError("confidence must be finite and between zero and one") + + +@dataclass(frozen=True, slots=True) +class EvidenceReference: + kind: EvidenceKind + record_id: EvidenceRecordId + digest: Sha256Digest + + +@dataclass(frozen=True, slots=True) +class RequiredOutcome: + check_id: EnvironmentCheckId + source_id: EnvironmentSourceId + description: str + + def __post_init__(self) -> None: + _require_text(self.description, "required outcome description") + + +@dataclass(frozen=True, slots=True) +class TaskConstraint: + kind: ConstraintKind + description: str + + def __post_init__(self) -> None: + _require_text(self.description, "task constraint description") + + +@dataclass(frozen=True, slots=True) +class MiningTask: + id: TaskId + intent: str + required_outcomes: tuple[RequiredOutcome, ...] + constraints: tuple[TaskConstraint, ...] = () + + def __post_init__(self) -> None: + _require_text(self.intent, "task intent") + if not self.required_outcomes: + raise ValueError("mining task requires a required outcome") + if len({item.check_id for item in self.required_outcomes}) != len( + self.required_outcomes + ): + raise ValueError("required outcome ids must be unique") + + +@dataclass(frozen=True, slots=True) +class ToolCapability: + name: ToolName + access: ToolAccess + + +@dataclass(frozen=True, slots=True) +class FailureSource: + id: FailureSourceId + kind: FailureSourceKind + trace_id: TraceId + observed_at: datetime + summary: str + evidence: tuple[EvidenceReference, ...] + + def __post_init__(self) -> None: + _require_text(self.summary, "failure source summary") + if not self.evidence or any( + item.kind is not EvidenceKind.PRODUCTION_SIGNAL for item in self.evidence + ): + raise ValueError("failure source requires production-signal evidence") + + +@dataclass(frozen=True, slots=True) +class EnvironmentSource: + id: EnvironmentSourceId + kind: EnvironmentSourceKind + summary: str + + def __post_init__(self) -> None: + _require_text(self.summary, "environment source summary") + + +@dataclass(frozen=True, slots=True) +class MiningNomination: + trace_id: TraceId + task: MiningTask + sources: tuple[FailureSource, ...] + environment_sources: tuple[EnvironmentSource, ...] + available_tools: tuple[ToolCapability, ...] = () + initial_state_evidence: tuple[EvidenceReference, ...] = () + + def __post_init__(self) -> None: + if not self.sources: + raise ValueError("mining nomination requires a failure source") + if any(source.trace_id != self.trace_id for source in self.sources): + raise ValueError("failure source trace does not match nomination") + if len({source.id for source in self.sources}) != len(self.sources): + raise ValueError("failure source ids must be unique") + source_ids = {source.id for source in self.environment_sources} + if len(source_ids) != len(self.environment_sources): + raise ValueError("environment source ids must be unique") + if any(outcome.source_id not in source_ids for outcome in self.task.required_outcomes): + raise ValueError("required outcome references an undeclared environment source") + if len({tool.name for tool in self.available_tools}) != len(self.available_tools): + raise ValueError("tool capabilities must be unique") + if any( + item.kind is not EvidenceKind.ENVIRONMENT + for item in self.initial_state_evidence + ): + raise ValueError("initial state requires environment evidence") + + +@dataclass(frozen=True, slots=True) +class MiningContext: + revision_id: HarnessRevisionId + trace_id: TraceId + trace_digest: Sha256Digest + observation_ids: tuple[ObservationId, ...] + session_id: str | None + environment_name: str | None + release: str | None + available_tools: tuple[ToolCapability, ...] + environment_sources: tuple[EnvironmentSource, ...] + initial_state_evidence: tuple[EvidenceReference, ...] + + def __post_init__(self) -> None: + if not self.observation_ids or len(set(self.observation_ids)) != len( + self.observation_ids + ): + raise ValueError("mining context requires unique observation ids") + + +@dataclass(frozen=True, slots=True) +class BehaviorObservation: + kind: FailureBehaviorKind + phase: FailurePhase + first_observation_id: ObservationId + last_observation_id: ObservationId | None + recovery_status: RecoveryStatus + evidence: tuple[EvidenceReference, ...] + + def __post_init__(self) -> None: + if not self.evidence or any( + item.kind is not EvidenceKind.TRAJECTORY for item in self.evidence + ): + raise ValueError("failure behavior observation requires trajectory evidence") + + +@dataclass(frozen=True, slots=True) +class FailureBehavior: + primary: FailureBehaviorKind + summary: str + observations: tuple[BehaviorObservation, ...] + + def __post_init__(self) -> None: + _require_text(self.summary, "failure behavior summary") + if not self.observations: + raise ValueError("failure behavior requires observations") + if not any(item.kind is self.primary for item in self.observations): + raise ValueError("primary failure behavior must appear in observations") + + +@dataclass(frozen=True, slots=True) +class TraceMiningCase: + task: MiningTask + context: MiningContext + sources: tuple[FailureSource, ...] + + @property + def trace_id(self) -> TraceId: + return self.context.trace_id + + +@dataclass(frozen=True, slots=True) +class EnvironmentCheckRequest: + source_id: EnvironmentSourceId + check_id: EnvironmentCheckId + + +@dataclass(frozen=True, slots=True) +class EnvironmentVerification: + status: CompletionStatus + observed_state: str | None + evidence: tuple[EvidenceReference, ...] + + def __post_init__(self) -> None: + if self.status is CompletionStatus.UNKNOWN: + if self.observed_state is not None or self.evidence: + raise ValueError("unknown environment state cannot carry evidence") + return + if self.observed_state is None or not self.observed_state.strip() or not self.evidence: + raise ValueError("known environment state requires observation and evidence") + if any(item.kind is not EvidenceKind.ENVIRONMENT for item in self.evidence): + raise ValueError("environment verification requires environment evidence") + + +@dataclass(frozen=True, slots=True) +class CompletionCheck: + check_id: EnvironmentCheckId + required_outcome: str + agent_claim: str | None + observed_state: str | None + status: CompletionStatus + evidence: tuple[EvidenceReference, ...] + + +@dataclass(frozen=True, slots=True) +class FailureMiningResult: + task: MiningTask + context: MiningContext | None + verdict: MiningVerdict + source_ids: tuple[FailureSourceId, ...] + completion_checks: tuple[CompletionCheck, ...] + failure_behavior: FailureBehavior | None + trajectory_evidence: tuple[EvidenceReference, ...] + environment_evidence: tuple[EvidenceReference, ...] + confidence: Confidence + unresolved_questions: tuple[str, ...] + invalid_reason: MiningInvalidReason | None + + def __post_init__(self) -> None: + if self.verdict is MiningVerdict.INVALID: + if self.invalid_reason is None or self.failure_behavior is not None: + raise ValueError("invalid result requires a reason and no failure behavior") + return + if self.context is None or self.invalid_reason is not None: + raise ValueError("valid result requires context and no invalid reason") + if not self.completion_checks: + raise ValueError("mining result requires completion checks") + self._validate_checks() + self._validate_behavior() + if self.verdict is MiningVerdict.CONFIRMED_FAILURE: + if ( + self.failure_behavior is None + or not any( + check.status is CompletionStatus.NOT_COMPLETED + for check in self.completion_checks + ) + or not self.trajectory_evidence + or not self.environment_evidence + ): + raise ValueError( + "confirmed failure requires behavior, failed completion, trajectory, " + "and environment evidence" + ) + elif self.verdict is MiningVerdict.NO_FAILURE: + if ( + self.failure_behavior is not None + or any( + check.status is not CompletionStatus.COMPLETED + for check in self.completion_checks + ) + or not self.trajectory_evidence + or not self.environment_evidence + ): + raise ValueError( + "no-failure verdict requires completed checks, no failure behavior, " + "and supporting evidence" + ) + elif ( + self.failure_behavior is not None + or not self.unresolved_questions + and not any( + check.status is CompletionStatus.UNKNOWN for check in self.completion_checks + ) + ): + raise ValueError("ambiguous verdict requires no behavior and an unresolved question") + + def _validate_checks(self) -> None: + outcomes = {item.check_id: item for item in self.task.required_outcomes} + if len(self.completion_checks) != len(outcomes): + raise ValueError("completion checks must cover every required outcome") + for check in self.completion_checks: + outcome = outcomes.get(check.check_id) + if outcome is None or check.required_outcome != outcome.description: + raise ValueError("completion check does not match the task") + + def _validate_behavior(self) -> None: + if self.context is None or self.failure_behavior is None: + return + ids = set(self.context.observation_ids) + evidence = set(self.trajectory_evidence) + for observation in self.failure_behavior.observations: + if ( + observation.first_observation_id not in ids + or observation.last_observation_id is not None + and observation.last_observation_id not in ids + or not set(observation.evidence).issubset(evidence) + ): + raise ValueError("failure behavior observation is outside the grounded context") + + +@dataclass(frozen=True, slots=True) +class FailureMiningRun: + revision_id: HarnessRevisionId + collection_digest: Sha256Digest + results: tuple[FailureMiningResult, ...] + + +@dataclass(frozen=True, slots=True) +class TrajectorySearchRequest: + text: str + field: ObservationContentField + limit: int + + def __post_init__(self) -> None: + _require_text(self.text, "trajectory search text") + if not 1 <= self.limit <= 100: + raise ValueError("trajectory search limit must be between 1 and 100") + + +@dataclass(frozen=True, slots=True) +class TrajectorySearchResult: + status: ToolStatus + summary: str + hits: tuple[ObservationContentHit, ...] + next_actions: tuple[ToolAction, ...] + artifacts: tuple[EvidenceReference, ...] + + +@dataclass(frozen=True, slots=True) +class TrajectoryPageRequest: + cursor: ObservationId | None + limit: int + + def __post_init__(self) -> None: + if not 1 <= self.limit <= 100: + raise ValueError("trajectory page limit must be between 1 and 100") + + +@dataclass(frozen=True, slots=True) +class TrajectoryObservation: + record: ObservationRecord + input_content: ObservationContent | None + output_content: ObservationContent | None + + +@dataclass(frozen=True, slots=True) +class TrajectoryPageResult: + status: ToolStatus + summary: str + observations: tuple[TrajectoryObservation, ...] + next_cursor: ObservationId | None + next_actions: tuple[ToolAction, ...] + artifacts: tuple[EvidenceReference, ...] + + +@dataclass(frozen=True, slots=True) +class EnvironmentCheckResult: + status: ToolStatus + summary: str + verification: EnvironmentVerification | None + next_actions: tuple[ToolAction, ...] + artifacts: tuple[EvidenceReference, ...] + + +@dataclass(frozen=True, slots=True) +class AdaptationRequest: + kinds: tuple[FailureSourceKind, ...] + limit: int + + def __post_init__(self) -> None: + if not self.kinds or len(set(self.kinds)) != len(self.kinds): + raise ValueError("adaptation requires unique signal kinds") + if not 1 <= self.limit <= 100: + raise ValueError("adaptation limit must be between 1 and 100") + + +@dataclass(frozen=True, slots=True) +class AdaptationResult: + status: ToolStatus + summary: str + signals: tuple[FailureSource, ...] + next_actions: tuple[ToolAction, ...] + + +class EnvironmentVerifier(Protocol): + def verify( + self, + request: EnvironmentCheckRequest, + source: EnvironmentSource, + outcome: RequiredOutcome, + ) -> EnvironmentVerification: ... + + +class FailureJudge(Protocol): + def investigate( + self, + case: TraceMiningCase, + tools: MiningTools, + ) -> FailureMiningResult: ... + + +@dataclass(slots=True) +class MiningTools: + case: TraceMiningCase + collection: CollectionResult + environment: EnvironmentVerifier + production_signals: tuple[FailureSource, ...] + _issued_evidence: list[EvidenceReference] = field(default_factory=list, init=False) + _read_trajectory_evidence: list[EvidenceReference] = field( + default_factory=list, init=False + ) + + @property + def issued_evidence(self) -> tuple[EvidenceReference, ...]: + return tuple(self._issued_evidence) + + @property + def read_trajectory_evidence(self) -> tuple[EvidenceReference, ...]: + return tuple(self._read_trajectory_evidence) + + def search_trajectory(self, request: TrajectorySearchRequest) -> TrajectorySearchResult: + return self._search(request, self.case.trace_id, False) + + def search_prior_trajectories( + self, request: TrajectorySearchRequest + ) -> TrajectorySearchResult: + return self._search(request, None, True) + + def _search( + self, + request: TrajectorySearchRequest, + trace_id: TraceId | None, + prior_only: bool, + ) -> TrajectorySearchResult: + store = CollectionStore(self.collection.store_path) + try: + limit = 100 if prior_only else request.limit + hits = self._search_phrase(store, request, trace_id, request.text, limit) + if not hits: + found: list[ObservationContentHit] = [] + for token in _search_tokens(request.text): + for hit in self._search_phrase(store, request, trace_id, token, limit): + if hit not in found: + found.append(hit) + if len(found) >= limit: + break + hits = tuple(found[:limit]) + hits = tuple(_focus_hit(store, self.collection, hit, request.text) for hit in hits) + except CollectionError: + return TrajectorySearchResult( + ToolStatus.ERROR, + "Trajectory search failed.", + (), + (ToolAction.RETURN_VERDICT,), + (), + ) + finally: + store.close() + if prior_only: + hits = tuple(hit for hit in hits if hit.trace_id != self.case.trace_id)[ + : request.limit + ] + artifacts = tuple( + EvidenceReference( + EvidenceKind.TRAJECTORY, + EvidenceRecordId(hit.observation_id.value), + hit.reference.digest, + ) + for hit in hits + ) + self._issued_evidence.extend(artifacts) + return TrajectorySearchResult( + ToolStatus.OK if hits else ToolStatus.NOT_FOUND, + f"Found {len(hits)} matching trajectory segments.", + hits, + (ToolAction.READ_TRAJECTORY,), + artifacts, + ) + + def _search_phrase( + self, + store: CollectionStore, + request: TrajectorySearchRequest, + trace_id: TraceId | None, + text: str, + limit: int, + ) -> tuple[ObservationContentHit, ...]: + return store.search_content( + self.collection.observation_sync_id, + ObservationContentQuery( + text=text, + match=ObservationContentMatch.TOKEN_PHRASE, + field=request.field, + trace_id=trace_id, + limit=limit, + maximum_excerpt_characters=1000, + ), + ) + + def read_trajectory(self, request: TrajectoryPageRequest) -> TrajectoryPageResult: + store = CollectionStore(self.collection.store_path) + try: + # ponytail: collection scan is simplest for local v0; add a trace SQL query + # if production profiles show this O(pages * records) path matters. + observations = tuple( + item + for item in store.observations(self.collection.observation_sync_id) + if item.trace_id == self.case.trace_id + ) + start = _page_start(observations, request.cursor) + if start is None: + return TrajectoryPageResult( + ToolStatus.BLOCKED, + "Cursor is outside the nominated trace.", + (), + None, + (ToolAction.RETURN_VERDICT,), + (), + ) + selected = observations[start : start + request.limit] + views = tuple(_read_observation(store, self.collection, item) for item in selected) + except CollectionError: + return TrajectoryPageResult( + ToolStatus.ERROR, + "Trajectory content could not be read.", + (), + None, + (ToolAction.RETURN_VERDICT,), + (), + ) + finally: + store.close() + has_more = start + len(selected) < len(observations) + next_cursor = observations[start + len(selected)].id if has_more else None + artifacts = tuple( + EvidenceReference( + EvidenceKind.TRAJECTORY, + EvidenceRecordId(observation.record.id.value), + observation.record.digest, + ) + for observation in views + ) + self._issued_evidence.extend(artifacts) + self._read_trajectory_evidence.extend(artifacts) + return TrajectoryPageResult( + ToolStatus.OK, + f"Read {len(views)} ordered trajectory observations.", + views, + next_cursor, + ( + (ToolAction.READ_TRAJECTORY,) + if next_cursor is not None + else ( + ToolAction.SEARCH_PRIOR_TRAJECTORIES, + ToolAction.VERIFY_ENVIRONMENT, + ToolAction.ADAPT, + ToolAction.RETURN_VERDICT, + ) + ), + artifacts, + ) + + def verify_environment(self, request: EnvironmentCheckRequest) -> EnvironmentCheckResult: + selected = _find_environment_check(self.case, request) + if selected is None: + return EnvironmentCheckResult( + ToolStatus.BLOCKED, + "Environment check is not declared for this mining case.", + None, + (ToolAction.RETURN_VERDICT,), + (), + ) + source, outcome = selected + verification = self.environment.verify(request, source, outcome) + self._issued_evidence.extend(verification.evidence) + return EnvironmentCheckResult( + ToolStatus.UNAVAILABLE + if verification.status is CompletionStatus.UNKNOWN + else ToolStatus.OK, + "Environment state is unavailable." + if verification.status is CompletionStatus.UNKNOWN + else "Environment state verified.", + verification, + (ToolAction.ADAPT, ToolAction.RETURN_VERDICT), + verification.evidence, + ) + + def adapt(self, request: AdaptationRequest) -> AdaptationResult: + kinds = set(request.kinds) + signals = tuple( + signal for signal in self.production_signals if signal.kind in kinds + )[: request.limit] + return AdaptationResult( + ToolStatus.OK if signals else ToolStatus.NOT_FOUND, + f"Found {len(signals)} human or production calibration signals.", + signals, + ( + ToolAction.SEARCH_PRIOR_TRAJECTORIES + if signals + else ToolAction.RETURN_VERDICT, + ), + ) + + +@dataclass(frozen=True, slots=True) +class Mine: + revision: HarnessRevision + collection: CollectionResult + nominations: tuple[MiningNomination, ...] + judge: FailureJudge + environment: EnvironmentVerifier + + def __post_init__(self) -> None: + if not self.nominations: + raise ValueError("mine requires at least one nomination") + + def run(self) -> FailureMiningRun: + signals = tuple(source for item in self.nominations for source in item.sources) + results = tuple(self._mine(nomination, signals) for nomination in self.nominations) + return FailureMiningRun(self.revision.id, self.collection.snapshot_digest, results) + + def _mine( + self, + nomination: MiningNomination, + production_signals: tuple[FailureSource, ...], + ) -> FailureMiningResult: + trace = next( + (item for item in self.collection.traces if item.id == nomination.trace_id), + None, + ) + if self.collection.revision_id != self.revision.id: + return _invalid(nomination, MiningInvalidReason.REVISION_MISMATCH) + if trace is None: + return _invalid(nomination, MiningInvalidReason.TRACE_NOT_FOUND) + if not _trace_is_complete(self.collection, trace): + return _invalid(nomination, MiningInvalidReason.CORRUPT_TRACE) + context = MiningContext( + revision_id=self.revision.id, + trace_id=trace.id, + trace_digest=trace.digest, + observation_ids=trace.observation_ids, + session_id=trace.session_id, + environment_name=trace.environment, + release=trace.release, + available_tools=nomination.available_tools, + environment_sources=nomination.environment_sources, + initial_state_evidence=nomination.initial_state_evidence, + ) + case = TraceMiningCase(nomination.task, context, nomination.sources) + tools = MiningTools(case, self.collection, self.environment, production_signals) + result = self.judge.investigate(case, tools) + if not _judge_result_matches( + case, + result, + tools.issued_evidence, + tools.read_trajectory_evidence, + ): + return _invalid(nomination, MiningInvalidReason.JUDGE_OUTPUT) + return result + + +def _trace_is_complete(collection: CollectionResult, trace: TraceRecord) -> bool: + if trace.attribution is not AttributionLevel.EXACT or trace.gaps: + return False + store = CollectionStore(collection.store_path) + try: + observations = tuple( + item + for item in store.observations(collection.observation_sync_id) + if item.trace_id == trace.id + ) + for observation in observations: + _read_observation(store, collection, observation) + except CollectionError: + return False + finally: + store.close() + return bool(observations) and tuple(item.id for item in observations) == trace.observation_ids + + +def _judge_result_matches( + case: TraceMiningCase, + result: FailureMiningResult, + issued_evidence: tuple[EvidenceReference, ...], + read_trajectory_evidence: tuple[EvidenceReference, ...], +) -> bool: + if result.context is None: + return False + read_ids = {ObservationId(item.record_id.value) for item in read_trajectory_evidence} + behavior_evidence = ( + () + if result.failure_behavior is None + else tuple( + evidence + for observation in result.failure_behavior.observations + for evidence in observation.evidence + ) + ) + completion_evidence = tuple( + evidence for check in result.completion_checks for evidence in check.evidence + ) + return ( + result.task == case.task + and result.context == case.context + and result.source_ids == tuple(source.id for source in case.sources) + and read_ids == set(case.context.observation_ids) + and all(item.kind is EvidenceKind.TRAJECTORY for item in result.trajectory_evidence) + and all(item in result.trajectory_evidence for item in read_trajectory_evidence) + and all(item.kind is EvidenceKind.ENVIRONMENT for item in result.environment_evidence) + and all(item.kind is EvidenceKind.ENVIRONMENT for item in completion_evidence) + and all(item in issued_evidence for item in result.trajectory_evidence) + and all(item in issued_evidence for item in result.environment_evidence) + and all(item in issued_evidence for item in completion_evidence) + and all(item in issued_evidence for item in behavior_evidence) + ) + + +def _find_environment_check( + case: TraceMiningCase, + request: EnvironmentCheckRequest, +) -> tuple[EnvironmentSource, RequiredOutcome] | None: + source = next( + ( + item + for item in case.context.environment_sources + if item.id == request.source_id + ), + None, + ) + outcome = next( + ( + item + for item in case.task.required_outcomes + if item.source_id == request.source_id and item.check_id == request.check_id + ), + None, + ) + return None if source is None or outcome is None else (source, outcome) + + +def _page_start( + observations: tuple[ObservationRecord, ...], + cursor: ObservationId | None, +) -> int | None: + if cursor is None: + return 0 + for index, observation in enumerate(observations): + if observation.id == cursor: + return index + return None + + +def _read_observation( + store: CollectionStore, + collection: CollectionResult, + observation: ObservationRecord, +) -> TrajectoryObservation: + input_content = ( + None + if observation.input_content is None + else store.read_content(collection.observation_sync_id, observation.input_content) + ) + output_content = ( + None + if observation.output_content is None + else store.read_content(collection.observation_sync_id, observation.output_content) + ) + return TrajectoryObservation(observation, input_content, output_content) + + +def _search_tokens(text: str) -> tuple[str, ...]: + tokens: list[str] = [] + for word in text.split(): + token = "".join(character for character in word if character.isalnum() or character in "-_") + if len(token) >= 3 and token.casefold() not in {item.casefold() for item in tokens}: + tokens.append(token) + return tuple(tokens) + + +def _focus_hit( + store: CollectionStore, + collection: CollectionResult, + hit: ObservationContentHit, + query: str, +) -> ObservationContentHit: + content = store.read_content(collection.observation_sync_id, hit.reference).text + lowered = content.casefold() + positions = tuple( + lowered.find(candidate.casefold()) + for candidate in (query, *_search_tokens(query)) + ) + position = next((item for item in positions if item >= 0), 0) + start = max(0, position - 300) + return ObservationContentHit( + hit.observation_id, + hit.trace_id, + hit.field, + hit.reference, + content[start : start + 1000], + ) + + +def _invalid( + nomination: MiningNomination, + reason: MiningInvalidReason, +) -> FailureMiningResult: + return FailureMiningResult( + task=nomination.task, + context=None, + verdict=MiningVerdict.INVALID, + source_ids=tuple(source.id for source in nomination.sources), + completion_checks=(), + failure_behavior=None, + trajectory_evidence=(), + environment_evidence=(), + confidence=Confidence(1.0), + unresolved_questions=(), + invalid_reason=reason, + ) + + +def _require_identifier(value: str) -> None: + if not value or not value.isascii() or any(character.isspace() for character in value): + raise ValueError("identifier must be non-empty ASCII without whitespace") + + +def _require_text(value: str, name: str) -> None: + if not value.strip() or "\0" in value: + raise ValueError(f"{name} must be non-empty text") diff --git a/src/ofw/repository.py b/src/ofw/repository.py new file mode 100644 index 0000000..362da10 --- /dev/null +++ b/src/ofw/repository.py @@ -0,0 +1,181 @@ +"""Turn a complete git repository into an immutable harness revision.""" + +from __future__ import annotations + +import hashlib +import os +import re +import subprocess # nosec B404 +import tempfile +from pathlib import Path + +from ofw.contracts import ( + GitCommit, + HarnessErrorCode, + HarnessRevision, + HarnessRevisionContent, + HarnessRevisionId, + HarnessSchemaVersion, + HarnessValidationError, + RepositorySnapshot, + Sha256Digest, +) +from ofw.observability.langfuse.contracts import LangfuseProject + +_NAME_PATTERN = re.compile(r"[a-z0-9]+(?:-[a-z0-9]+)*") + + +def process_repository( + name: str, + root: Path, + *, + traces: LangfuseProject | None = None, +) -> HarnessRevision: + """Snapshot a whole agent-harness repository without component mapping.""" + if _NAME_PATTERN.fullmatch(name) is None: + raise HarnessValidationError(HarnessErrorCode.INVALID_NAME, name) + selected_root = _resolve_root(root) + content = HarnessRevisionContent( + schema_version=HarnessSchemaVersion.V1, + harness_name=name, + repository=_snapshot_repository(selected_root), + observability=None if traces is None else traces.manifest(), + ) + revision_id = _revision_id(content.repository) + revision = HarnessRevision( + schema_version=content.schema_version, + id=revision_id, + harness_name=content.harness_name, + root=selected_root, + repository=content.repository, + observability=content.observability, + ) + _write_manifest(revision) + return revision + + +def _revision_id(repository: RepositorySnapshot) -> HarnessRevisionId: + if repository.dirty_digest is None: + return HarnessRevisionId(str(repository.commit)) + return HarnessRevisionId( + f"{repository.commit.value}-dirty-{repository.dirty_digest.value[7:23]}" + ) + + +def _resolve_root(root: Path) -> Path: + if not isinstance(root, Path): + raise HarnessValidationError(HarnessErrorCode.INVALID_SOURCE, repr(root)) + try: + resolved = root.expanduser().resolve(strict=True) + except FileNotFoundError as error: + raise HarnessValidationError(HarnessErrorCode.ROOT_NOT_FOUND, str(root)) from error + if not resolved.is_dir(): + raise HarnessValidationError(HarnessErrorCode.ROOT_NOT_DIRECTORY, str(resolved)) + return resolved + + +def _snapshot_repository(root: Path) -> RepositorySnapshot: + top_level = _run_git(root, "rev-parse", "--show-toplevel", repository_probe=True) + try: + git_root = Path(top_level.decode().strip()).resolve(strict=True) + except (UnicodeDecodeError, FileNotFoundError) as error: + raise HarnessValidationError( + HarnessErrorCode.GIT_REPOSITORY_REQUIRED, + str(root), + ) from error + if git_root != root: + raise HarnessValidationError(HarnessErrorCode.GIT_REPOSITORY_REQUIRED, str(root)) + commit_bytes = _run_git(root, "rev-parse", "HEAD") + try: + commit = GitCommit(commit_bytes.decode().strip()) + except UnicodeDecodeError as error: + raise HarnessValidationError( + HarnessErrorCode.GIT_COMMAND_FAILED, + "rev-parse HEAD", + ) from error + dirty = _dirty_payload(root) + return RepositorySnapshot( + commit=commit, + is_dirty=bool(dirty), + dirty_digest=None if not dirty else _digest(dirty), + ) + + +def _dirty_payload(root: Path) -> bytes: + payload = bytearray(_run_git(root, "diff", "--binary", "--no-ext-diff", "HEAD", "--")) + untracked = _run_git(root, "ls-files", "--others", "--exclude-standard", "-z") + for encoded_path in sorted(item for item in untracked.split(b"\0") if item): + relative = Path(os.fsdecode(encoded_path)) + if _ignored_internal_path(relative): + continue + path = root / relative + content = os.fsencode(os.readlink(path)) if path.is_symlink() else path.read_bytes() + payload.extend(b"\0untracked\0") + payload.extend(encoded_path) + payload.extend(b"\0") + payload.extend(content) + return bytes(payload) + + +def _ignored_internal_path(path: Path) -> bool: + return bool(path.parts) and ( + path.parts[0] == ".ofw" + or any(part == ".env" or part.startswith(".env.") for part in path.parts) + ) + + +def _run_git( + root: Path, + *arguments: str, + repository_probe: bool = False, +) -> bytes: + try: + result: subprocess.CompletedProcess[bytes] = subprocess.run( # nosec B603 + ("git", "-C", str(root), *arguments), + check=False, + capture_output=True, + ) + except OSError as error: + raise HarnessValidationError( + HarnessErrorCode.GIT_COMMAND_FAILED, + arguments[0], + ) from error + if result.returncode != 0: + code = ( + HarnessErrorCode.GIT_REPOSITORY_REQUIRED + if repository_probe + else HarnessErrorCode.GIT_COMMAND_FAILED + ) + raise HarnessValidationError(code, arguments[0]) + return result.stdout + + +def _digest(value: bytes) -> Sha256Digest: + return Sha256Digest(f"sha256:{hashlib.sha256(value).hexdigest()}") + + +def _write_manifest(revision: HarnessRevision) -> None: + path = revision.manifest_path + payload = f"{revision.to_json()}\n" + try: + path.parent.mkdir(parents=True, exist_ok=True) + descriptor, temporary_name = tempfile.mkstemp( + dir=path.parent, + prefix=f".{path.stem}-", + suffix=".json", + text=True, + ) + temporary_path = Path(temporary_name) + try: + with os.fdopen(descriptor, "w", encoding="utf-8") as stream: + stream.write(payload) + stream.flush() + os.fsync(stream.fileno()) + temporary_path.replace(path) + finally: + temporary_path.unlink(missing_ok=True) + except OSError as error: + raise HarnessValidationError( + HarnessErrorCode.MANIFEST_WRITE_FAILED, + str(path), + ) from error diff --git a/src/ofw/runtime.py b/src/ofw/runtime.py index 4c12b06..9d93825 100644 --- a/src/ofw/runtime.py +++ b/src/ofw/runtime.py @@ -21,7 +21,7 @@ from pydantic import TypeAdapter -from ofw.contracts import HarnessRevision, RuntimeConfiguration, Sha256Digest +from ofw.contracts import HarnessRevision, Sha256Digest _NAME_PATTERN = re.compile(r"[a-z][a-z0-9_-]*") _ENVIRONMENT_PATTERN = re.compile(r"[A-Z_][A-Z0-9_]*") @@ -393,19 +393,6 @@ def to_json(self) -> str: _CANARY_ADAPTER: TypeAdapter[CanaryReport] = TypeAdapter(CanaryReport) -def runtime_configuration( - root: Path, - execution: ExecutionEnvironment, - lifecycle: LifecycleAdapter, - verifiers: tuple[VerifierAdapter, ...], -) -> RuntimeConfiguration: - return RuntimeConfiguration( - execution.fingerprint(root), - lifecycle.fingerprint(root), - tuple(verifier.fingerprint(root) for verifier in verifiers), - ) - - def run_canary( revision: HarnessRevision, case: CanaryCase, diff --git a/tests/test_harness.py b/tests/test_harness.py deleted file mode 100644 index bf66af0..0000000 --- a/tests/test_harness.py +++ /dev/null @@ -1,422 +0,0 @@ -"""Component-observable harness revision behavior.""" - -from __future__ import annotations - -import subprocess -import sys -from dataclasses import FrozenInstanceError -from pathlib import Path - -import pytest - -from ofw import ( - AssetAccess, - ComponentKind, - EditableFile, - Harness, - HarnessAsset, - HarnessComponent, - HarnessErrorCode, - HarnessRevision, - HarnessValidationError, - Subagent, - Tool, - WorkspaceFile, - ofw, -) - - -def _run_git(root: Path, *arguments: str) -> None: - subprocess.run( - ("git", "-C", str(root), *arguments), - check=True, - capture_output=True, - text=True, - ) - - -def _repository(tmp_path: Path) -> Path: - root = tmp_path / "fixtureco-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") - return root - - -def _configured_harness(root: Path) -> Harness: - harness = Harness("fixtureco-research-agent", root=root) - harness.connect_prompt(ofw.editable(Path("prompt.md"))) - return harness - - -def _required_component(revision: HarnessRevision, kind: ComponentKind) -> HarnessComponent: - component = revision.component(kind) - assert component is not None - return component - - -def test_process_creates_typed_immutable_revision_and_manifest(tmp_path: Path) -> None: - root = _repository(tmp_path) - harness = _configured_harness(root) - - assert not (root / ".ofw").exists() - revision = harness.process() - - assert isinstance(revision, HarnessRevision) - assert all(isinstance(component, HarnessComponent) for component in revision.components) - assert all(isinstance(asset, HarnessAsset) for asset in revision.assets) - assert all(isinstance(asset.source, WorkspaceFile) for asset in revision.assets) - assert revision.editable_files == (Path("prompt.md"),) - assert ( - revision.manifest_path == root / ".ofw" / "revisions" / str(revision.id) / "manifest.json" - ) - assert revision.manifest_path.read_text(encoding="utf-8") == f"{revision.to_json()}\n" - subprocess.run( - (sys.executable, "-m", "json.tool", str(revision.manifest_path)), - check=True, - capture_output=True, - text=True, - ) - with pytest.raises(FrozenInstanceError): - revision.harness_name = "changed" # type: ignore[misc] - - -def test_process_records_five_file_level_components_for_polyglot_agent(tmp_path: Path) -> None: - root = _repository(tmp_path) - files = ( - Path("tools/search.ts"), - Path("tools/worker.go"), - Path("skills/research/SKILL.md"), - Path("subagents/reviewer.yaml"), - Path("middleware/retry.ts"), - ) - for relative_path in files: - path = root / relative_path - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(f"fixture: {relative_path.as_posix()}\n", encoding="utf-8") - - harness = _configured_harness(root) - harness.connect_tools( - Tool(name="search", source=ofw.editable(Path("tools/search.ts"))), - Tool(name="worker", source=ofw.editable(Path("tools/worker.go"))), - ) - harness.connect_skills(ofw.editable(Path("skills/research/SKILL.md"))) - harness.connect_subagents( - Subagent( - name="reviewer", - source=ofw.editable(Path("subagents/reviewer.yaml")), - ) - ) - harness.connect_middleware(ofw.editable(Path("middleware/retry.ts"))) - - revision = harness.process() - - assert {component.kind for component in revision.components} == set(ComponentKind) - assert len(revision.components) == 5 - assert all(component.assets for component in revision.components) - assert all(str(component.digest).startswith("sha256:") for component in revision.components) - assert Path("tools/search.ts") in revision.editable_files - assert Path("tools/worker.go") in revision.editable_files - assert _required_component(revision, ComponentKind.SUBAGENT).assets[0].name == "reviewer" - - -def test_tool_object_preserves_name_and_source(tmp_path: Path) -> None: - root = _repository(tmp_path) - implementation = root / "search.ts" - implementation.write_text("export const search = () => [];\n", encoding="utf-8") - harness = _configured_harness(root) - harness.connect_tools(Tool(name="search", source=ofw.editable(Path("search.ts")))) - - revision = harness.process() - - tool_component = _required_component(revision, ComponentKind.TOOL) - assert tool_component.assets[0].name == "search" - assert tool_component.assets[0].source.relative_path == Path("search.ts") - - -def test_component_fingerprint_localizes_a_tool_change(tmp_path: Path) -> None: - root = _repository(tmp_path) - implementation = root / "search.ts" - implementation.write_text("export const search = () => 1;\n", encoding="utf-8") - first_harness = _configured_harness(root) - first_harness.connect_tools( - Tool(name="search", source=ofw.editable(Path("search.ts"))), - ) - first = first_harness.process() - - implementation.write_text("export const search = () => 2;\n", encoding="utf-8") - second_harness = _configured_harness(root) - second_harness.connect_tools( - Tool(name="search", source=ofw.editable(Path("search.ts"))), - ) - second = second_harness.process() - - assert ( - _required_component(first, ComponentKind.TOOL).digest - != _required_component(second, ComponentKind.TOOL).digest - ) - assert ( - _required_component(first, ComponentKind.PROMPT).digest - == _required_component(second, ComponentKind.PROMPT).digest - ) - - -def test_adding_middleware_does_not_change_prompt_component(tmp_path: Path) -> None: - root = _repository(tmp_path) - middleware = root / "middleware.ts" - middleware.write_text("export const beforeCall = () => {};\n", encoding="utf-8") - first = _configured_harness(root).process() - second_harness = _configured_harness(root) - second_harness.connect_middleware(ofw.editable(Path("middleware.ts"))) - second = second_harness.process() - - assert ( - _required_component(first, ComponentKind.PROMPT).digest - == _required_component(second, ComponentKind.PROMPT).digest - ) - assert second.component(ComponentKind.MIDDLEWARE) is not None - - -def test_file_cannot_be_owned_by_two_components(tmp_path: Path) -> None: - root = _repository(tmp_path) - harness = _configured_harness(root) - harness.connect_skills(ofw.editable(Path("prompt.md"))) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - - assert raised.value.code is HarnessErrorCode.COMPONENT_OVERLAP - - -@pytest.mark.parametrize("name", ("", "contains spaces", "UPPERCASE")) -def test_invalid_tool_name_fails(name: str) -> None: - with pytest.raises(HarnessValidationError) as raised: - Tool(name=name, source=Path("tool.py")) - assert raised.value.code is HarnessErrorCode.INVALID_TOOL_NAME - - -def test_duplicate_tool_name_fails(tmp_path: Path) -> None: - root = _repository(tmp_path) - first = root / "first.py" - second = root / "second.ts" - first.write_text("def run(): pass\n", encoding="utf-8") - second.write_text("export const run = () => {};\n", encoding="utf-8") - harness = _configured_harness(root) - - with pytest.raises(HarnessValidationError) as raised: - harness.connect_tools( - Tool(name="run", source=Path("first.py")), - Tool(name="run", source=Path("second.ts")), - ) - - assert raised.value.code is HarnessErrorCode.DUPLICATE_TOOL - - -def test_multiple_named_tools_may_share_one_source_file(tmp_path: Path) -> None: - root = _repository(tmp_path) - source = root / "tools.py" - source.write_text("def read(): pass\n\ndef write(): pass\n", encoding="utf-8") - harness = _configured_harness(root) - harness.connect_tools( - Tool(name="read", source=ofw.editable(Path("tools.py"))), - Tool(name="write", source=ofw.editable(Path("tools.py"))), - ) - - revision = harness.process() - - tool_component = _required_component(revision, ComponentKind.TOOL) - assert tuple(asset.name for asset in tool_component.assets) == ("read", "write") - - -@pytest.mark.parametrize("name", ("", "contains spaces", "UPPERCASE")) -def test_invalid_subagent_name_fails(name: str) -> None: - with pytest.raises(HarnessValidationError) as raised: - Subagent(name=name, source=Path("subagent.py")) - assert raised.value.code is HarnessErrorCode.INVALID_SUBAGENT_NAME - - -def test_duplicate_subagent_name_fails(tmp_path: Path) -> None: - root = _repository(tmp_path) - source = root / "subagents.py" - source.write_text("reviewer = 1\n", encoding="utf-8") - harness = _configured_harness(root) - - with pytest.raises(HarnessValidationError) as raised: - harness.connect_subagents( - Subagent(name="reviewer", source=Path("subagents.py")), - Subagent(name="reviewer", source=Path("subagents.py")), - ) - - assert raised.value.code is HarnessErrorCode.DUPLICATE_SUBAGENT - - -def test_multiple_named_subagents_may_share_one_source_file(tmp_path: Path) -> None: - root = _repository(tmp_path) - source = root / "subagents.py" - source.write_text("reviewer = 1\nresearcher = 2\n", encoding="utf-8") - harness = _configured_harness(root) - harness.connect_subagents( - Subagent(name="reviewer", source=ofw.editable(Path("subagents.py"))), - Subagent(name="researcher", source=ofw.editable(Path("subagents.py"))), - ) - - revision = harness.process() - - component = _required_component(revision, ComponentKind.SUBAGENT) - assert tuple(asset.name for asset in component.assets) == ("researcher", "reviewer") - - -def test_assets_are_frozen_unless_explicitly_editable(tmp_path: Path) -> None: - root = _repository(tmp_path) - skill = root / "SKILL.md" - skill.write_text("# Skill\n", encoding="utf-8") - harness = Harness("fixtureco-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - harness.connect_skills(ofw.editable(Path("SKILL.md"))) - - revision = harness.process() - - assert ( - _required_component(revision, ComponentKind.PROMPT).assets[0].access is AssetAccess.FROZEN - ) - assert ( - _required_component(revision, ComponentKind.SKILL).assets[0].access - is AssetAccess.FIT_EDITABLE - ) - - -def test_environment_secret_file_is_never_fingerprinted(tmp_path: Path) -> None: - root = _repository(tmp_path) - (root / ".env").write_text("SECRET=do-not-read\n", encoding="utf-8") - harness = _configured_harness(root) - harness.connect_prompt(Path(".env")) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - - assert raised.value.code is HarnessErrorCode.SENSITIVE_ASSET - - -def test_same_inputs_produce_same_revision(tmp_path: Path) -> None: - root = _repository(tmp_path) - assert _configured_harness(root).process() == _configured_harness(root).process() - - -def test_file_change_produces_new_revision(tmp_path: Path) -> None: - root = _repository(tmp_path) - first = _configured_harness(root).process() - (root / "prompt.md").write_text("Be accurate and concise.\n", encoding="utf-8") - second = _configured_harness(root).process() - - assert second.id != first.id - assert second.repository.is_dirty - assert second.repository.dirty_digest is not None - - -def test_new_git_commit_produces_new_revision(tmp_path: Path) -> None: - root = _repository(tmp_path) - first = _configured_harness(root).process() - (root / "README.md").write_text("Fixture repository.\n", encoding="utf-8") - _run_git(root, "add", "README.md") - _run_git(root, "commit", "-qm", "document fixture") - second = _configured_harness(root).process() - - assert second.id != first.id - assert second.repository.commit != first.repository.commit - assert not second.repository.is_dirty - - -@pytest.mark.parametrize("name", ("", "contains spaces", "UPPERCASE")) -def test_invalid_harness_name_fails(name: str, tmp_path: Path) -> None: - with pytest.raises(HarnessValidationError) as raised: - Harness(name, root=tmp_path) - assert raised.value.code is HarnessErrorCode.INVALID_NAME - - -@pytest.mark.parametrize( - ("source", "code"), - ( - (Path("missing.md"), HarnessErrorCode.MISSING_ASSET), - (Path("folder"), HarnessErrorCode.NOT_A_FILE), - ), -) -def test_invalid_workspace_file_fails( - source: Path, - code: HarnessErrorCode, - tmp_path: Path, -) -> None: - root = _repository(tmp_path) - (root / "folder").mkdir() - harness = _configured_harness(root) - harness.connect_skills(source) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - assert raised.value.code is code - - -def test_path_and_symlink_escape_fail(tmp_path: Path) -> None: - root = _repository(tmp_path) - outside = tmp_path / "outside.md" - outside.write_text("outside", encoding="utf-8") - (root / "linked.md").symlink_to(outside) - - for source in (outside, Path("linked.md")): - harness = _configured_harness(root) - harness.connect_skills(source) - with pytest.raises(HarnessValidationError) as raised: - harness.process() - assert raised.value.code is HarnessErrorCode.PATH_OUTSIDE_ROOT - - -@pytest.mark.parametrize( - ("sources", "code"), - ( - ((Path("prompt.md"), Path("prompt.md")), HarnessErrorCode.DUPLICATE_ASSET), - ( - (Path("prompt.md"), ofw.editable(Path("prompt.md"))), - HarnessErrorCode.CONFLICTING_ACCESS, - ), - ), -) -def test_duplicate_component_asset_fails( - sources: tuple[Path | EditableFile, ...], - code: HarnessErrorCode, - tmp_path: Path, -) -> None: - root = _repository(tmp_path) - harness = Harness("fixtureco-agent", root=root) - first, second = sources - harness.connect_prompt(first, second) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - assert raised.value.code is code - - -def test_process_requires_prompt(tmp_path: Path) -> None: - root = _repository(tmp_path) - skill = root / "SKILL.md" - skill.write_text("# Skill\n", encoding="utf-8") - harness = Harness("fixtureco-agent", root=root) - harness.connect_skills(Path("SKILL.md")) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - assert raised.value.code is HarnessErrorCode.PROMPT_REQUIRED - - -def test_root_must_be_git_repository(tmp_path: Path) -> None: - root = tmp_path / "not-git" - root.mkdir() - (root / "prompt.md").write_text("hello", encoding="utf-8") - harness = Harness("fixtureco-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - assert raised.value.code is HarnessErrorCode.GIT_REPOSITORY_REQUIRED diff --git a/tests/test_langfuse_collection.py b/tests/test_langfuse_collection.py index f33a111..9aee130 100644 --- a/tests/test_langfuse_collection.py +++ b/tests/test_langfuse_collection.py @@ -16,15 +16,14 @@ from ofw import ( CollectionError, CollectionErrorCode, - Harness, HarnessRevision, LangfuseProject, ObservationContentField, ObservationContentMatch, ObservationContentQuery, - Tool, TraceWindow, ofw, + process_repository, ) from ofw.observability.langfuse.domain import ( AttributionLevel, @@ -258,11 +257,7 @@ def _revision( base_url=server.base_url, allow_private_network=True, ) - harness = Harness("fixture-agent", root=root) - harness.connect_prompt(ofw.editable(Path("prompt.md"))) - harness.connect_tools(Tool(name="run", source=ofw.editable(Path("tool.py")))) - harness.connect_observability(project) - return harness.process() + return process_repository("fixture-agent", root, traces=project) def _window() -> TraceWindow: diff --git a/tests/test_langfuse_contracts.py b/tests/test_langfuse_contracts.py index b4d8379..f1ec1a4 100644 --- a/tests/test_langfuse_contracts.py +++ b/tests/test_langfuse_contracts.py @@ -12,10 +12,10 @@ from ofw import ( CollectionError, CollectionErrorCode, - Harness, LangfuseProject, TraceWindow, ofw, + process_repository, ) @@ -105,27 +105,21 @@ def test_trace_window_requires_aware_utc_ordering() -> None: assert reversed_window.value.code is CollectionErrorCode.INVALID_WINDOW -def test_observability_connection_changes_harness_revision_without_persisting_secrets( +def test_observability_connection_does_not_change_code_revision_or_persist_secrets( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: root = _harness_root(tmp_path) - baseline_harness = Harness("fixture-agent", root=root) - baseline_harness.connect_prompt(ofw.editable(Path("prompt.md"))) - baseline = baseline_harness.process() + baseline = process_repository("fixture-agent", root) monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-sensitive") monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-sensitive") project = LangfuseProject.from_env( environment="production", base_url="https://us.cloud.langfuse.com", ) - connected_harness = Harness("fixture-agent", root=root) - connected_harness.connect_prompt(ofw.editable(Path("prompt.md"))) - connected_harness.connect_observability(project) + connected = process_repository("fixture-agent", root, traces=project) - connected = connected_harness.process() - - assert connected.id != baseline.id + assert connected.id == baseline.id assert connected.observability is not None assert connected.observability == project.manifest() assert "pk-sensitive" not in connected.to_json() @@ -134,9 +128,7 @@ def test_observability_connection_changes_harness_revision_without_persisting_se def test_collect_requires_connected_observability(tmp_path: Path) -> None: root = _harness_root(tmp_path) - harness = Harness("fixture-agent", root=root) - harness.connect_prompt(ofw.editable(Path("prompt.md"))) - revision = harness.process() + revision = process_repository("fixture-agent", root) start = datetime(2026, 8, 22, tzinfo=UTC) with pytest.raises(CollectionError) as raised: diff --git a/tests/test_mine.py b/tests/test_mine.py new file mode 100644 index 0000000..aa40774 --- /dev/null +++ b/tests/test_mine.py @@ -0,0 +1,870 @@ +"""Failure mining over complete Langfuse trajectories and verified final state.""" + +from __future__ import annotations + +from dataclasses import dataclass, replace +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest + +from ofw.contracts import ( + GitCommit, + HarnessRevision, + HarnessRevisionId, + HarnessSchemaVersion, + RepositorySnapshot, + Sha256Digest, +) +from ofw.mine import ( + AdaptationRequest, + BehaviorObservation, + CompletionCheck, + CompletionStatus, + Confidence, + ConstraintKind, + EnvironmentCheckId, + EnvironmentCheckRequest, + EnvironmentSource, + EnvironmentSourceId, + EnvironmentSourceKind, + EnvironmentVerification, + EvidenceKind, + EvidenceRecordId, + EvidenceReference, + FailureBehavior, + FailureBehaviorKind, + FailureMiningResult, + FailurePhase, + FailureSource, + FailureSourceId, + FailureSourceKind, + Mine, + MiningContext, + MiningInvalidReason, + MiningNomination, + MiningTask, + MiningTools, + MiningVerdict, + RecoveryStatus, + RequiredOutcome, + TaskConstraint, + TaskId, + ToolAccess, + ToolCapability, + ToolName, + ToolStatus, + TraceMiningCase, + TrajectoryPageRequest, + TrajectorySearchRequest, +) +from ofw.observability.langfuse.contracts import LangfuseConnectionId, TraceWindow +from ofw.observability.langfuse.domain import ( + AttributionLevel, + CollectionCapabilityReason, + CollectionResult, + CollectionSyncId, + JsonDocument, + ObservationContent, + ObservationContentField, + ObservationContentReference, + ObservationId, + ObservationPage, + ObservationRecord, + ObservationType, + ProjectId, + ScorePage, + TraceGap, + TraceId, + TraceRecord, +) +from ofw.observability.langfuse.store import CollectionStore + +NOW = datetime(2026, 8, 25, tzinfo=UTC) +TRACE_ID = TraceId("trace-1") +CHECK_ID = EnvironmentCheckId("ticket-closed") +SOURCE_ID = EnvironmentSourceId("itsm-production") + + +def _digest(value: str) -> Sha256Digest: + return Sha256Digest(f"sha256:{value}") + + +def _revision(tmp_path: Path, value: str = "revision-1") -> HarnessRevision: + return HarnessRevision( + schema_version=HarnessSchemaVersion.V1, + id=HarnessRevisionId(value), + harness_name="support-agent", + root=tmp_path, + repository=RepositorySnapshot(GitCommit("abc123"), False, None), + observability=None, + ) + + +def _observation( + revision: HarnessRevision, + index: int, + name: str, + output: str, + *, + level: str | None = None, +) -> tuple[ObservationRecord, ObservationContent]: + reference = ObservationContentReference.for_text(output) + observation_id = ObservationId(f"observation-{index}") + return ( + ObservationRecord( + id=observation_id, + trace_id=TRACE_ID, + start_time=NOW + timedelta(seconds=index), + end_time=NOW + timedelta(seconds=index, milliseconds=100), + project_id=ProjectId("project-1"), + parent_observation_id=None if index == 0 else ObservationId("observation-0"), + type=ObservationType.AGENT if index in (0, 3) else ObservationType.TOOL, + is_root=index == 0, + name=name, + level=None, + version="v1", + environment="production", + user_id="user-1", + session_id="session-1", + created_at=NOW + timedelta(seconds=index), + updated_at=NOW + timedelta(seconds=index, milliseconds=100), + metadata=JsonDocument(f'{{"revision":"{revision.id}"}}'), + usage=None, + costs=None, + total_cost=None, + tags=(), + release=str(revision.id), + trace_name="close-ticket", + raw=JsonDocument(f'{{"name":"{name}","output":"{output}"}}'), + digest=_digest(f"observation-{index}"), + status_message=level, + output_content=reference, + ), + ObservationContent(reference, output), + ) + + +def _collection( + tmp_path: Path, + revision: HarnessRevision, + outputs: tuple[str, ...], + *, + attribution: AttributionLevel = AttributionLevel.EXACT, + gaps: tuple[TraceGap, ...] = (), +) -> CollectionResult: + named = tuple( + _observation( + revision, + index, + "agent" if index in (0, len(outputs) - 1) else "update-ticket", + output, + level="failed" if "failed" in output else None, + ) + for index, output in enumerate(outputs) + ) + observations = tuple(item[0] for item in named) + contents = tuple(item[1] for item in named) + observation_sync_id = CollectionSyncId("observations-mine") + score_sync_id = CollectionSyncId("scores-mine") + store_path = tmp_path / "collection.sqlite" + store = CollectionStore(store_path) + try: + store.commit_observation_page( + "connection-1", + observation_sync_id, + ObservationPage(observations, None, contents), + ) + store.commit_score_page("connection-1", score_sync_id, ScorePage((), None)) + finally: + store.close() + trace = TraceRecord( + id=TRACE_ID, + observation_ids=tuple(observation.id for observation in observations), + root_observation_ids=(observations[0].id,), + score_ids=(), + session_id="session-1", + environment="production", + release=str(revision.id), + attribution=attribution, + gaps=gaps, + digest=_digest("trace-1"), + ) + return CollectionResult( + revision_id=revision.id, + connection_id=LangfuseConnectionId("connection-1"), + window=TraceWindow(NOW - timedelta(minutes=1), NOW + timedelta(minutes=1)), + observation_sync_id=observation_sync_id, + score_sync_id=score_sync_id, + traces=(trace,), + observation_count=len(observations), + score_count=0, + gap_count=len(gaps), + snapshot_digest=_digest("collection"), + capability=( + CollectionCapabilityReason.READY + if attribution is AttributionLevel.EXACT and not gaps + else CollectionCapabilityReason.INCOMPLETE_TRACE + ), + store_path=store_path, + ) + + +def _evidence(kind: EvidenceKind, record_id: str, digest: str) -> EvidenceReference: + return EvidenceReference(kind, EvidenceRecordId(record_id), _digest(digest)) + + +def _nomination(kind: FailureSourceKind) -> MiningNomination: + source = FailureSource( + id=FailureSourceId("signal-1"), + kind=kind, + trace_id=TRACE_ID, + observed_at=NOW, + summary="The ticket may still be open.", + evidence=(_evidence(EvidenceKind.PRODUCTION_SIGNAL, "signal-1", "signal-1"),), + ) + environment = EnvironmentSource( + id=SOURCE_ID, + kind=EnvironmentSourceKind.PRODUCTION_API, + summary="Read-only ITSM production state.", + ) + return MiningNomination( + trace_id=TRACE_ID, + task=_task(), + sources=(source,), + environment_sources=(environment,), + available_tools=( + ToolCapability(ToolName("update-ticket"), ToolAccess.MUTATING), + ), + ) + + +def _task() -> MiningTask: + return MiningTask( + id=TaskId("close-ticket"), + intent="Close the customer ticket.", + required_outcomes=( + RequiredOutcome( + check_id=CHECK_ID, + source_id=SOURCE_ID, + description="The ticket is closed.", + ), + ), + constraints=( + TaskConstraint( + kind=ConstraintKind.POLICY, + description="Do not claim completion without verifying the ticket state.", + ), + ), + ) + + +def _context(observation_ids: tuple[ObservationId, ...]) -> MiningContext: + return MiningContext( + revision_id=HarnessRevisionId("revision-1"), + trace_id=TRACE_ID, + trace_digest=_digest("trace-1"), + observation_ids=observation_ids, + session_id="session-1", + environment_name="production", + release="revision-1", + available_tools=( + ToolCapability(ToolName("search_trajectory"), ToolAccess.READ_ONLY), + ToolCapability(ToolName("read_trajectory"), ToolAccess.READ_ONLY), + ToolCapability(ToolName("verify_environment"), ToolAccess.READ_ONLY), + ), + environment_sources=( + EnvironmentSource( + id=SOURCE_ID, + kind=EnvironmentSourceKind.PRODUCTION_API, + summary="Read-only ITSM production state.", + ), + ), + initial_state_evidence=(), + ) + + +@dataclass(frozen=True, slots=True) +class RecordedEnvironmentVerifier: + status: CompletionStatus + observed_state: str | None + + def verify( + self, + request: EnvironmentCheckRequest, + source: EnvironmentSource, + outcome: RequiredOutcome, + ) -> EnvironmentVerification: + assert request.source_id == source.id + assert request.check_id == outcome.check_id + evidence = ( + () + if self.status is CompletionStatus.UNKNOWN + else (_evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"),) + ) + return EnvironmentVerification( + status=self.status, + observed_state=self.observed_state, + evidence=evidence, + ) + + +@dataclass(frozen=True, slots=True) +class FakeFailureJudge: + search_text: str + + def investigate( + self, + case: TraceMiningCase, + tools: MiningTools, + ) -> FailureMiningResult: + assert case == tools.case + search = tools.search_trajectory( + TrajectorySearchRequest( + text=self.search_text, + field=ObservationContentField.ANY, + limit=10, + ) + ) + assert search.status is ToolStatus.OK + focused = tools.read_trajectory( + TrajectoryPageRequest(cursor=search.hits[0].observation_id, limit=1) + ) + assert focused.observations[0].record.id == search.hits[0].observation_id + + trajectory_evidence: tuple[EvidenceReference, ...] = () + cursor: ObservationId | None = None + while True: + page = tools.read_trajectory(TrajectoryPageRequest(cursor=cursor, limit=2)) + assert page.status is ToolStatus.OK + trajectory_evidence = (*trajectory_evidence, *page.artifacts) + if page.next_cursor is None: + break + cursor = page.next_cursor + + verification = tools.verify_environment( + EnvironmentCheckRequest(SOURCE_ID, CHECK_ID) + ) + unresolved: tuple[str, ...] + if verification.status is ToolStatus.UNAVAILABLE: + completion = CompletionStatus.UNKNOWN + verdict = MiningVerdict.AMBIGUOUS + unresolved = ("Production state was unavailable.",) + environment_evidence: tuple[EvidenceReference, ...] = () + observed_state = None + else: + assert verification.verification is not None + completion = verification.verification.status + verdict = ( + MiningVerdict.CONFIRMED_FAILURE + if completion is CompletionStatus.NOT_COMPLETED + else MiningVerdict.NO_FAILURE + ) + unresolved = () + environment_evidence = verification.artifacts + observed_state = verification.verification.observed_state + + adapted = tools.adapt( + AdaptationRequest( + ( + FailureSourceKind.HUMAN_FEEDBACK, + FailureSourceKind.USER_CORRECTION, + FailureSourceKind.DOWNSTREAM_FAILURE, + FailureSourceKind.AGENT_ERROR, + ), + 10, + ) + ) + assert adapted.status is ToolStatus.OK + + behavior_evidence = next( + item + for item in trajectory_evidence + if item.record_id.value == search.hits[0].observation_id.value + ) + + failure_behavior = ( + None + if verdict is not MiningVerdict.CONFIRMED_FAILURE + else FailureBehavior( + primary=( + FailureBehaviorKind.FALSE_COMPLETION + if "successfully" in search.hits[0].excerpt.lower() + else FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE + ), + summary=( + "The agent claimed success while the ticket stayed open." + if "successfully" in search.hits[0].excerpt.lower() + else "The tool failure was not recovered before completion." + ), + observations=( + BehaviorObservation( + kind=( + FailureBehaviorKind.FALSE_COMPLETION + if "successfully" in search.hits[0].excerpt.lower() + else FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE + ), + phase=( + FailurePhase.COMPLETION + if "successfully" in search.hits[0].excerpt.lower() + else FailurePhase.RECOVERY + ), + first_observation_id=search.hits[0].observation_id, + last_observation_id=None, + recovery_status=RecoveryStatus.NOT_RECOVERED, + evidence=(behavior_evidence,), + ), + ), + ) + ) + + return FailureMiningResult( + task=tools.case.task, + context=tools.case.context, + verdict=verdict, + source_ids=tuple(source.id for source in tools.case.sources), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Ticket closed successfully.", + observed_state=observed_state, + status=completion, + evidence=environment_evidence, + ), + ), + failure_behavior=failure_behavior, + trajectory_evidence=trajectory_evidence, + environment_evidence=environment_evidence, + confidence=Confidence(0.95 if completion is not CompletionStatus.UNKNOWN else 0.4), + unresolved_questions=unresolved, + invalid_reason=None, + ) + + +@dataclass(frozen=True, slots=True) +class ForgingFailureJudge: + delegate: FakeFailureJudge + + def investigate( + self, + case: TraceMiningCase, + tools: MiningTools, + ) -> FailureMiningResult: + result = self.delegate.investigate(case, tools) + forged = _evidence(EvidenceKind.ENVIRONMENT, "invented-state", "invented") + check = replace(result.completion_checks[0], evidence=(forged,)) + return replace( + result, + completion_checks=(check,), + environment_evidence=(forged,), + ) + + +@dataclass(frozen=True, slots=True) +class PartialTraceFailureJudge: + def investigate( + self, + case: TraceMiningCase, + tools: MiningTools, + ) -> FailureMiningResult: + page = tools.read_trajectory(TrajectoryPageRequest(cursor=None, limit=1)) + verification = tools.verify_environment( + EnvironmentCheckRequest(SOURCE_ID, CHECK_ID) + ) + assert verification.verification is not None + return FailureMiningResult( + task=case.task, + context=case.context, + verdict=MiningVerdict.CONFIRMED_FAILURE, + source_ids=tuple(source.id for source in case.sources), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Ticket closed successfully.", + observed_state=verification.verification.observed_state, + status=CompletionStatus.NOT_COMPLETED, + evidence=verification.artifacts, + ), + ), + failure_behavior=FailureBehavior( + primary=FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE, + summary="The tool failure was not recovered before completion.", + observations=( + BehaviorObservation( + kind=FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE, + phase=FailurePhase.RECOVERY, + first_observation_id=page.observations[0].record.id, + last_observation_id=None, + recovery_status=RecoveryStatus.NOT_RECOVERED, + evidence=page.artifacts, + ), + ), + ), + trajectory_evidence=page.artifacts, + environment_evidence=verification.artifacts, + confidence=Confidence(0.9), + unresolved_questions=(), + invalid_reason=None, + ) + + +@pytest.mark.parametrize( + ("outputs", "source_kind", "state", "observed", "search", "expected"), + ( + ( + ("Close ticket", "update failed", "continuing", "Ticket closed successfully"), + FailureSourceKind.DOWNSTREAM_FAILURE, + CompletionStatus.NOT_COMPLETED, + "Ticket remains open", + "failed", + MiningVerdict.CONFIRMED_FAILURE, + ), + ( + ("Close ticket", "update failed", "retry succeeded", "Ticket closed successfully"), + FailureSourceKind.AGENT_ERROR, + CompletionStatus.COMPLETED, + "Ticket is closed", + "failed", + MiningVerdict.NO_FAILURE, + ), + ( + ("Close ticket", "update accepted", "Ticket closed successfully"), + FailureSourceKind.USER_CORRECTION, + CompletionStatus.NOT_COMPLETED, + "Ticket remains open", + "successfully", + MiningVerdict.CONFIRMED_FAILURE, + ), + ( + ("Close ticket", "update accepted", "Ticket closed successfully"), + FailureSourceKind.HUMAN_FEEDBACK, + CompletionStatus.UNKNOWN, + None, + "successfully", + MiningVerdict.AMBIGUOUS, + ), + ), +) +def test_mine_uses_complete_trajectory_and_verified_state( + tmp_path: Path, + outputs: tuple[str, ...], + source_kind: FailureSourceKind, + state: CompletionStatus, + observed: str | None, + search: str, + expected: MiningVerdict, +) -> None: + revision = _revision(tmp_path) + run = Mine( + revision=revision, + collection=_collection(tmp_path, revision, outputs), + nominations=(_nomination(source_kind),), + judge=FakeFailureJudge(search), + environment=RecordedEnvironmentVerifier(state, observed), + ).run() + + assert run.results[0].verdict is expected + assert len(run.results[0].trajectory_evidence) == len(outputs) + + +@pytest.mark.parametrize( + ("collection_revision", "attribution", "gaps", "reason"), + ( + ("other-revision", AttributionLevel.EXACT, (), MiningInvalidReason.REVISION_MISMATCH), + ( + "revision-1", + AttributionLevel.MISSING, + (TraceGap.MISSING_ROOT,), + MiningInvalidReason.CORRUPT_TRACE, + ), + ), +) +def test_wrong_revision_or_corrupt_trace_is_invalid( + tmp_path: Path, + collection_revision: str, + attribution: AttributionLevel, + gaps: tuple[TraceGap, ...], + reason: MiningInvalidReason, +) -> None: + revision = _revision(tmp_path) + foreign_revision = _revision(tmp_path, collection_revision) + collection = _collection( + tmp_path, + foreign_revision, + ("Close ticket", "Ticket closed successfully"), + attribution=attribution, + gaps=gaps, + ) + + result = Mine( + revision=revision, + collection=collection, + nominations=(_nomination(FailureSourceKind.AGENT_ERROR),), + judge=FakeFailureJudge("successfully"), + environment=RecordedEnvironmentVerifier(CompletionStatus.COMPLETED, "closed"), + ).run().results[0] + + assert result.verdict is MiningVerdict.INVALID + assert result.invalid_reason is reason + + +def test_confirmed_failure_requires_trajectory_and_environment_evidence() -> None: + with pytest.raises(ValueError, match="confirmed failure requires"): + FailureMiningResult( + task=_task(), + context=_context((ObservationId("observation-1"),)), + verdict=MiningVerdict.CONFIRMED_FAILURE, + source_ids=(FailureSourceId("signal-1"),), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Ticket closed successfully.", + observed_state="Ticket remains open.", + status=CompletionStatus.NOT_COMPLETED, + evidence=(), + ), + ), + failure_behavior=FailureBehavior( + primary=FailureBehaviorKind.FALSE_COMPLETION, + summary="The agent claimed success while the ticket stayed open.", + observations=( + BehaviorObservation( + kind=FailureBehaviorKind.FALSE_COMPLETION, + phase=FailurePhase.COMPLETION, + first_observation_id=ObservationId("observation-1"), + last_observation_id=None, + recovery_status=RecoveryStatus.NOT_RECOVERED, + evidence=( + _evidence( + EvidenceKind.TRAJECTORY, + "observation-1", + "observation-1", + ), + ), + ), + ), + ), + trajectory_evidence=( + _evidence(EvidenceKind.TRAJECTORY, "observation-1", "observation-1"), + ), + environment_evidence=(), + confidence=Confidence(0.9), + unresolved_questions=(), + invalid_reason=None, + ) + + +def test_judge_cannot_invent_evidence_that_no_tool_returned(tmp_path: Path) -> None: + revision = _revision(tmp_path) + result = Mine( + revision=revision, + collection=_collection( + tmp_path, + revision, + ("Close ticket", "update failed", "Ticket closed successfully"), + ), + nominations=(_nomination(FailureSourceKind.DOWNSTREAM_FAILURE),), + judge=ForgingFailureJudge(FakeFailureJudge("failed")), + environment=RecordedEnvironmentVerifier( + CompletionStatus.NOT_COMPLETED, + "Ticket remains open", + ), + ).run().results[0] + + assert result.verdict is MiningVerdict.INVALID + assert result.invalid_reason is MiningInvalidReason.JUDGE_OUTPUT + + +def test_judge_must_read_the_full_trace_before_returning_verdict(tmp_path: Path) -> None: + revision = _revision(tmp_path) + result = Mine( + revision=revision, + collection=_collection( + tmp_path, + revision, + ("Close ticket", "update failed", "Ticket closed successfully"), + ), + nominations=(_nomination(FailureSourceKind.DOWNSTREAM_FAILURE),), + judge=PartialTraceFailureJudge(), + environment=RecordedEnvironmentVerifier( + CompletionStatus.NOT_COMPLETED, + "Ticket remains open", + ), + ).run().results[0] + + assert result.verdict is MiningVerdict.INVALID + assert result.invalid_reason is MiningInvalidReason.JUDGE_OUTPUT + + +def test_confirmed_failure_requires_failure_behavior() -> None: + with pytest.raises(ValueError, match="confirmed failure requires"): + FailureMiningResult( + task=_task(), + context=_context((ObservationId("observation-1"),)), + verdict=MiningVerdict.CONFIRMED_FAILURE, + source_ids=(FailureSourceId("signal-1"),), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Ticket closed successfully.", + observed_state="Ticket remains open.", + status=CompletionStatus.NOT_COMPLETED, + evidence=(_evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"),), + ), + ), + failure_behavior=None, + trajectory_evidence=( + _evidence(EvidenceKind.TRAJECTORY, "observation-1", "observation-1"), + ), + environment_evidence=( + _evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"), + ), + confidence=Confidence(0.9), + unresolved_questions=(), + invalid_reason=None, + ) + + +def test_no_failure_rejects_failure_behavior() -> None: + with pytest.raises(ValueError, match="no-failure verdict"): + FailureMiningResult( + task=_task(), + context=_context((ObservationId("observation-1"),)), + verdict=MiningVerdict.NO_FAILURE, + source_ids=(FailureSourceId("signal-1"),), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Retry succeeded.", + observed_state="Ticket is closed.", + status=CompletionStatus.COMPLETED, + evidence=(_evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"),), + ), + ), + failure_behavior=FailureBehavior( + primary=FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE, + summary="A failure was incorrectly retained after recovery.", + observations=( + BehaviorObservation( + kind=FailureBehaviorKind.UNRECOVERED_ACTION_FAILURE, + phase=FailurePhase.RECOVERY, + first_observation_id=ObservationId("observation-1"), + last_observation_id=None, + recovery_status=RecoveryStatus.RECOVERED, + evidence=( + _evidence( + EvidenceKind.TRAJECTORY, + "observation-1", + "observation-1", + ), + ), + ), + ), + ), + trajectory_evidence=( + _evidence(EvidenceKind.TRAJECTORY, "observation-1", "observation-1"), + ), + environment_evidence=( + _evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"), + ), + confidence=Confidence(0.9), + unresolved_questions=(), + invalid_reason=None, + ) + + +def test_failure_behavior_observation_must_belong_to_context() -> None: + with pytest.raises(ValueError, match="failure behavior observation"): + FailureMiningResult( + task=_task(), + context=_context((ObservationId("observation-1"),)), + verdict=MiningVerdict.CONFIRMED_FAILURE, + source_ids=(FailureSourceId("signal-1"),), + completion_checks=( + CompletionCheck( + check_id=CHECK_ID, + required_outcome="The ticket is closed.", + agent_claim="Ticket closed successfully.", + observed_state="Ticket remains open.", + status=CompletionStatus.NOT_COMPLETED, + evidence=(_evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"),), + ), + ), + failure_behavior=FailureBehavior( + primary=FailureBehaviorKind.FALSE_COMPLETION, + summary="The agent claimed success while the ticket stayed open.", + observations=( + BehaviorObservation( + kind=FailureBehaviorKind.FALSE_COMPLETION, + phase=FailurePhase.COMPLETION, + first_observation_id=ObservationId("observation-2"), + last_observation_id=None, + recovery_status=RecoveryStatus.NOT_RECOVERED, + evidence=( + _evidence( + EvidenceKind.TRAJECTORY, + "observation-2", + "observation-2", + ), + ), + ), + ), + ), + trajectory_evidence=( + _evidence(EvidenceKind.TRAJECTORY, "observation-1", "observation-1"), + ), + environment_evidence=( + _evidence(EvidenceKind.ENVIRONMENT, "ticket-123", "state-1"), + ), + confidence=Confidence(0.9), + unresolved_questions=(), + invalid_reason=None, + ) + + +def test_search_focuses_natural_query_and_adapt_compares_signals(tmp_path: Path) -> None: + revision = _revision(tmp_path) + collection = _collection( + tmp_path, + revision, + ("x" * 1500 + " update failed after retry", "Ticket remains open"), + ) + nomination = _nomination(FailureSourceKind.DOWNSTREAM_FAILURE) + prior_signal = replace( + nomination.sources[0], + id=FailureSourceId("signal-2"), + trace_id=TraceId("trace-2"), + kind=FailureSourceKind.HUMAN_FEEDBACK, + ) + tools = MiningTools( + TraceMiningCase( + nomination.task, + _context(collection.traces[0].observation_ids), + nomination.sources, + ), + collection, + RecordedEnvironmentVerifier(CompletionStatus.UNKNOWN, None), + (*nomination.sources, prior_signal), + ) + + search = tools.search_trajectory( + TrajectorySearchRequest("completion failed retries", ObservationContentField.ANY, 5) + ) + adapted = tools.adapt( + AdaptationRequest( + (FailureSourceKind.DOWNSTREAM_FAILURE, FailureSourceKind.HUMAN_FEEDBACK), + 5, + ) + ) + + assert search.status is ToolStatus.OK + assert "update failed" in search.hits[0].excerpt + assert tuple(signal.id for signal in adapted.signals) == ( + FailureSourceId("signal-1"), + FailureSourceId("signal-2"), + ) diff --git a/tests/test_repository.py b/tests/test_repository.py new file mode 100644 index 0000000..e97bff2 --- /dev/null +++ b/tests/test_repository.py @@ -0,0 +1,103 @@ +"""Whole-repository harness revision behavior.""" + +from __future__ import annotations + +import subprocess +from dataclasses import FrozenInstanceError +from pathlib import Path + +import pytest + +from ofw import ( + HarnessErrorCode, + HarnessRevision, + HarnessValidationError, + LangfuseProject, + process_repository, +) + + +def _run_git(root: Path, *arguments: str) -> None: + subprocess.run( + ("git", "-C", str(root), *arguments), + check=True, + capture_output=True, + text=True, + ) + + +def _repository(tmp_path: Path) -> Path: + root = tmp_path / "fixture-agent" + root.mkdir() + (root / "agent.py").write_text("PROMPT = '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") + return root + + +def test_process_repository_creates_an_immutable_revision(tmp_path: Path) -> None: + root = _repository(tmp_path) + + revision = process_repository("fixture-agent", root) + + assert isinstance(revision, HarnessRevision) + assert str(revision.id) == str(revision.repository.commit) + assert revision.root == root + assert revision.manifest_path.is_file() + assert revision.manifest_path.read_text(encoding="utf-8") == f"{revision.to_json()}\n" + with pytest.raises(FrozenInstanceError): + revision.harness_name = "changed" # type: ignore[misc] + + +def test_repository_commit_dirty_diff_and_untracked_files_change_revision( + tmp_path: Path, +) -> None: + root = _repository(tmp_path) + clean = process_repository("fixture-agent", root) + (root / "agent.py").write_text("PROMPT = 'be concise'\n", encoding="utf-8") + dirty = process_repository("fixture-agent", root) + (root / "new_skill.md").write_text("# New skill\n", encoding="utf-8") + untracked = process_repository("fixture-agent", root) + + assert clean.id != dirty.id != untracked.id + assert not clean.repository.is_dirty + assert dirty.repository.is_dirty + assert untracked.repository.is_dirty + + +def test_observability_connection_changes_revision_without_storing_secrets( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + root = _repository(tmp_path) + baseline = process_repository("fixture-agent", root) + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-sensitive") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-sensitive") + project = LangfuseProject.from_env(environment="production") + + connected = process_repository("fixture-agent", root, traces=project) + + assert connected.id == baseline.id + assert connected.observability == project.manifest() + assert "pk-sensitive" not in connected.to_json() + assert "sk-sensitive" not in connected.to_json() + + +@pytest.mark.parametrize("name", ("", "contains spaces", "UPPERCASE")) +def test_invalid_repository_name_fails(name: str, tmp_path: Path) -> None: + with pytest.raises(HarnessValidationError) as raised: + process_repository(name, tmp_path) + assert raised.value.code is HarnessErrorCode.INVALID_NAME + + +def test_root_must_be_a_git_repository(tmp_path: Path) -> None: + root = tmp_path / "not-git" + root.mkdir() + + with pytest.raises(HarnessValidationError) as raised: + process_repository("fixture-agent", root) + + assert raised.value.code is HarnessErrorCode.GIT_REPOSITORY_REQUIRED diff --git a/tests/test_runtime.py b/tests/test_runtime.py index ac4e6e6..5146f5b 100644 --- a/tests/test_runtime.py +++ b/tests/test_runtime.py @@ -20,10 +20,7 @@ CommandLoop, CommandVerifier, E2BSandbox, - Harness, - HarnessErrorCode, HarnessRevision, - HarnessValidationError, ModelFingerprint, ProcessCommand, ProcessLimits, @@ -31,6 +28,7 @@ RunResult, RunStatus, VerifierVerdict, + process_repository, ) @@ -192,9 +190,7 @@ def _repository(tmp_path: Path) -> Path: def _revision(root: Path) -> HarnessRevision: - harness = Harness("runtime-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - return harness.process() + return process_repository("runtime-agent", root) def _verifier(mode: str, name: str = "verifier") -> CommandVerifier: @@ -208,77 +204,36 @@ def _command_loop() -> CommandLoop: ) -def _runtime_harness(root: Path) -> Harness: - harness = Harness("runtime-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - harness.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=2)))) - harness.connect_lifecycle(_command_loop()) - harness.connect_verifiers(_verifier("uppercase", "uppercase")) - return harness - - -def test_process_runs_e2b_command_canary_and_records_frozen_evidence(tmp_path: Path) -> None: - root = _repository(tmp_path) - harness = Harness("runtime-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - harness.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=2)))) - harness.connect_lifecycle(_command_loop()) - harness.connect_verifiers(_verifier("uppercase", "uppercase")) - - revision = harness.process(canary=CanaryCase(CaseId("smoke"), "ship")) - - assert revision.runtime is not None - assert revision.canary_digest is not None - assert revision.canary_path.is_file() - assert "pass" in revision.canary_path.read_text(encoding="utf-8") - - -def test_canary_evidence_does_not_change_runtime_revision_identity(tmp_path: Path) -> None: - root = _repository(tmp_path) - - without_canary = _runtime_harness(root).process() - with_canary = _runtime_harness(root).process(canary=CanaryCase(CaseId("identity"), "ship")) - - assert with_canary.id == without_canary.id - assert with_canary.canary_digest is not None - - -def test_partial_runtime_configuration_is_rejected(tmp_path: Path) -> None: +def test_run_canary_returns_frozen_evidence(tmp_path: Path) -> None: root = _repository(tmp_path) - harness = Harness("runtime-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - harness.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=1)))) - - with pytest.raises(HarnessValidationError) as raised: - harness.process() - - assert raised.value.code is HarnessErrorCode.RUNTIME_INCOMPLETE - - -def test_duplicate_verifier_name_is_rejected(tmp_path: Path) -> None: - harness = Harness("runtime-agent", root=_repository(tmp_path)) + revision = _revision(root) - with pytest.raises(HarnessValidationError) as raised: - harness.connect_verifiers( - _verifier("uppercase", "duplicate"), - _verifier("reject", "duplicate"), - ) + report = runtime_module.run_canary( + revision, + CanaryCase(CaseId("smoke"), "ship"), + E2BSandbox(ProcessLimits(timedelta(seconds=2))), + _command_loop(), + (_verifier("uppercase", "uppercase"),), + ) - assert raised.value.code is HarnessErrorCode.DUPLICATE_VERIFIER + assert report.passed + assert str(report.digest).startswith("sha256:") -def test_failed_canary_blocks_revision_creation(tmp_path: Path) -> None: +def test_failed_canary_is_reported_without_mutating_revision(tmp_path: Path) -> None: root = _repository(tmp_path) - harness = Harness("runtime-agent", root=root) - harness.connect_prompt(Path("prompt.md")) - harness.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=2)))) - harness.connect_lifecycle(_command_loop()) - harness.connect_verifiers(_verifier("reject", "reject")) + revision = _revision(root) - with pytest.raises(HarnessValidationError) as raised: - harness.process(canary=CanaryCase(CaseId("rejected"), "ship")) + report = runtime_module.run_canary( + revision, + CanaryCase(CaseId("rejected"), "ship"), + E2BSandbox(ProcessLimits(timedelta(seconds=2))), + _command_loop(), + (_verifier("reject", "reject"),), + ) - assert raised.value.code is HarnessErrorCode.CANARY_FAILED + assert not report.passed + assert _revision(root).id == revision.id def test_e2b_reports_timeout_and_nonzero_exit(tmp_path: Path) -> None: @@ -487,26 +442,6 @@ def test_parallel_e2b_environments_are_distinct(tmp_path: Path) -> None: environment.destroy(second) -def test_runtime_fingerprint_changes_without_changing_assets(tmp_path: Path) -> None: - root = _repository(tmp_path) - first = Harness("runtime-agent", root=root) - first.connect_prompt(Path("prompt.md")) - first.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=1)))) - first.connect_lifecycle(_command_loop()) - first.connect_verifiers(_verifier("uppercase", "uppercase")) - second = Harness("runtime-agent", root=root) - second.connect_prompt(Path("prompt.md")) - second.connect_execute(E2BSandbox(ProcessLimits(timedelta(seconds=2)))) - second.connect_lifecycle(_command_loop()) - second.connect_verifiers(_verifier("uppercase", "uppercase")) - - first_revision = first.process() - second_revision = second.process() - - assert first_revision.id != second_revision.id - assert first_revision.components == second_revision.components - - def test_command_verifier_is_terminated_at_environment_timeout(tmp_path: Path) -> None: root = _repository(tmp_path) revision = _revision(root) diff --git a/tests/test_typing.py b/tests/test_typing.py index 45c54da..015c964 100644 --- a/tests/test_typing.py +++ b/tests/test_typing.py @@ -11,6 +11,6 @@ def test_package_declares_inline_types_and_namespace_methods() -> None: assert package_file is not None assert Path(package_file).with_name("py.typed").is_file() assert callable(ofw.collect) - assert callable(ofw.editable) + assert callable(package.process_repository) assert callable(ofw.E2BSandbox) assert callable(ofw.ProcessLimits)