diff --git a/.github/workflows/ucf-foundation.yml b/.github/workflows/ucf-foundation.yml index 1fcec51..5f2c96e 100644 --- a/.github/workflows/ucf-foundation.yml +++ b/.github/workflows/ucf-foundation.yml @@ -3,33 +3,33 @@ name: UCF foundation on: pull_request: paths: - - "foundry/adapters/**" - - "foundry/contracts/transition_models.py" - - "foundry/environments/**" - - "foundry/providers/base.py" - - "foundry/storage/artifact_store.py" - - "foundry/storage/transition_journal.py" - - "foundry/verification/policy.py" - - "foundry/runtime/**" - - "foundry/orchestration/agent_runner.py" - - "foundry/orchestration/run_engine.py" - - "foundry/contracts/task_types.py" - - "tests/unit/runtime/**" - - "tests/unit/test_transition_models.py" - - "tests/unit/orchestration/test_agent_runner.py" + - "foundry/**" + - "tests/**" + - "pyproject.toml" - ".github/workflows/ucf-foundation.yml" push: + branches: [main] paths: - - "foundry/contracts/transition_models.py" - - "foundry/environments/**" - - "foundry/providers/base.py" - - "foundry/runtime/**" + - "foundry/**" + - "tests/**" + - "pyproject.toml" + - ".github/workflows/ucf-foundation.yml" + +permissions: + contents: read + +concurrency: + group: ucf-foundation-${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: true jobs: foundation: runs-on: ubuntu-latest + timeout-minutes: 15 steps: - uses: actions/checkout@v4 + with: + persist-credentials: false - uses: actions/setup-python@v5 with: python-version: "3.12" @@ -39,10 +39,16 @@ jobs: - name: Compile run: python -m compileall -q foundry tests - name: Lint changed foundation - run: | - ruff check foundry/adapters foundry/contracts/transition_models.py foundry/contracts/task_types.py foundry/environments foundry/providers/base.py foundry/runtime foundry/storage/artifact_store.py foundry/storage/transition_journal.py foundry/verification/policy.py foundry/orchestration/agent_runner.py foundry/orchestration/run_engine.py tests/unit/test_transition_models.py tests/unit/runtime/test_transition_engine.py tests/unit/runtime/test_foundry_transition_adapter.py tests/unit/orchestration/test_agent_runner.py + run: >- + ruff check foundry/adapters foundry/contracts/transition_models.py + foundry/contracts/task_types.py foundry/environments foundry/providers/base.py + foundry/runtime foundry/storage/artifact_store.py foundry/storage/transition_journal.py + foundry/verification/policy.py foundry/orchestration/agent_runner.py + foundry/orchestration/run_engine.py tests/unit/test_transition_models.py + tests/unit/runtime tests/unit/orchestration/test_agent_runner.py - name: Test foundation - run: | - pytest -q tests/unit/test_transition_models.py tests/unit/runtime/test_transition_engine.py tests/unit/runtime/test_foundry_transition_adapter.py tests/unit/orchestration/test_agent_runner.py + run: >- + pytest -q tests/unit/test_transition_models.py tests/unit/runtime + tests/unit/orchestration/test_agent_runner.py - name: Test full suite run: pytest -q diff --git a/docs/evidence-retention.md b/docs/evidence-retention.md new file mode 100644 index 0000000..293de0e --- /dev/null +++ b/docs/evidence-retention.md @@ -0,0 +1,89 @@ +# Evidence retention before runtime convergence + +Engineering pass: 25 September 2026. Audited baseline: +`2639b8ebeeb485950a8faed17491b03ff440b354` (after PR #5). +This is new implementation work, not a claim about the original experiment. + +## The gap + +The first adapter stored an outcome containing a patch checksum, but not the +patch itself. It then deleted the worktree. A fingerprint was therefore being +used where retrievable evidence was needed. Processing exceptions also bypassed +outcome journaling, and a supplied before-state bypassed fresh observation. +Passing the earlier tests did not establish those properties. + +## This change + +`FoundryTransitionRuntime` now wires the artifact store into its planner, +executor and verifier. It saves the plan before implementation, independently +captures and saves the real Git-visible patch, and stores complete verification +and review reports. Outcome references contain retrievable `artifact://runs/...` +paths and checksums. Individual adapters retain an optional non-persisting mode +for compatibility; use the runtime factory for the evidence-preserving path. + +Git capture uses a temporary index rather than changing the real staging area. +It includes tracked edits/deletions and non-ignored additions, including binary +patches and literal filenames. The observer and executor use the same capture +mechanism. The persistent verifier checks the stored patch against the workspace +rather than trusting the agent's returned diff. An executor that changes HEAD +is rejected for inspection, with its workspace retained. + +The transition engine freshly observes before-state even when a prior snapshot +is supplied. That supplied snapshot is an explicit precondition. It checks +observation environment identity and detects state changes during verification. +Snapshot `state` must therefore be a stable, comparable representation; this is +not a proof that an arbitrary representation is accurate or complete. + +Processing failures and cooperative cancellation attempt to journal an +unaccepted outcome with phase, exception type, and retained workspace. Raw +exception messages are not copied into that record. The original exception is +re-raised. If journaling also fails, the original exception gains an explicit +note; no durable record is claimed. No automatic retry or destructive cleanup +is attempted on these paths. + +A normal decision is journaled before cleanup. Cleanup errors are recorded as +resource warnings without converting a verified decision into a different +verdict. Atomic same-filesystem artifact replacement avoids exposing partial +new files, and `ArtifactStoreTransitionJournal.load()` allows a new process to +read and validate an outcome. This is local persistence, not an execution log +with exactly-once or crash-resume guarantees. + +## Evidence to run + +```bash +pytest -q tests/unit/runtime/test_transition_integrity.py \ + tests/unit/runtime/test_git_evidence_retention.py +pytest -q +``` + +The new tests use actual files and Git repositories. They cover two counter +transitions separated by a fresh-process journal reload; stale/wrong state; +failures across processing phases; cancellation and cleanup/persistence errors; +atomic-write failure; verifier mutation; and patch replay after worktree removal. +Patch tests include staged edits, deletion, binary additions, spaces, newlines, +and literal pathspec-shaped names. Model behavior is deterministic test code, +not a live LLM evaluation. The full regression suite remains the compatibility +gate; a count of tests is not a completeness or production-readiness claim. + +## Still not established + +- The historical WorktreeManager still creates from HEAD and ignores the + adapter's requested base_ref. Fix and test that before runtime convergence. +- The historical RunEngine has not been routed through this path. PR creation, + adoption of a worktree result, and deployment remain separate. A saved patch + is not an applied change in the source repository. +- Historical verifier coverage and REQUEST_CHANGES-as-advisory policy remain. + An accepted outcome is not a claim of comprehensive verification or approval + for production publication. +- Abrupt process death, concurrent execution, duplicate transition IDs, + distributed storage, and exactly-once actions are not solved here. Outcome + records can be replaced under the same ID; do not retry blindly. +- Ignored files, external LFS objects and submodule contents are not a + self-contained workspace backup. Git filters are not sandboxed. The + secret-shaped filename check is not a content-based secret scanner. +- Retained workspaces consume resources and require explicit inspection and + cleanup. Keeping evidence is not rollback. + +These limitations are prerequisites for the convergence plan in +[runtime-decoupling-audit.md](runtime-decoupling-audit.md), not reasons to rename +or rewrite the historical implementation again. diff --git a/foundry/adapters/foundry_transition.py b/foundry/adapters/foundry_transition.py index 908bab3..c655645 100644 --- a/foundry/adapters/foundry_transition.py +++ b/foundry/adapters/foundry_transition.py @@ -1,10 +1,14 @@ -"""Historical Foundry capabilities exposed through UCF transition interfaces.""" +"""Historical Foundry capabilities exposed through UCF transition interfaces. + +Use FoundryTransitionRuntime for persisted evidence. The optional stores on +individual adapters retain compatibility with isolated, non-persisting callers. +""" from __future__ import annotations -import asyncio import hashlib -from typing import Literal +import json +from typing import TYPE_CHECKING, Literal from uuid import UUID, uuid4 from foundry.contracts.review_models import ReviewVerdict @@ -19,25 +23,42 @@ VerificationDecision, ) from foundry.environments.git_observer import GitStateObserver +from foundry.environments.git_patch import capture_patch from foundry.environments.git_worktree import GitWorktreeEnvironment from foundry.git.branch import generate_branch_name -from foundry.git.worktree import WorktreeManager -from foundry.orchestration.agent_runner import AgentRunner -from foundry.runtime.transition_engine import TransitionEngine -from foundry.storage.artifact_store import ArtifactStore +from foundry.runtime.transition_engine import StateConflictError, TransitionEngine +from foundry.storage.artifact_store import ArtifactStore, ArtifactType from foundry.storage.transition_journal import ArtifactStoreTransitionJournal from foundry.verification.policy import ( MIGRATION_GUARD_ALLOWED_TASK_TYPES, match_protected_paths, ) -from foundry.verification.runner import VerificationRunner + +if TYPE_CHECKING: + from foundry.git.worktree import WorktreeManager + from foundry.orchestration.agent_runner import AgentRunner + from foundry.verification.runner import VerificationRunner def _task_from_transition(request: TransitionRequest) -> TaskRequest: raw = request.metadata.get("foundry_task") if not isinstance(raw, dict): raise ValueError("TransitionRequest metadata is missing foundry_task") - return TaskRequest.model_validate(raw) + task = TaskRequest.model_validate(raw) + return task.model_copy(update={"metadata": {**task.metadata, "run_id": str(request.id)}}) + + +async def _store_evidence( + store: ArtifactStore, request: TransitionRequest, kind: ArtifactType, + filename: str, data: str | bytes, +) -> EvidenceRef: + result = await store.store(request.id, kind, data, filename=filename) + return EvidenceRef( + kind=kind.value, + uri=f"artifact://{result['storage_path']}", + checksum=result["checksum"], + description=f"{filename} ({result['size_bytes']} bytes)", + ) def foundry_task_to_transition_request( @@ -67,241 +88,226 @@ def foundry_task_to_transition_request( class FoundryPlannerAdapter: - """Expose AgentRunner planning through the UCF TransitionPlanner contract.""" + """Expose planning, persisting the plan before implementation when configured.""" - def __init__(self, agent_runner: AgentRunner) -> None: + def __init__( + self, agent_runner: AgentRunner, artifact_store: ArtifactStore | None = None, + ) -> None: self.agent_runner = agent_runner + self.artifact_store = artifact_store async def plan( - self, - request: TransitionRequest, - current_state: StateSnapshot, - workspace: str, + self, request: TransitionRequest, current_state: StateSnapshot, workspace: str, ) -> ActionProposal: task = _task_from_transition(request) plan = await self.agent_runner.run_planner(task, workspace) + payload = {"plan": plan.model_dump(), "base_head": current_state.state.get("head")} + if self.artifact_store is not None: + ref = await _store_evidence( + self.artifact_store, request, ArtifactType.PLAN, "plan.json", + plan.model_dump_json(indent=2), + ) + payload["plan_evidence"] = ref.model_dump() return ActionProposal( kind="foundry-plan", - description=( - f"Apply {len(plan.steps)} planned steps to {request.environment} " - f"from {current_state.state.get('head', 'unknown state')}" - ), - payload={"plan": plan.model_dump()}, + description=f"Apply {len(plan.steps)} planned steps to {request.environment}", + payload=payload, ) class FoundryExecutorAdapter: - """Expose AgentRunner implementation through the UCF ActionExecutor contract.""" + """Persist the actual workspace patch, not just a provider's claimed diff.""" - def __init__(self, agent_runner: AgentRunner) -> None: + def __init__( + self, agent_runner: AgentRunner, artifact_store: ArtifactStore | None = None, + ) -> None: self.agent_runner = agent_runner + self.artifact_store = artifact_store async def execute( - self, - request: TransitionRequest, - action: ActionProposal, - workspace: str, + self, request: TransitionRequest, action: ActionProposal, workspace: str, ) -> list[EvidenceRef]: if action.kind != "foundry-plan": raise ValueError(f"Unsupported Foundry action kind: {action.kind}") - task = _task_from_transition(request) plan = PlanArtifact.model_validate(action.payload["plan"]) language = request.metadata.get("language", "go") if language not in {"go", "typescript"}: raise ValueError(f"Unsupported Foundry implementation language: {language}") - - diff = await self.agent_runner.run_implementer( - plan=plan, - task_request=task, - worktree_path=workspace, - language=language, + claimed_diff = await self.agent_runner.run_implementer( + plan=plan, task_request=task, worktree_path=workspace, language=language, ) - checksum = hashlib.sha256(diff.encode("utf-8")).hexdigest() - return [ - EvidenceRef( - kind="git-diff", - uri=f"sha256:{checksum}", - description=f"Implementation produced {len(diff.encode('utf-8'))} diff bytes", - checksum=checksum, - ) - ] + if self.artifact_store is None: + checksum = hashlib.sha256(claimed_diff.encode("utf-8")).hexdigest() + return [EvidenceRef(kind="git-diff", uri=f"sha256:{checksum}", checksum=checksum)] + base = action.payload.get("base_head") + if not isinstance(base, str) or not base: + raise ValueError("Persistent Git execution requires an observed base HEAD") + patch = await capture_patch(workspace, base) + ref = await _store_evidence( + self.artifact_store, request, ArtifactType.DIFF, "diff.patch", patch.data, + ) + refs = [ref] + if "plan_evidence" in action.payload: + refs.insert(0, EvidenceRef.model_validate(action.payload["plan_evidence"])) + return refs async def _read_diff(workspace: str) -> str: - proc = await asyncio.create_subprocess_exec( - "git", - "diff", - "HEAD", - cwd=workspace, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() - if proc.returncode != 0: - raise RuntimeError( - f"git diff failed in {workspace}: {stderr.decode(errors='replace').strip()}" - ) - return stdout.decode(errors="replace") + patch = await capture_patch(workspace) + return patch.data.decode("utf-8", errors="replace") def _verification_evidence(check_type: str, passed: bool, output: str) -> EvidenceRef: status = "pass" if passed else "fail" return EvidenceRef( - kind="verification", - uri=f"verification://{check_type}/{status}", + kind="verification", uri=f"verification://{check_type}/{status}", description=output[:500] or f"{check_type}: {status}", ) def _review_evidence(kind: str, review: ReviewVerdict) -> EvidenceRef: return EvidenceRef( - kind=kind, - uri=f"{kind}://{review.verdict.value}", - description=review.summary, + kind=kind, uri=f"{kind}://{review.verdict.value}", description=review.summary, ) class FoundryVerifierAdapter: - """Apply Foundry verification and blind review to a UCF transition.""" + """Persist full verification and review results; preserve historical verdict policy.""" def __init__( - self, - verification_runner: VerificationRunner, - agent_runner: AgentRunner, + self, verification_runner: VerificationRunner, agent_runner: AgentRunner, + artifact_store: ArtifactStore | None = None, ) -> None: self.verification_runner = verification_runner self.agent_runner = agent_runner + self.artifact_store = artifact_store + + async def _review_ref( + self, request: TransitionRequest, kind: str, review: ReviewVerdict, + ) -> EvidenceRef: + if self.artifact_store is None: + return _review_evidence(kind, review) + return await _store_evidence( + self.artifact_store, request, ArtifactType.REVIEW, + f"{kind}.json", review.model_dump_json(indent=2), + ) async def verify( - self, - request: TransitionRequest, - before: StateSnapshot, - after: StateSnapshot, - action: ActionProposal, - action_evidence: list[EvidenceRef], - workspace: str, + self, request: TransitionRequest, before: StateSnapshot, after: StateSnapshot, + action: ActionProposal, action_evidence: list[EvidenceRef], workspace: str, ) -> VerificationDecision: - del before, action, action_evidence task = _task_from_transition(request) - changed_files = [ - str(path) for path in after.state.get("changed_files", []) - ] + changed_files = [str(path) for path in after.state.get("changed_files", [])] evidence: list[EvidenceRef] = [] + if self.artifact_store is not None: + if before.state.get("head") != after.state.get("head"): + raise StateConflictError("Executor changed Git HEAD; workspace retained for review") + patch = await capture_patch(workspace) + refs = [ref for ref in action_evidence if ref.kind == ArtifactType.DIFF.value] + if len(refs) != 1 or not refs[0].uri.startswith("artifact://"): + raise StateConflictError("Missing retrievable execution patch") + saved = await self.artifact_store.retrieve(refs[0].uri.removeprefix("artifact://")) + if hashlib.sha256(saved).hexdigest() != refs[0].checksum or saved != patch.data: + raise StateConflictError("Saved patch does not match the verification workspace") + diff = patch.data.decode("utf-8", errors="replace") + else: + diff = await _read_diff(workspace) if task.verify: results, passed = await self.verification_runner.run_all( - workspace, - changed_files, - run_id=request.id, - ) - evidence.extend( - _verification_evidence(result.check_type, result.passed, result.output) - for result in results + workspace, changed_files, run_id=request.id, ) - if not passed: - failed = ", ".join( - result.check_type for result in results if not result.passed + if self.artifact_store is None: + evidence.extend( + _verification_evidence(result.check_type, result.passed, result.output) + for result in results ) + else: + report = {"passed": passed, "checks": [ + {"check_type": result.check_type, "passed": result.passed, + "output": result.output, "duration_ms": result.duration_ms} + for result in results + ]} + evidence.append(await _store_evidence( + self.artifact_store, request, ArtifactType.VERIFICATION, + "verification.json", json.dumps(report, indent=2), + )) + if not passed: + failed = ", ".join(result.check_type for result in results if not result.passed) return VerificationDecision( - accepted=False, - reason=f"Deterministic verification failed: {failed}", + accepted=False, reason=f"Deterministic verification failed: {failed}", evidence=evidence, ) - diff = await _read_diff(workspace) protected_files = match_protected_paths(changed_files) if protected_files: if task.task_type == TaskType.BUG_FIX: return VerificationDecision( - accepted=False, - reason=( - "Protected-path policy rejected bug-fix transition: " - + ", ".join(protected_files) - ), - evidence=evidence, + accepted=False, reason="Protected-path policy rejected bug-fix transition: " + + ", ".join(protected_files), evidence=evidence, ) - if task.task_type not in MIGRATION_GUARD_ALLOWED_TASK_TYPES: return VerificationDecision( accepted=False, reason=( - f"Task type {task.task_type.value} is not authorized " - f"for protected paths: {', '.join(protected_files)}" + f"Task type {task.task_type.value} is not authorized for protected paths" ), evidence=evidence, ) - guard = await self.agent_runner.run_migration_guard( - diff=diff, - changed_files=protected_files, + diff=diff, changed_files=protected_files, ) - evidence.append(_review_evidence("migration-guard", guard)) + evidence.append(await self._review_ref(request, "migration-guard", guard)) if guard.verdict == ReviewVerdictType.REJECT: return VerificationDecision( - accepted=False, - reason=f"Migration guard rejected transition: {guard.summary}", + accepted=False, reason=f"Migration guard rejected transition: {guard.summary}", evidence=evidence, ) changed_summary = ", ".join(changed_files[:10]) or "no files detected" if len(changed_files) > 10: changed_summary += f" (+{len(changed_files) - 10} more)" - review = await self.agent_runner.run_reviewer( diff=diff, pr_title=f"[Foundry] {task.task_type.value}: {task.title}", pr_description=f"{task.prompt[:500]}\n\nChanged files: {changed_summary}", changed_files=changed_files, ) - evidence.append(_review_evidence("review", review)) - - accepted = review.verdict != ReviewVerdictType.REJECT + evidence.append(await self._review_ref(request, "review", review)) + reason = review.summary if review.verdict == ReviewVerdictType.REQUEST_CHANGES: reason = ( "Independent review requested changes; historical Foundry policy " "treats this verdict as advisory." ) - else: - reason = review.summary - return VerificationDecision( - accepted=accepted, - reason=reason, - evidence=evidence, + accepted=review.verdict != ReviewVerdictType.REJECT, + reason=reason, evidence=evidence, ) class FoundryTransitionRuntime: - """Run historical Foundry work through the provider-neutral UCF loop.""" + """Run Foundry work with persisted plans, patches, reports and outcomes.""" def __init__( - self, - *, - worktree_manager: WorktreeManager, - agent_runner: AgentRunner, - verification_runner: VerificationRunner, - artifact_store: ArtifactStore, + self, *, worktree_manager: WorktreeManager, agent_runner: AgentRunner, + verification_runner: VerificationRunner, artifact_store: ArtifactStore, ) -> None: self.engine = TransitionEngine( environment=GitWorktreeEnvironment(worktree_manager), observer=GitStateObserver(), - planner=FoundryPlannerAdapter(agent_runner), - executor=FoundryExecutorAdapter(agent_runner), - verifier=FoundryVerifierAdapter(verification_runner, agent_runner), + planner=FoundryPlannerAdapter(agent_runner, artifact_store), + executor=FoundryExecutorAdapter(agent_runner, artifact_store), + verifier=FoundryVerifierAdapter(verification_runner, agent_runner, artifact_store), journal=ArtifactStoreTransitionJournal(artifact_store), ) async def execute_task( - self, - task: TaskRequest, - *, - transition_id: UUID | None = None, + self, task: TaskRequest, *, transition_id: UUID | None = None, language: Literal["go", "typescript"] = "go", ) -> TransitionOutcome: request = foundry_task_to_transition_request( - task, - transition_id=transition_id, - language=language, + task, transition_id=transition_id, language=language, ) return await self.engine.execute(request) diff --git a/foundry/environments/git_observer.py b/foundry/environments/git_observer.py index 3f3f934..364da6a 100644 --- a/foundry/environments/git_observer.py +++ b/foundry/environments/git_observer.py @@ -1,72 +1,39 @@ -"""Git-backed state observation for UCF execution environments.""" +"""Observe Git-visible state using the same patch capture as durable evidence.""" from __future__ import annotations -import asyncio import hashlib from datetime import UTC, datetime from foundry.contracts.transition_models import EvidenceRef, StateSnapshot, TransitionRequest - - -async def _git(workspace: str, *args: str) -> str: - proc = await asyncio.create_subprocess_exec( - "git", - *args, - cwd=workspace, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() - if proc.returncode != 0: - raise RuntimeError( - f"git {' '.join(args)} failed in {workspace}: " - f"{stderr.decode(errors='replace').strip()}" - ) - return stdout.decode(errors="replace") - - -def _changed_paths(status: str) -> list[str]: - paths: list[str] = [] - for raw_line in status.splitlines(): - if len(raw_line) < 4: - continue - path = raw_line[3:] - if " -> " in path: - path = path.split(" -> ", 1)[1] - paths.append(path) - return paths +from foundry.environments.git_patch import capture_patch class GitStateObserver: - """Represent the current Git workspace as explicit UCF state.""" + """Represent tracked and non-ignored untracked changes, not only tracked diffs.""" async def observe(self, request: TransitionRequest, workspace: str) -> StateSnapshot: - head = (await _git(workspace, "rev-parse", "HEAD")).strip() - status = await _git(workspace, "status", "--porcelain") - diff = await _git(workspace, "diff", "HEAD") - checksum = hashlib.sha256(diff.encode("utf-8")).hexdigest() - changed_files = _changed_paths(status) - + patch = await capture_patch(workspace) + checksum = hashlib.sha256(patch.data).hexdigest() return StateSnapshot( environment=request.environment, observed_at=datetime.now(UTC), state={ - "head": head, - "dirty": bool(status.strip()), - "changed_files": changed_files, + "head": patch.base_commit, + "dirty": bool(patch.changed_files), + "changed_files": list(patch.changed_files), "diff_checksum": checksum, }, evidence=[ EvidenceRef( kind="git-head", - uri=f"git://commit/{head}", + uri=f"git://commit/{patch.base_commit}", description="Observed repository HEAD", ), EvidenceRef( kind="git-diff", uri=f"sha256:{checksum}", - description="Checksum of working tree diff against HEAD", + description="Fingerprint; retrievable patch is stored by the executor", checksum=checksum, ), ], diff --git a/foundry/environments/git_patch.py b/foundry/environments/git_patch.py new file mode 100644 index 0000000..25256a5 --- /dev/null +++ b/foundry/environments/git_patch.py @@ -0,0 +1,91 @@ +"""Capture a replayable Git-visible patch without changing the user's index. + +Includes tracked edits/deletions and non-ignored additions, including binaries. +Ignored files are outside this contract. This is not a sandbox for Git filters. +""" + +from __future__ import annotations + +import asyncio +import fnmatch +import os +import tempfile +from dataclasses import dataclass +from pathlib import Path + + +@dataclass(frozen=True) +class GitPatch: + """A patch relative to one pinned commit, with literal UTF-8 changed paths.""" + + base_commit: str + data: bytes + changed_files: tuple[str, ...] + + +async def _git( + workspace: str, + *args: str, + env: dict[str, str] | None = None, + stdin: bytes | None = None, +) -> bytes: + proc = await asyncio.create_subprocess_exec( + "git", "--literal-pathspecs", *args, + cwd=workspace, + env=env, + stdin=asyncio.subprocess.PIPE if stdin is not None else asyncio.subprocess.DEVNULL, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + try: + stdout, _ = await asyncio.wait_for(proc.communicate(stdin), timeout=30) + except (TimeoutError, asyncio.CancelledError): + if proc.returncode is None: + proc.kill() + await proc.wait() + raise + if proc.returncode != 0: + # Do not copy arbitrary Git/filter stderr into persisted exception messages. + raise RuntimeError(f"Git evidence command failed: {args[0]} (exit {proc.returncode})") + return stdout + + +def _sensitive_path(path: str) -> bool: + lowered = path.lower() + patterns = ("*.env", "*.env.*", "*.pem", "*.key", "*secrets*", "*credentials*", + "*service-account*") + return any(fnmatch.fnmatch(lowered, pattern) for pattern in patterns) + + +async def capture_patch(workspace: str, base_ref: str = "HEAD") -> GitPatch: + """Capture a complete patch against base_ref, failing before secret-path reads. + + A temporary index stages only the candidate paths. The real index remains + untouched. The caller must persist data before deleting the workspace. + """ + base = (await _git( + workspace, "rev-parse", "--verify", "--end-of-options", f"{base_ref}^{{commit}}" + )).decode("ascii").strip() + tracked = await _git( + workspace, "diff", "--no-ext-diff", "--no-textconv", "--no-renames", + "--name-only", "-z", base, "--", + ) + untracked = await _git(workspace, "ls-files", "--others", "--exclude-standard", "-z") + paths = sorted({part.decode("utf-8") for part in (tracked + untracked).split(b"\0") if part}) + if any(_sensitive_path(path) for path in paths): + raise PermissionError("Refusing to capture a secret-shaped changed path") + with tempfile.TemporaryDirectory(prefix="ucf-index-") as staging: + env = os.environ.copy() + env["GIT_INDEX_FILE"] = str(Path(staging) / "index") + await _git(workspace, "read-tree", base, env=env) + if paths: + literals = b"".join(path.encode("utf-8") + b"\0" for path in paths) + await _git( + workspace, "add", "--all", "--pathspec-from-file=-", "--pathspec-file-nul", + env=env, stdin=literals, + ) + data = await _git( + workspace, "diff", "--cached", "--binary", "--full-index", "--no-ext-diff", + "--no-textconv", "--no-renames", base, "--", env=env, + ) + return GitPatch(base_commit=base, data=data, changed_files=tuple(paths)) diff --git a/foundry/runtime/transition_engine.py b/foundry/runtime/transition_engine.py index f754df4..f451043 100644 --- a/foundry/runtime/transition_engine.py +++ b/foundry/runtime/transition_engine.py @@ -1,12 +1,18 @@ -"""Minimal provider-neutral state transition engine. +"""Evidence-preserving transition execution, not automatic crash recovery. -This runs alongside the historical Foundry RunEngine. It exists to prove the -general UCF loop before legacy Git/PR terminology is migrated. +Outcomes describe observed work, not publication or deployment. Processing and +journal failures retain the workspace; a failed operation never becomes accepted. """ from __future__ import annotations +import asyncio +import logging + from foundry.contracts.transition_models import ( + ActionProposal, + EvidenceRef, + StateSnapshot, TransitionObservation, TransitionOutcome, TransitionRequest, @@ -20,9 +26,25 @@ TransitionVerifier, ) +logger = logging.getLogger(__name__) + + +class StateConflictError(ValueError): + """Observed state does not match the transition's environment or precondition.""" + + +def _validate_snapshot(snapshot: StateSnapshot, request: TransitionRequest) -> None: + if snapshot.environment != request.environment: + raise StateConflictError("Observation belongs to a different environment") + class TransitionEngine: - """Coordinate one explicit state transition from observation to outcome.""" + """Execute, observe, verify and journal; preserve evidence on failure. + + Adapters are trusted capabilities, not a security sandbox. A supplied + before_state is a precondition to compare with a fresh observation, not a + replacement for observation. State equality is adapter-defined via `state`. + """ def __init__( self, @@ -41,47 +63,123 @@ def __init__( self.journal = journal async def execute(self, request: TransitionRequest) -> TransitionOutcome: - workspace = await self.environment.prepare( - target=request.environment, - base_ref=str(request.metadata.get("base_ref", "current")), - transition_id=request.id, - transition_name=str(request.metadata.get("transition_name", request.id)), - ) + """Return a journaled decision, or re-raise a recorded processing failure. + No automatic retry is attempted. A retained workspace needs explicit + inspection and cleanup. Persistence failures are surfaced to the caller. + """ + phase = "prepare" + workspace: str | None = None + before: StateSnapshot | None = None + after: StateSnapshot | None = None + action: ActionProposal | None = None + evidence: list[EvidenceRef] = [] try: - before = request.before_state - if before is None: - before = await self.observer.observe(request, workspace) - - action = await self.planner.plan(request, before, workspace) - action_evidence = await self.executor.execute(request, action, workspace) - after = await self.observer.observe(request, workspace) + workspace = await self.environment.prepare( + target=request.environment, + base_ref=str(request.metadata.get("base_ref", "current")), + transition_id=request.id, + transition_name=str(request.metadata.get("transition_name", request.id)), + ) + phase = "observe_before" + before = (await self.observer.observe(request, workspace)).model_copy(deep=True) + _validate_snapshot(before, request) + if request.before_state is not None: + _validate_snapshot(request.before_state, request) + if request.before_state.state != before.state: + raise StateConflictError("Supplied before_state is stale") + phase = "plan" + action = await self.planner.plan(request, before.model_copy(deep=True), workspace) + phase = "execute" + evidence = await self.executor.execute(request, action, workspace) + phase = "observe_after" + after = (await self.observer.observe(request, workspace)).model_copy(deep=True) + _validate_snapshot(after, request) + phase = "verify" decision = await self.verifier.verify( request=request, - before=before, - after=after, - action=action, - action_evidence=action_evidence, + before=before.model_copy(deep=True), + after=after.model_copy(deep=True), + action=action.model_copy(deep=True), + action_evidence=[item.model_copy(deep=True) for item in evidence], workspace=workspace, ) - - observation = TransitionObservation( - transition_id=request.id, - observed_state=after, - action_evidence=action_evidence, - verification_evidence=decision.evidence, - ) - outcome = TransitionOutcome( + phase = "observe_verified" + verified = await self.observer.observe(request, workspace) + _validate_snapshot(verified, request) + if verified.state != after.state: + after = verified.model_copy(deep=True) + raise StateConflictError("Environment changed during verification") + except (Exception, asyncio.CancelledError) as error: + # Failure is not a verification rejection. Do not erase partial work. + failed = TransitionOutcome( transition_id=request.id, - accepted=decision.accepted, + accepted=False, before_state=before, after_state=after, - observation=observation, - reason=decision.reason, - metadata={"action_kind": action.kind}, + observation=( + TransitionObservation( + transition_id=request.id, + observed_state=after, + action_evidence=evidence, + ) + if after is not None else None + ), + reason=f"Transition failed during {phase}", + metadata={ + "status": ( + "cancelled" if isinstance(error, asyncio.CancelledError) else "failed" + ), + "phase": phase, + "error_type": type(error).__name__, + "environment": request.environment, + "retained_workspace": workspace, + }, ) + try: + await self.journal.record(failed) + except Exception as journal_error: + # Keep the original exception and make the missing record explicit. + error.add_note(f"Failure journal also failed: {type(journal_error).__name__}") + logger.error("Failure journal unavailable for transition %s", request.id) + error.add_note(f"UCF retained workspace: {workspace!r}; no automatic retry") + raise + + outcome = TransitionOutcome( + transition_id=request.id, + accepted=decision.accepted, + before_state=before, + after_state=after, + observation=TransitionObservation( + transition_id=request.id, + observed_state=after, + action_evidence=evidence, + verification_evidence=decision.evidence, + ), + reason=decision.reason, + metadata={ + "status": "accepted" if decision.accepted else "rejected", + "action_kind": action.kind, + }, + ) + # A successful call to record is required BEFORE destructive cleanup. + try: await self.journal.record(outcome) - return outcome - finally: + except Exception as error: + error.add_note(f"Outcome not confirmed durable; retained workspace: {workspace!r}") + raise + try: await self.environment.cleanup(workspace) + except (Exception, asyncio.CancelledError) as error: + outcome.metadata.update({ + "cleanup_status": "failed", + "cleanup_error_type": type(error).__name__, + "retained_workspace": workspace, + }) + # Persist the resource warning without changing the verification decision. + await self.journal.record(outcome) + if isinstance(error, asyncio.CancelledError): + raise + logger.warning("Cleanup failed for transition %s", request.id) + return outcome diff --git a/foundry/storage/artifact_store.py b/foundry/storage/artifact_store.py index aec1b47..59ea5e7 100644 --- a/foundry/storage/artifact_store.py +++ b/foundry/storage/artifact_store.py @@ -7,6 +7,8 @@ import hashlib import logging +import os +import tempfile from datetime import UTC, datetime from enum import StrEnum from pathlib import Path @@ -64,17 +66,11 @@ async def store( data: bytes | str, filename: str | None = None, ) -> StoreResult: - """Store an artifact for a run. + """Store an artifact using same-filesystem atomic replacement. - Args: - run_id: The run that produced this artifact. - artifact_type: Type of artifact (plan, diff, review, etc.). - data: Raw artifact data as bytes or string. - filename: Optional custom filename. Defaults to - '{artifact_type}.json' (or '.patch' for diffs). - - Returns: - Dict with storage_path (relative), size_bytes, and SHA-256 checksum. + The temporary file is flushed and synced before replacement, then its + directory is synced. This is local file persistence, not a distributed + transaction or a guarantee of exactly-once action execution. """ if filename is None: ext = ".patch" if artifact_type == ArtifactType.DIFF else ".json" @@ -85,7 +81,18 @@ async def store( full_path.parent.mkdir(parents=True, exist_ok=True) content = data if isinstance(data, bytes) else data.encode("utf-8") - full_path.write_bytes(content) + with tempfile.TemporaryDirectory(prefix=".ucf-write-", dir=full_path.parent) as staging: + temporary = Path(staging) / "artifact" + with temporary.open("wb") as handle: + handle.write(content) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temporary, full_path) + directory_fd = os.open(full_path.parent, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) + try: + os.fsync(directory_fd) + finally: + os.close(directory_fd) size_bytes = len(content) checksum = hashlib.sha256(content).hexdigest() @@ -101,42 +108,21 @@ async def store( ) async def retrieve(self, storage_path: str) -> bytes: - """Retrieve an artifact by its storage path. - - Args: - storage_path: The path returned by store(). - - Returns: - Raw artifact data as bytes. - - Raises: - FileNotFoundError: If the artifact does not exist at the given path. - """ + """Retrieve an artifact by its storage path.""" full_path = self.base_path / storage_path if not full_path.exists(): raise FileNotFoundError(f"Artifact not found: {storage_path}") return full_path.read_bytes() async def delete(self, storage_path: str) -> None: - """Delete an artifact from storage. - - Args: - storage_path: The path of the artifact to delete. - """ + """Delete an artifact from storage.""" full_path = self.base_path / storage_path if full_path.exists(): full_path.unlink() logger.info("Deleted artifact: %s", storage_path) async def list_artifacts(self, run_id: UUID) -> list[ArtifactInfo]: - """List all artifacts for a given run with metadata. - - Args: - run_id: The run to list artifacts for. - - Returns: - List of dicts with filename, size_bytes, and modified (ISO timestamp). - """ + """List all artifacts for a given run with metadata.""" run_dir = self.base_path / "runs" / str(run_id) if not run_dir.exists(): return [] diff --git a/foundry/storage/transition_journal.py b/foundry/storage/transition_journal.py index 80d239a..5d6e9ae 100644 --- a/foundry/storage/transition_journal.py +++ b/foundry/storage/transition_journal.py @@ -1,21 +1,34 @@ -"""Transition journal backed by the existing Foundry artifact store.""" +"""Persist and reload outcomes independently from model or process memory.""" from __future__ import annotations +from uuid import UUID + from foundry.contracts.transition_models import TransitionOutcome from foundry.storage.artifact_store import ArtifactStore, ArtifactType class ArtifactStoreTransitionJournal: - """Persist verified UCF outcomes using Foundry's durable artifact storage.""" + """Local atomic outcome storage; not a resumable execution scheduler.""" def __init__(self, artifact_store: ArtifactStore) -> None: self.artifact_store = artifact_store async def record(self, outcome: TransitionOutcome) -> None: + """Persist a decision, failure, or later cleanup warning under its ID.""" await self.artifact_store.store( outcome.transition_id, ArtifactType.TRANSITION, outcome.model_dump_json(indent=2), filename="transition_outcome.json", ) + + async def load(self, transition_id: UUID) -> TransitionOutcome: + """Load and validate a persisted record after constructing a new journal.""" + data = await self.artifact_store.retrieve( + f"runs/{transition_id}/transition_outcome.json" + ) + outcome = TransitionOutcome.model_validate_json(data) + if outcome.transition_id != transition_id: + raise ValueError("Stored outcome has a different transition ID") + return outcome diff --git a/tests/unit/runtime/test_git_evidence_retention.py b/tests/unit/runtime/test_git_evidence_retention.py new file mode 100644 index 0000000..68401f7 --- /dev/null +++ b/tests/unit/runtime/test_git_evidence_retention.py @@ -0,0 +1,244 @@ +"""Real Git/index/artifact tests; model and verifier behavior is deterministic.""" + +from __future__ import annotations + +import hashlib +import subprocess +from dataclasses import dataclass +from pathlib import Path +from uuid import UUID, uuid4 + +import pytest + +from foundry.adapters.foundry_transition import FoundryTransitionRuntime +from foundry.contracts.review_models import ReviewVerdict +from foundry.contracts.shared import Complexity, ReviewVerdictType, TaskType +from foundry.contracts.task_types import PlanArtifact, PlanStep, TaskRequest +from foundry.environments.git_patch import capture_patch +from foundry.git.worktree import WorktreeManager +from foundry.runtime.transition_engine import StateConflictError +from foundry.storage.artifact_store import ArtifactStore +from foundry.storage.transition_journal import ArtifactStoreTransitionJournal + + +def git(root: Path, *args: str, data: bytes | None = None) -> bytes: + return subprocess.run( + ["git", *args], cwd=root, input=data, check=True, capture_output=True, + ).stdout + + +def repository(root: Path) -> Path: + source = root / "source" + source.mkdir() + git(source, "init", "-b", "main") + git(source, "config", "user.email", "ucf-test@example.invalid") + git(source, "config", "user.name", "UCF Test") + (source / "main.py").write_text("before\n", encoding="utf-8") + (source / "remove.txt").write_text("remove this\n", encoding="utf-8") + git(source, "add", ".") + git(source, "commit", "-m", "Initial fixture") + return source + + +@pytest.mark.parametrize("name", ["new file.txt", "new\nline.txt", ":(glob)*.txt", "--flag.txt"]) +async def test_patch_replays_additions_deletions_and_binary_without_changing_index( + tmp_path: Path, name: str, +) -> None: + source = repository(tmp_path) + (source / "main.py").write_text("after\n", encoding="utf-8") + git(source, "add", "main.py") + index_before = git(source, "diff", "--cached", "--binary") + (source / "remove.txt").unlink() + (source / name).write_text("added\n", encoding="utf-8") + (source / "new.bin").write_bytes(b"\x00\xff\x00new") + patch = await capture_patch(str(source)) + assert git(source, "diff", "--cached", "--binary") == index_before + assert set(patch.changed_files) == {"main.py", "remove.txt", "new.bin", name} + assert b"GIT binary patch" in patch.data + git(source, "reset", "--hard") + git(source, "clean", "-fd") + git(source, "apply", "--binary", "-", data=patch.data) + assert (source / "main.py").read_text(encoding="utf-8") == "after\n" + assert not (source / "remove.txt").exists() + assert (source / name).read_text(encoding="utf-8") == "added\n" + assert (source / "new.bin").read_bytes() == b"\x00\xff\x00new" + + +async def test_secret_shaped_paths_are_rejected_before_staging( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + source = repository(tmp_path) + (source / ".env").write_text("synthetic fixture, not a credential", encoding="utf-8") + + def unexpected_staging(*args: object, **kwargs: object) -> None: + raise AssertionError("Sensitive-path rejection must happen before opening an index") + + monkeypatch.setattr( + "foundry.environments.git_patch.tempfile.TemporaryDirectory", unexpected_staging, + ) + with pytest.raises(PermissionError, match="secret-shaped"): + await capture_patch(str(source)) + assert git(source, "diff", "--cached") == b"" + + +async def test_clean_patch_and_invalid_refs_are_not_confused(tmp_path: Path) -> None: + source = repository(tmp_path) + clean = await capture_patch(str(source)) + assert clean.data == b"" and not clean.changed_files + with pytest.raises(RuntimeError, match="Git evidence command failed"): + await capture_patch(str(source), "missing-reference") + + +@dataclass +class Check: + check_type: str = "fixture-check" + passed: bool = True + output: str = "Checked actual fixture files" + duration_ms: int = 1 + + +class DeterministicRunner: + """Writes actual files but returns a deliberately false claimed diff.""" + + def __init__(self, fail: str = "") -> None: + self.fail = fail + self.implementations = 0 + + async def run_planner(self, task: TaskRequest, workspace: str) -> PlanArtifact: + return PlanArtifact( + task_id=UUID(task.metadata["run_id"]), + steps=[PlanStep(file_path="main.py", action="modify", rationale="Fixture edit")], + risks=[], open_questions=[], estimated_complexity=Complexity.SMALL, + ) + + async def run_implementer( + self, plan: PlanArtifact, task_request: TaskRequest, worktree_path: str, language: str, + ) -> str: + self.implementations += 1 + workspace = Path(worktree_path) + (workspace / "main.py").write_text("after\n", encoding="utf-8") + (workspace / "new file.txt").write_text("added\n", encoding="utf-8") + (workspace / "new.bin").write_bytes(b"\0\xfffixture") + if self.fail == "execute": + raise RuntimeError("Synthetic partial execution error") + if self.fail == "commit": + git(workspace, "add", ".") + git(workspace, "commit", "-m", "Unexpected executor commit") + return "This is not the actual Git diff" + + async def run_reviewer( + self, diff: str, pr_title: str, pr_description: str, changed_files: list[str], + ) -> ReviewVerdict: + assert "new file.txt" in diff and "GIT binary patch" in diff + assert "This is not the actual Git diff" not in diff + if self.fail == "review": + raise RuntimeError("Synthetic review failure") + return ReviewVerdict( + verdict=ReviewVerdictType.APPROVE, issues=[], + summary="Inspected fixture patch", confidence=1.0, + ) + + +class DeterministicVerifier: + def __init__(self, mutate: bool = False) -> None: + self.mutate = mutate + + async def run_all( + self, workspace: str, changed_files: list[str], *, run_id: UUID, + ) -> tuple[list[Check], bool]: + root = Path(workspace) + assert (root / "main.py").read_text(encoding="utf-8") == "after\n" + assert (root / "new.bin").read_bytes() == b"\0\xfffixture" + if self.mutate: + (root / "main.py").write_text("changed during verification\n", encoding="utf-8") + return [Check()], True + + +def runtime( + root: Path, runner: DeterministicRunner, mutate: bool = False, +) -> FoundryTransitionRuntime: + source = repository(root) + return FoundryTransitionRuntime( + worktree_manager=WorktreeManager(str(source), str(root / "worktrees")), + agent_runner=runner, + verification_runner=DeterministicVerifier(mutate), + artifact_store=ArtifactStore(str(root / "artifacts")), + ) + + +def task() -> TaskRequest: + return TaskRequest( + task_type=TaskType.BUG_FIX, repo="fixture", title="Evidence retention", + prompt="Apply fixture edit", open_pr=False, + ) + + +async def test_full_adapter_evidence_is_retrievable_after_workspace_cleanup(tmp_path: Path) -> None: + identity = uuid4() + result = await runtime(tmp_path, DeterministicRunner()).execute_task( + task(), transition_id=identity, + ) + assert result.accepted + assert not (tmp_path / "worktrees" / str(identity)).exists() + journal = ArtifactStoreTransitionJournal(ArtifactStore(str(tmp_path / "artifacts"))) + reloaded = await journal.load(identity) + assert result == reloaded + refs = reloaded.observation.action_evidence + reloaded.observation.verification_evidence + assert len(refs) == 4 # Plan, patch, full check output, full review. + for ref in refs: + assert ref.uri.startswith("artifact://") + data = await journal.artifact_store.retrieve(ref.uri.removeprefix("artifact://")) + assert hashlib.sha256(data).hexdigest() == ref.checksum + patch = await journal.artifact_store.retrieve(f"runs/{identity}/diff.patch") + assert b"This is not the actual Git diff" not in patch + source = tmp_path / "source" + assert (source / "main.py").read_text(encoding="utf-8") == "before\n" # Not deployed. + git(source, "apply", "--binary", "-", data=patch) + assert (source / "main.py").read_text(encoding="utf-8") == "after\n" + assert (source / "new.bin").read_bytes() == b"\0\xfffixture" + + +@pytest.mark.parametrize("failure,phase", [("execute", "execute"), ("review", "verify")]) +async def test_real_partial_failure_keeps_workspace_and_records_failure( + tmp_path: Path, failure: str, phase: str, +) -> None: + identity = uuid4() + subject = runtime(tmp_path, DeterministicRunner(failure)) + with pytest.raises(RuntimeError): + await subject.execute_task(task(), transition_id=identity) + record = await ArtifactStoreTransitionJournal( + ArtifactStore(str(tmp_path / "artifacts")) + ).load(identity) + assert not record.accepted and record.metadata["phase"] == phase + retained = Path(record.metadata["retained_workspace"]) + assert (retained / "main.py").read_text(encoding="utf-8") == "after\n" + assert (tmp_path / "artifacts" / "runs" / str(identity) / "plan.json").exists() + if failure == "review": + assert (tmp_path / "artifacts" / "runs" / str(identity) / "diff.patch").exists() + + +async def test_mutating_verifier_cannot_leave_an_accepted_outcome(tmp_path: Path) -> None: + identity = uuid4() + with pytest.raises(StateConflictError, match="during verification"): + await runtime(tmp_path, DeterministicRunner(), mutate=True).execute_task( + task(), transition_id=identity, + ) + record = await ArtifactStoreTransitionJournal( + ArtifactStore(str(tmp_path / "artifacts")) + ).load(identity) + assert not record.accepted and record.metadata["phase"] == "observe_verified" + assert Path(record.metadata["retained_workspace"]).exists() + + +async def test_executor_commit_requires_inspection_not_silent_cleanup(tmp_path: Path) -> None: + identity = uuid4() + with pytest.raises(StateConflictError, match="Git HEAD"): + await runtime(tmp_path, DeterministicRunner("commit")).execute_task( + task(), transition_id=identity, + ) + record = await ArtifactStoreTransitionJournal( + ArtifactStore(str(tmp_path / "artifacts")) + ).load(identity) + assert not record.accepted and Path(record.metadata["retained_workspace"]).exists() + patch = await ArtifactStore(str(tmp_path / "artifacts")).retrieve(f"runs/{identity}/diff.patch") + assert b"GIT binary patch" in patch # Captured against the original observed HEAD. diff --git a/tests/unit/runtime/test_transition_integrity.py b/tests/unit/runtime/test_transition_integrity.py new file mode 100644 index 0000000..bd95bdb --- /dev/null +++ b/tests/unit/runtime/test_transition_integrity.py @@ -0,0 +1,287 @@ +"""Real-file journal tests; deterministic capabilities, no paid model calls.""" + +from __future__ import annotations + +import asyncio +import shutil +import subprocess +import sys +from datetime import UTC, datetime +from pathlib import Path +from uuid import UUID, uuid4 + +import pytest +from pydantic import ValidationError + +from foundry.contracts.transition_models import ( + ActionProposal, + EvidenceRef, + StateSnapshot, + TransitionOutcome, + TransitionRequest, + VerificationDecision, +) +from foundry.runtime.interfaces import TransitionJournal +from foundry.runtime.transition_engine import StateConflictError, TransitionEngine +from foundry.storage.artifact_store import ArtifactStore, ArtifactType +from foundry.storage.transition_journal import ArtifactStoreTransitionJournal + + +class FileWorld: + """An actual persisted counter; observing cannot fabricate progress.""" + + def __init__(self, root: Path, fail: str = "") -> None: + self.root = root + self.file = root / "world.json" + if not self.file.exists(): + self.file.write_text("0", encoding="utf-8") + self.fail = fail + self.observations = 0 + self.executions = 0 + self.cleanups = 0 + self.workspace: Path | None = None + + def value(self) -> int: + return int(self.file.read_text(encoding="utf-8")) + + async def prepare( + self, target: str, base_ref: str, transition_id: UUID, transition_name: str, + ) -> str: + if self.fail == "prepare": + raise RuntimeError("Synthetic prepare failure") + self.workspace = self.root / str(transition_id) + self.workspace.mkdir() + return str(self.workspace) + + async def observe_changes(self, workspace: str) -> str: + return str(self.value()) + + async def cleanup(self, workspace: str) -> None: + self.cleanups += 1 + if self.fail == "cleanup": + raise OSError("Synthetic cleanup failure") + shutil.rmtree(workspace) + + async def observe(self, request: TransitionRequest, workspace: str) -> StateSnapshot: + self.observations += 1 + phase = {1: "observe_before", 2: "observe_after", 3: "observe_verified"}.get( + self.observations + ) + if self.fail == phase: + raise RuntimeError("Synthetic observer failure") + return StateSnapshot( + environment="wrong" if self.fail == "wrong_environment" else request.environment, + observed_at=datetime.now(UTC), state={"value": self.value()}, + ) + + async def plan( + self, request: TransitionRequest, current_state: StateSnapshot, workspace: str, + ) -> ActionProposal: + if self.fail == "plan": + raise RuntimeError("Synthetic planner failure") + if self.fail == "mutate_snapshot": + current_state.state["value"] = 999 + return ActionProposal(kind="increment", description="Increment persisted counter") + + async def execute( + self, request: TransitionRequest, action: ActionProposal, workspace: str, + ) -> list[EvidenceRef]: + self.executions += 1 + self.file.write_text(str(self.value() + 1), encoding="utf-8") + if self.fail == "execute": + raise RuntimeError("Synthetic error text must not enter persisted journal") + if self.fail == "cancel": + raise asyncio.CancelledError() + return [] + + async def verify( + self, request: TransitionRequest, before: StateSnapshot, after: StateSnapshot, + action: ActionProposal, action_evidence: list[EvidenceRef], workspace: str, + ) -> VerificationDecision: + if self.fail == "verify": + raise RuntimeError("Synthetic verifier failure") + if self.fail == "mutate_during_verify": + self.file.write_text("100", encoding="utf-8") + return VerificationDecision( + accepted=self.fail != "reject" and after.state["value"] == before.state["value"] + 1, + reason="Compared actual persisted states", + ) + + +def engine(world: FileWorld, journal: TransitionJournal) -> TransitionEngine: + return TransitionEngine(world, world, world, world, world, journal) + + +def journal_at(root: Path) -> ArtifactStoreTransitionJournal: + return ArtifactStoreTransitionJournal(ArtifactStore(str(root / "evidence"))) + + +async def test_journal_reload_in_new_process_and_second_real_transition(tmp_path: Path) -> None: + request = TransitionRequest(environment="counter", objective="increment") + first = await engine(FileWorld(tmp_path), journal_at(tmp_path)).execute(request) + assert first.accepted and first.after_state.state == {"value": 1} + script = ( + "import asyncio,sys; from uuid import UUID; " + "from foundry.storage.artifact_store import ArtifactStore; " + "from foundry.storage.transition_journal import ArtifactStoreTransitionJournal; " + "j=ArtifactStoreTransitionJournal(ArtifactStore(sys.argv[1])); " + "print(asyncio.run(j.load(UUID(sys.argv[2]))).model_dump_json())" + ) + process = subprocess.run( + [sys.executable, "-c", script, str(tmp_path / "evidence"), str(request.id)], + check=True, capture_output=True, text=True, + ) + loaded = TransitionOutcome.model_validate_json(process.stdout) + second_request = TransitionRequest( + environment="counter", objective="increment again", before_state=loaded.after_state, + ) + second = await engine(FileWorld(tmp_path), journal_at(tmp_path)).execute(second_request) + assert second.accepted + assert second.before_state.state == {"value": 1} + assert second.after_state.state == {"value": 2} + assert (await journal_at(tmp_path).load(request.id)).after_state.state == {"value": 1} + + +@pytest.mark.parametrize("phase", [ + "prepare", "observe_before", "plan", "execute", "observe_after", "verify", "observe_verified", +]) +async def test_each_processing_failure_is_recorded_and_not_cleaned( + tmp_path: Path, phase: str, +) -> None: + world = FileWorld(tmp_path, phase) + request = TransitionRequest(environment="counter", objective="increment") + with pytest.raises(RuntimeError): + await engine(world, journal_at(tmp_path)).execute(request) + record = await journal_at(tmp_path).load(request.id) + assert record.accepted is False + assert record.metadata["status"] == "failed" + assert record.metadata["phase"] == phase + assert record.metadata["error_type"] == "RuntimeError" + assert world.cleanups == 0 + if phase != "prepare": + assert Path(record.metadata["retained_workspace"]).is_dir() + assert "Synthetic error text" not in record.model_dump_json() + + +async def test_cancellation_retains_partial_action_and_failure_record(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "cancel") + request = TransitionRequest(environment="counter", objective="increment") + with pytest.raises(asyncio.CancelledError): + await engine(world, journal_at(tmp_path)).execute(request) + record = await journal_at(tmp_path).load(request.id) + assert record.metadata["status"] == "cancelled" + assert not record.accepted and world.value() == 1 and world.cleanups == 0 + + +async def test_stale_state_prevents_an_action(tmp_path: Path) -> None: + world = FileWorld(tmp_path) + stale = StateSnapshot(environment="counter", observed_at=datetime.now(UTC), state={"value": 7}) + request = TransitionRequest(environment="counter", objective="increment", before_state=stale) + with pytest.raises(StateConflictError, match="stale"): + await engine(world, journal_at(tmp_path)).execute(request) + assert world.executions == 0 and world.value() == 0 + assert not (await journal_at(tmp_path).load(request.id)).accepted + + +async def test_wrong_environment_cannot_supply_state(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "wrong_environment") + request = TransitionRequest(environment="counter", objective="increment") + with pytest.raises(StateConflictError, match="different environment"): + await engine(world, journal_at(tmp_path)).execute(request) + assert world.executions == 0 + + +async def test_verifier_mutation_invalidates_the_decision(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "mutate_during_verify") + request = TransitionRequest(environment="counter", objective="increment") + with pytest.raises(StateConflictError, match="during verification"): + await engine(world, journal_at(tmp_path)).execute(request) + assert world.value() == 100 and world.cleanups == 0 + assert not (await journal_at(tmp_path).load(request.id)).accepted + + +async def test_planner_cannot_mutate_recorded_before_snapshot(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "mutate_snapshot") + result = await engine(world, journal_at(tmp_path)).execute( + TransitionRequest(environment="counter", objective="increment") + ) + assert result.accepted and result.before_state.state == {"value": 0} + + +async def test_rejection_is_not_exception_and_does_not_claim_rollback(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "reject") + request = TransitionRequest(environment="counter", objective="increment") + result = await engine(world, journal_at(tmp_path)).execute(request) + assert not result.accepted and result.metadata["status"] == "rejected" + assert world.value() == 1 # Rejected actions can still change reality. + assert world.cleanups == 1 + assert (await journal_at(tmp_path).load(request.id)) == result + + +async def test_journal_failure_cannot_trigger_cleanup(tmp_path: Path) -> None: + class BrokenJournal: + async def record(self, outcome: TransitionOutcome) -> None: + raise OSError("Synthetic persistence failure") + + world = FileWorld(tmp_path) + with pytest.raises(OSError, match="persistence"): + await engine(world, BrokenJournal()).execute( + TransitionRequest(environment="counter", objective="increment") + ) + assert world.cleanups == 0 and world.workspace.exists() + + +async def test_primary_error_survives_secondary_journal_error(tmp_path: Path) -> None: + class BrokenJournal: + async def record(self, outcome: TransitionOutcome) -> None: + raise OSError("Synthetic persistence failure") + + world = FileWorld(tmp_path, "execute") + with pytest.raises(RuntimeError) as caught: + await engine(world, BrokenJournal()).execute( + TransitionRequest(environment="counter", objective="increment") + ) + assert any("Failure journal also failed: OSError" in note for note in caught.value.__notes__) + assert world.cleanups == 0 + + +async def test_cleanup_failure_preserves_decision_and_records_warning(tmp_path: Path) -> None: + world = FileWorld(tmp_path, "cleanup") + request = TransitionRequest(environment="counter", objective="increment") + result = await engine(world, journal_at(tmp_path)).execute(request) + assert result.accepted and result.metadata["cleanup_status"] == "failed" + loaded = await journal_at(tmp_path).load(request.id) + assert loaded == result and Path(loaded.metadata["retained_workspace"]).exists() + + +async def test_atomic_write_failure_preserves_previous_bytes( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ArtifactStore(str(tmp_path / "artifacts")) + identity = uuid4() + old = await store.store(identity, ArtifactType.TRANSITION, b"previous") + + def refuse_replace(*args: object) -> None: + raise OSError("Synthetic replace failure") + + monkeypatch.setattr("foundry.storage.artifact_store.os.replace", refuse_replace) + with pytest.raises(OSError): + await store.store(identity, ArtifactType.TRANSITION, b"replacement") + assert await store.retrieve(old["storage_path"]) == b"previous" + assert len(await store.list_artifacts(identity)) == 1 + + +async def test_journal_refuses_corrupt_or_wrong_identity(tmp_path: Path) -> None: + journal = journal_at(tmp_path) + identity = uuid4() + await journal.artifact_store.store( + identity, ArtifactType.TRANSITION, "{", "transition_outcome.json", + ) + with pytest.raises(ValidationError): + await journal.load(identity) + other = TransitionOutcome(transition_id=uuid4(), accepted=False) + await journal.artifact_store.store( + identity, ArtifactType.TRANSITION, other.model_dump_json(), "transition_outcome.json" + ) + with pytest.raises(ValueError, match="different transition ID"): + await journal.load(identity)