Sfoglia il codice sorgente

测试:覆盖动态工具、阶段续跑、边界收口与接口安全

SamLee 2 giorni fa
parent
commit
fc025046c6

+ 219 - 0
script_build_host/tests/test_agent_surface.py

@@ -1,6 +1,8 @@
 from __future__ import annotations
 
+from dataclasses import replace
 from datetime import UTC, datetime
+from hashlib import sha256
 from types import SimpleNamespace
 
 import pytest
@@ -11,7 +13,12 @@ from script_build_host.agents.model_resolver import (
     SnapshotModelManifestSource,
 )
 from script_build_host.agents.presets import register_script_presets
+from script_build_host.agents.prompt_resolver import SnapshotRoleSystemPromptResolver
+from script_build_host.agents.prompts import script_build_prompt_manifest
+from script_build_host.domain.canonical_json import canonical_sha256
+from script_build_host.domain.errors import ProtocolViolation
 from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
+from script_build_host.tools.gateway import LegacyScriptToolGateway
 from script_build_host.tools.registry import register_script_tools
 
 
@@ -26,6 +33,12 @@ def test_script_presets_are_exact_and_never_expose_legacy_control_tools() -> Non
         "script_knowledge_retrieval_worker",
         "script_retrieval_validator",
         "script_candidate_validator",
+        "script_structure_worker",
+        "script_paragraph_worker",
+        "script_element_set_worker",
+        "script_candidate_compare_worker",
+        "script_compose_worker",
+        "script_candidate_portfolio_worker",
     }
     forbidden = {
         "agent",
@@ -43,6 +56,46 @@ def test_script_presets_are_exact_and_never_expose_legacy_control_tools() -> Non
         assert forbidden.isdisjoint(set(preset.allowed_tools or []))
         assert set(preset.allowed_tools or []).isdisjoint(set(preset.denied_tools or []))
 
+    planner = get_preset("script_planner")
+    assert set(planner.allowed_tools or []) == {
+        "plan_script_tasks",
+        "decide_script_task",
+        "inspect_script_plan",
+        "dispatch_script_tasks",
+        "validate_attempt",
+        "read_input_snapshot",
+    }
+    compose = get_preset("script_compose_worker")
+    assert "read_active_frontier" in (compose.allowed_tools or [])
+    assert "read_accepted_artifact" not in (compose.allowed_tools or [])
+    compare = get_preset("script_candidate_compare_worker")
+    assert "read_pinned_candidates" in (compare.allowed_tools or [])
+    assert "create_script_paragraph" not in (compare.allowed_tools or [])
+    portfolio = get_preset("script_candidate_portfolio_worker")
+    assert "save_candidate_portfolio" in (portfolio.allowed_tools or [])
+    assert "save_structured_script_candidate" not in (portfolio.allowed_tools or [])
+    for validator_name in ("script_retrieval_validator", "script_candidate_validator"):
+        assert "view_frozen_images" in (get_preset(validator_name).allowed_tools or [])
+
+
+def test_new_build_prompt_manifest_freezes_every_phase_one_and_phase_two_role() -> None:
+    manifest = script_build_prompt_manifest()
+    assert len(manifest) == 14
+    assert len({str(item["preset"]) for item in manifest}) == len(manifest)
+    assert {str(item["preset"]) for item in manifest} >= {
+        "script_planner",
+        "script_structure_worker",
+        "script_paragraph_worker",
+        "script_element_set_worker",
+        "script_candidate_compare_worker",
+        "script_compose_worker",
+        "script_candidate_portfolio_worker",
+    }
+    for item in manifest:
+        assert str(item["content"])
+        assert str(item["content_sha256"]).startswith("sha256:")
+        assert len(str(item["content_sha256"])) == 71
+
 
 class _Gateway:
     coordinator = SimpleNamespace(task_store=None)
@@ -66,6 +119,20 @@ def test_legacy_retrieval_tool_schemas_keep_old_argument_shapes() -> None:
     assert knowledge["properties"]["max_count"]["default"] == 3
     assert schemas["query_pattern_qa"]["required"] == ["message"]
     assert schemas["load_images"]["required"] == ["image_urls"]
+    assert schemas["load_frozen_strategy"]["required"] == ["strategy_ref"]
+    assert schemas["submit_attempt"].get("properties") == {}
+    assert schemas["submit_validation"]["required"] == [
+        "verdict",
+        "criterion_results",
+        "defects",
+    ]
+    assert "summary" not in schemas["submit_attempt"].get("properties", {})
+    assert "artifact_refs" not in schemas["submit_attempt"].get("properties", {})
+    assert {item.value for item in registry.get_capabilities("search_script_decode_case")} == {
+        "external_send",
+        "read",
+        "write",
+    }
 
 
 class _Bindings:
@@ -117,3 +184,155 @@ async def test_role_model_resolver_reads_each_roots_frozen_manifest() -> None:
     assert value.model == "frozen-direction-model"
     assert value.temperature == 0.2
     assert value.max_iterations == 18
+
+
+@pytest.mark.asyncio
+async def test_role_prompt_resolver_uses_task_pinned_snapshot_and_rechecks_digest() -> None:
+    content = "Frozen paragraph worker policy"
+    snapshot = ScriptBuildInputSnapshotV1(
+        snapshot_id="3",
+        script_build_id=9,
+        execution_id=1,
+        topic_build_id=2,
+        topic_id=3,
+        topic={},
+        account={},
+        persona_points=(),
+        section_patterns=(),
+        strategies=(),
+        prompt_manifest=(
+            {
+                "preset": "script_paragraph_worker",
+                "role": "worker",
+                "content": content,
+                "content_sha256": f"sha256:{sha256(content.encode('utf-8')).hexdigest()}",
+            },
+        ),
+        datasource_manifest={},
+        model_manifest={},
+        canonical_sha256="sha256:" + "a" * 64,
+        created_at=datetime.now(UTC),
+    )
+
+    class Snapshots:
+        async def get(self, snapshot_id: str, *, script_build_id: int):
+            assert (snapshot_id, script_build_id) == ("3", 9)
+            return snapshot
+
+    resolver = SnapshotRoleSystemPromptResolver(_Bindings(), Snapshots())
+    value = await resolver.resolve(
+        role=AgentRole.WORKER,
+        preset="script_paragraph_worker",
+        context={
+            "root_trace_id": "root-a",
+            "task_spec": {"context_refs": ["script-build://inputs/3"]},
+        },
+    )
+    assert value is not None
+    assert value.content == content
+
+    tampered = replace(
+        snapshot,
+        prompt_manifest=(
+            {
+                **snapshot.prompt_manifest[0],
+                "content_sha256": "sha256:" + "f" * 64,
+            },
+        ),
+    )
+
+    class TamperedSnapshots:
+        async def get(self, _snapshot_id: str, *, script_build_id: int):
+            assert script_build_id == 9
+            return tampered
+
+    with pytest.raises(ProtocolViolation, match="digest"):
+        await SnapshotRoleSystemPromptResolver(_Bindings(), TamperedSnapshots()).resolve(
+            role=AgentRole.WORKER,
+            preset="script_paragraph_worker",
+            context={
+                "root_trace_id": "root-a",
+                "task_spec": {"context_refs": ["script-build://inputs/3"]},
+            },
+        )
+
+
+@pytest.mark.asyncio
+async def test_snapshot_tool_hides_on_demand_body_until_explicit_digest_verified_load() -> None:
+    strategy_ref = "script-build://strategies/8/versions/2"
+    snapshot = ScriptBuildInputSnapshotV1(
+        snapshot_id="3",
+        script_build_id=9,
+        execution_id=1,
+        topic_build_id=2,
+        topic_id=3,
+        topic={},
+        account={},
+        persona_points=(),
+        section_patterns=(),
+        strategies=(
+            {
+                "strategy_id": 7,
+                "version": 1,
+                "mode": "always_on",
+                "content": "always body",
+                "content_sha256": "sha256:" + "1" * 64,
+            },
+            {
+                "strategy_id": 8,
+                "version": 2,
+                "mode": "on_demand",
+                "description": "optional",
+                "content": "on demand body",
+                "content_sha256": canonical_sha256("on demand body").wire,
+            },
+        ),
+        prompt_manifest=(),
+        datasource_manifest={},
+        model_manifest={},
+        canonical_sha256="sha256:" + "a" * 64,
+        created_at=datetime.now(UTC),
+    )
+
+    class Snapshots:
+        async def get(self, _snapshot_id: str, *, script_build_id: int):
+            assert script_build_id == 9
+            return snapshot
+
+    task = SimpleNamespace(current_spec=SimpleNamespace(context_refs=["script-build://inputs/3"]))
+    coordinator = SimpleNamespace(
+        task_store=SimpleNamespace(load=lambda _root: _async(SimpleNamespace(tasks={"task": task})))
+    )
+    gateway = LegacyScriptToolGateway(
+        bindings=_Bindings(),
+        snapshots=Snapshots(),
+        artifacts=SimpleNamespace(),
+        coordinator=coordinator,
+    )
+    context = {"root_trace_id": "root-a", "task_id": "task"}
+    summary = await gateway.read_input_snapshot(context)
+    assert summary["strategies"][0]["content"] == "always body"
+    assert "content" not in summary["strategies"][1]
+    loaded = await gateway.load_frozen_strategy(strategy_ref, context)
+    assert loaded["content"] == "on demand body"
+
+    tampered = replace(
+        snapshot,
+        strategies=(
+            snapshot.strategies[0],
+            {**snapshot.strategies[1], "content": "changed after freeze"},
+        ),
+    )
+
+    class TamperedSnapshots:
+        async def get(self, _snapshot_id: str, *, script_build_id: int):
+            assert script_build_id == 9
+            return tampered
+
+    gateway.snapshots = TamperedSnapshots()
+    with pytest.raises(ProtocolViolation, match="strategy digest"):
+        await gateway.load_frozen_strategy(strategy_ref, context)
+
+
+async def _async(value):
+    return value

+ 192 - 8
script_build_host/tests/test_api_security.py

@@ -1,9 +1,9 @@
 from __future__ import annotations
 
 import json
+from datetime import UTC, datetime
 from pathlib import Path
 from types import SimpleNamespace
-from typing import ClassVar
 
 import httpx
 import pytest
@@ -11,7 +11,17 @@ from fastapi.testclient import TestClient
 from starlette.websockets import WebSocketDisconnect
 
 from script_build_host.api import ApiSecurity, create_app
-from script_build_host.application.mission_service import ScriptMissionStartResult
+from script_build_host.application.mission_service import (
+    PhaseAdvanceResult,
+    ScriptMissionStartResult,
+)
+from script_build_host.domain.artifacts import (
+    ArtifactKind,
+    ArtifactState,
+    ArtifactVersion,
+    EvidenceRecordV1,
+)
+from script_build_host.domain.errors import ArtifactNotFound
 from script_build_host.domain.records import BuildStatus, Principal
 from script_build_host.infrastructure.outbound import OutboundPolicy
 
@@ -44,6 +54,7 @@ class _Authorizer:
 class _MissionService:
     def __init__(self) -> None:
         self.command = None
+        self.advanced = False
 
     async def start(self, command: object) -> ScriptMissionStartResult:
         self.command = command
@@ -57,6 +68,13 @@ class _MissionService:
             stopped_operation_ids=(),
         )
 
+    async def advance_to_phase_two(
+        self, script_build_id: int, principal: Principal
+    ) -> PhaseAdvanceResult:
+        assert principal.subject == "owner"
+        self.advanced = True
+        return PhaseAdvanceResult(script_build_id, BuildStatus.RUNNING, "root-7", "11")
+
 
 class _Bindings:
     async def get_by_build(self, script_build_id: int) -> object:
@@ -64,9 +82,40 @@ class _Bindings:
 
 
 class _TaskStore:
+    def __init__(self) -> None:
+        criterion = SimpleNamespace(criterion_id="criterion-1", description="safe", hard=True)
+        current_spec = SimpleNamespace(
+            version=1,
+            objective="bounded task",
+            acceptance_criteria=(criterion,),
+            context_refs=("script-build://task-kinds/paragraph",),
+        )
+        self.task = SimpleNamespace(
+            task_id="task-7",
+            parent_task_id="root-task",
+            display_path="Root/task-7",
+            status="completed",
+            current_spec=current_spec,
+            child_task_ids=(),
+            attempt_ids=(),
+            validation_ids=(),
+            decision_ids=(),
+            blocked_reason=None,
+            superseded_by=None,
+            created_at="2026-07-19T00:00:00Z",
+            updated_at="2026-07-19T00:00:00Z",
+        )
+
     async def load(self, root_trace_id: str) -> object:
         del root_trace_id
-        return SimpleNamespace(to_dict=lambda: {"root_trace_id": "root-7"})
+        return SimpleNamespace(
+            tasks={"task-7": self.task},
+            to_dict=lambda: {
+                "root_trace_id": "root-7",
+                "protected_context": {"token": "must-not-leak"},
+                "command_records": [{"arguments": "secret"}],
+            },
+        )
 
 
 def _app(
@@ -74,6 +123,7 @@ def _app(
     *,
     trace_store: object | None = None,
     uploaded_topics: object | None = None,
+    business_artifacts: object | None = None,
 ) -> tuple[object, _MissionService]:
     mission = _MissionService()
     app = create_app(
@@ -84,7 +134,7 @@ def _app(
             websocket_allowed_origins=("https://ui.example",),
         ),
         bindings=_Bindings(),
-        business_artifacts=SimpleNamespace(),
+        business_artifacts=business_artifacts or SimpleNamespace(),
         publications=SimpleNamespace(),
         coordinator=SimpleNamespace(task_store=_TaskStore()),
         trace_store=trace_store or SimpleNamespace(),
@@ -138,7 +188,8 @@ async def test_start_requires_auth_and_preserves_old_response_fields() -> None:
         assert body["script_build_id"] == 7
         assert body["status"] == "running"
         assert mission.command.strategies_always_on == (13,)
-        assert len(mission.command.runtime_prompt_manifest) == 8
+        assert len(mission.command.prompt_requests) >= 8
+        assert mission.command.runtime_prompt_manifest == ()
 
         secret_config = dict(legacy_request)
         secret_config["agent_config"] = {
@@ -180,6 +231,8 @@ async def test_observation_is_build_scoped_and_control_routes_are_not_mounted()
     ) as client:
         allowed = await client.get("/api/pattern/script_builds/7/mission")
         assert allowed.status_code == 200
+        assert allowed.json() == {"root_trace_id": "root-7"}
+        assert "must-not-leak" not in allowed.text
         denied = await client.get("/api/pattern/script_builds/8/mission")
         assert denied.status_code == 404
         unsafe = await client.post(
@@ -204,6 +257,40 @@ async def test_observation_is_build_scoped_and_control_routes_are_not_mounted()
         assert secret_query.status_code == 400
 
 
+@pytest.mark.asyncio
+async def test_phase_two_advance_auth_and_hidden_not_found() -> None:
+    unauthenticated, _ = _app(None)
+    async with httpx.AsyncClient(
+        transport=httpx.ASGITransport(app=unauthenticated), base_url="http://test"
+    ) as client:
+        response = await client.post("/api/pattern/script_builds/7/phase-two/advance")
+        assert response.status_code == 401
+
+    forbidden, forbidden_mission = _app(Principal("intruder"))
+    async with httpx.AsyncClient(
+        transport=httpx.ASGITransport(app=forbidden), base_url="http://test"
+    ) as client:
+        response = await client.post("/api/pattern/script_builds/7/phase-two/advance")
+        assert response.status_code == 404
+        assert response.json()["detail"]["error_code"] == "BUILD_NOT_FOUND"
+        assert forbidden_mission.advanced is False
+
+    allowed, mission = _app(Principal("owner"))
+    async with httpx.AsyncClient(
+        transport=httpx.ASGITransport(app=allowed), base_url="http://test"
+    ) as client:
+        response = await client.post("/api/pattern/script_builds/7/phase-two/advance")
+        assert response.status_code == 200
+        assert response.json() == {
+            "success": True,
+            "script_build_id": 7,
+            "status": "running",
+            "root_trace_id": "root-7",
+            "input_snapshot_id": "11",
+        }
+        assert mission.advanced is True
+
+
 def test_websocket_rejects_unauthenticated_subscription_before_accept() -> None:
     app, _ = _app(None)
     with TestClient(app) as client:
@@ -214,8 +301,9 @@ def test_websocket_rejects_unauthenticated_subscription_before_accept() -> None:
 
 
 class _Trace:
-    trace_id = "root-7"
-    context: ClassVar[dict[str, str]] = {}
+    def __init__(self, trace_id: str, root_trace_id: str | None = None) -> None:
+        self.trace_id = trace_id
+        self.context = {"root_trace_id": root_trace_id} if root_trace_id is not None else {}
 
     def to_dict(self) -> dict[str, str]:
         return {"trace_id": self.trace_id}
@@ -223,7 +311,16 @@ class _Trace:
 
 class _TraceStore:
     async def get_trace(self, trace_id: str) -> _Trace | None:
-        return _Trace() if trace_id == "root-7" else None
+        values = {
+            "root-7": _Trace("root-7"),
+            "worker-7": _Trace("worker-7", "root-7"),
+            "worker-8": _Trace("worker-8", "root-8"),
+        }
+        return values.get(trace_id)
+
+    async def get_goal_tree(self, trace_id: str) -> None:
+        del trace_id
+        return None
 
     async def get_events(self, trace_id: str, since: int) -> list[dict[str, object]]:
         assert trace_id == "root-7"
@@ -231,6 +328,43 @@ class _TraceStore:
         return []
 
 
+class _BusinessArtifacts:
+    def __init__(self) -> None:
+        now = datetime.now(UTC)
+        evidence = EvidenceRecordV1(
+            evidence_id="evidence-api-1",
+            source_type="decode",
+            tool_name="retrieve_decode",
+            query={"query": "safe query"},
+            source_refs=("decode:item:1",),
+            raw_artifact_ref=None,
+            summary="authorized phase-two artifact content",
+            supports=("scope:opening",),
+            confidence="high",
+            limitations=(),
+            content_sha256="sha256:" + "a" * 64,
+            created_at=now,
+        )
+        self.value = ArtifactVersion(
+            artifact_version_id=71,
+            script_build_id=7,
+            task_id="task-7",
+            attempt_id="attempt-7",
+            spec_version=1,
+            artifact_type=ArtifactKind.EVIDENCE,
+            canonical_sha256="sha256:" + "b" * 64,
+            state=ArtifactState.FROZEN,
+            artifact=evidence,
+            created_at=now,
+            frozen_at=now,
+        )
+
+    async def get_by_id(self, artifact_version_id: int, *, script_build_id: int) -> ArtifactVersion:
+        if artifact_version_id != 71 or script_build_id != 7:
+            raise ArtifactNotFound()
+        return self.value
+
+
 def test_websocket_requires_origin_and_accepts_authorized_build_trace() -> None:
     app, _ = _app(Principal("owner"), trace_store=_TraceStore())
     with TestClient(app) as client:
@@ -249,6 +383,56 @@ def test_websocket_requires_origin_and_accepts_authorized_build_trace() -> None:
             websocket.send_text("ping")
             assert websocket.receive_json() == {"event": "pong"}
 
+        for foreign_trace_id in ("worker-8", "root-8"):
+            with pytest.raises(WebSocketDisconnect) as caught:
+                with client.websocket_connect(
+                    f"/api/pattern/script_builds/7/traces/{foreign_trace_id}/watch",
+                    headers={"origin": "https://ui.example"},
+                ):
+                    pass
+            assert caught.value.code == 4404
+
+
+@pytest.mark.asyncio
+async def test_phase_two_observation_hides_foreign_resources_and_returns_owned_artifact() -> None:
+    app, _ = _app(
+        Principal("owner"),
+        trace_store=_TraceStore(),
+        business_artifacts=_BusinessArtifacts(),
+    )
+    async with httpx.AsyncClient(
+        transport=httpx.ASGITransport(app=app), base_url="http://test"
+    ) as client:
+        owned_task = await client.get("/api/pattern/script_builds/7/tasks/task-7")
+        assert owned_task.status_code == 200
+        assert owned_task.json()["task_id"] == "task-7"
+
+        foreign_task = await client.get("/api/pattern/script_builds/7/tasks/task-8")
+        assert foreign_task.status_code == 404
+        assert foreign_task.json()["detail"]["error_code"] == "TASK_NOT_FOUND"
+
+        owned_artifact = await client.get("/api/pattern/script_builds/7/artifacts/71")
+        assert owned_artifact.status_code == 200
+        assert (
+            owned_artifact.json()["artifact"]["summary"] == "authorized phase-two artifact content"
+        )
+
+        foreign_artifact = await client.get("/api/pattern/script_builds/7/artifacts/81")
+        assert foreign_artifact.status_code == 404
+        assert foreign_artifact.json()["error_code"] == "ARTIFACT_NOT_FOUND"
+
+        owned_trace = await client.get("/api/pattern/script_builds/7/traces/worker-7")
+        assert owned_trace.status_code == 200
+        assert owned_trace.json()["trace"]["trace_id"] == "worker-7"
+
+        foreign_trace = await client.get("/api/pattern/script_builds/7/traces/worker-8")
+        assert foreign_trace.status_code == 404
+        assert foreign_trace.json()["detail"]["error_code"] == "TRACE_NOT_FOUND"
+
+        forged_root = await client.get("/api/v2/script-builds/7/roots/root-8/snapshot")
+        assert forged_root.status_code == 404
+        assert forged_root.json()["detail"]["error_code"] == "ROOT_NOT_FOUND"
+
 
 class _UploadedTopics:
     async def parse(self, _value: object) -> dict[str, object]:

+ 352 - 2
script_build_host/tests/test_mission_service_boundaries.py

@@ -4,6 +4,13 @@ import asyncio
 from types import SimpleNamespace
 
 import pytest
+from agent.orchestration import (
+    ArtifactRef,
+    DecisionAction,
+    OperationStatus,
+    TaskStatus,
+    ValidationVerdict,
+)
 
 from script_build_host.application.mission_service import (
     BuildTransitionGate,
@@ -11,8 +18,18 @@ from script_build_host.application.mission_service import (
     ScriptMissionService,
     StartScriptBuildCommand,
 )
-from script_build_host.domain.errors import InputRelationMismatch, ProtocolViolation
-from script_build_host.domain.records import BuildStatus, Principal
+from script_build_host.domain.artifacts import (
+    ArtifactKind,
+    ArtifactState,
+    DirectionGoal,
+    ScriptDirectionArtifactV1,
+)
+from script_build_host.domain.errors import (
+    DirectionProjectionConflict,
+    InputRelationMismatch,
+    ProtocolViolation,
+)
+from script_build_host.domain.records import BuildStatus, Principal, PublicationState
 
 
 class _RejectingInputs:
@@ -162,6 +179,227 @@ async def test_stop_is_idempotent_and_publication_is_gated_by_durable_status() -
     assert publications.prepared is False
 
 
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("publication_state", "projected", "active"),
+    [
+        (PublicationState.PENDING, "# direction", None),
+        (PublicationState.FAILED, "# direction", 3),
+        (PublicationState.PUBLISHED, "# direction", 3),
+    ],
+)
+async def test_direction_reconciler_reenters_each_publication_crash_point(
+    publication_state: PublicationState,
+    projected: str | None,
+    active: int | None,
+) -> None:
+    ref = ArtifactRef(
+        "script-build://artifact-versions/3",
+        ArtifactKind.DIRECTION.value,
+        "3",
+        "sha256:" + "a" * 64,
+    )
+    direction = ScriptDirectionArtifactV1(
+        goals=(DirectionGoal("goal", "statement"),),
+        evidence_refs=("script-build://artifact-versions/2",),
+        legacy_markdown="# direction",
+    )
+    version = SimpleNamespace(
+        artifact_version_id=3,
+        canonical_sha256=ref.digest,
+        state=(
+            ArtifactState.PUBLISHED
+            if publication_state is PublicationState.PUBLISHED
+            else ArtifactState.FROZEN
+        ),
+        artifact=direction,
+    )
+    task = SimpleNamespace(
+        task_id="direction-task",
+        current_spec=SimpleNamespace(context_refs=("script-build://task-kinds/direction",)),
+        status=TaskStatus.COMPLETED,
+        decision_ids=["accept"],
+    )
+    attempt = SimpleNamespace(
+        attempt_id="attempt",
+        snapshot_id="snapshot",
+        submission=SimpleNamespace(artifact_refs=(ref,)),
+    )
+    ledger = SimpleNamespace(
+        tasks={task.task_id: task},
+        decisions={
+            "accept": SimpleNamespace(
+                action=DecisionAction.ACCEPT,
+                decision_id="accept",
+                attempt_id="attempt",
+                validation_id="validation",
+            )
+        },
+        attempts={"attempt": attempt},
+        validations={
+            "validation": SimpleNamespace(
+                attempt_id="attempt",
+                verdict=ValidationVerdict.PASSED,
+                snapshot_id="snapshot",
+            )
+        },
+    )
+
+    class Bindings:
+        current = active
+
+        async def get_by_build(self, _build: int):
+            return SimpleNamespace(active_direction_artifact_version_id=self.current)
+
+        async def set_active_direction(self, *, artifact_version_id: int, **_kwargs):
+            self.current = artifact_version_id
+
+    class Publications:
+        state = publication_state
+        published_calls = 0
+
+        async def prepare(self, **_kwargs):
+            return SimpleNamespace(publication_id=1, state=self.state)
+
+        async def mark_published(self, _publication_id: int):
+            self.state = PublicationState.PUBLISHED
+            self.published_calls += 1
+
+        async def mark_failed(self, *_args, **_kwargs):
+            self.state = PublicationState.FAILED
+
+    class State(_StopState):
+        direction = projected
+
+        async def get_direction(self, _build: int):
+            return self.direction
+
+        async def project_direction(self, _build: int, value: str):
+            self.direction = value
+
+    bindings = Bindings()
+    publications = Publications()
+    state = State(BuildStatus.PARTIAL)
+    reconciler = DirectionReconciler(
+        coordinator=SimpleNamespace(
+            task_store=SimpleNamespace(load=lambda _root: _async_value(ledger))
+        ),
+        bindings=bindings,
+        artifacts=SimpleNamespace(read_by_ref=lambda *_args, **_kwargs: _async_value(version)),
+        publications=publications,
+        legacy_state=state,
+    )
+
+    assert await reconciler.reconcile(7, "root") == 3
+    assert state.direction == "# direction"
+    assert bindings.current == 3
+    assert publications.state is PublicationState.PUBLISHED
+    assert publications.published_calls == (
+        0 if publication_state is PublicationState.PUBLISHED else 1
+    )
+
+
+@pytest.mark.asyncio
+async def test_direction_reconciler_fails_closed_on_projection_conflict() -> None:
+    ref = ArtifactRef(
+        "script-build://artifact-versions/3",
+        ArtifactKind.DIRECTION.value,
+        "3",
+        "sha256:" + "a" * 64,
+    )
+    direction = ScriptDirectionArtifactV1(
+        goals=(DirectionGoal("goal", "statement"),),
+        evidence_refs=("script-build://artifact-versions/2",),
+        legacy_markdown="# accepted direction",
+    )
+    task = SimpleNamespace(
+        task_id="direction-task",
+        current_spec=SimpleNamespace(context_refs=("script-build://task-kinds/direction",)),
+        status=TaskStatus.COMPLETED,
+        decision_ids=["accept"],
+    )
+    ledger = SimpleNamespace(
+        tasks={task.task_id: task},
+        decisions={
+            "accept": SimpleNamespace(
+                action=DecisionAction.ACCEPT,
+                decision_id="accept",
+                attempt_id="attempt",
+                validation_id="validation",
+            )
+        },
+        attempts={
+            "attempt": SimpleNamespace(
+                attempt_id="attempt",
+                snapshot_id="snapshot",
+                submission=SimpleNamespace(artifact_refs=(ref,)),
+            )
+        },
+        validations={
+            "validation": SimpleNamespace(
+                attempt_id="attempt",
+                verdict=ValidationVerdict.PASSED,
+                snapshot_id="snapshot",
+            )
+        },
+    )
+
+    class Publications:
+        failed = False
+
+        async def prepare(self, **_kwargs):
+            return SimpleNamespace(publication_id=1, state=PublicationState.PENDING)
+
+        async def mark_failed(self, *_args, **_kwargs):
+            self.failed = True
+
+    publications = Publications()
+    reconciler = DirectionReconciler(
+        coordinator=SimpleNamespace(
+            task_store=SimpleNamespace(load=lambda _root: _async_value(ledger))
+        ),
+        bindings=SimpleNamespace(),
+        artifacts=SimpleNamespace(
+            read_by_ref=lambda *_args, **_kwargs: _async_value(
+                SimpleNamespace(
+                    artifact_version_id=3,
+                    canonical_sha256=ref.digest,
+                    state=ArtifactState.FROZEN,
+                    artifact=direction,
+                )
+            )
+        ),
+        publications=publications,
+        legacy_state=SimpleNamespace(
+            get_status=lambda _build: _async_value(BuildStatus.PARTIAL),
+            get_direction=lambda _build: _async_value("# conflicting direction"),
+        ),
+    )
+
+    with pytest.raises(DirectionProjectionConflict):
+        await reconciler.reconcile(7, "root")
+    assert publications.failed is True
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("status", [BuildStatus.FAILED, BuildStatus.SUCCESS])
+async def test_stop_does_not_rewrite_terminal_failed_or_success(status: BuildStatus) -> None:
+    state = _StopState(status)
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+    result = await service.stop(7, Principal("owner"))
+    assert result.status is status
+    assert state.status is status
+
+
 @pytest.mark.asyncio
 async def test_stop_intent_and_direction_reconcile_are_serialized() -> None:
     gate = BuildTransitionGate()
@@ -209,5 +447,117 @@ async def test_stop_intent_and_direction_reconcile_are_serialized() -> None:
     assert (await stop).status is BuildStatus.STOPPED
 
 
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    "initial_status",
+    [BuildStatus.RUNNING, BuildStatus.STOPPING, BuildStatus.PARTIAL],
+)
+async def test_stop_closes_all_nonterminal_states_when_no_execution_remains(
+    initial_status: BuildStatus,
+) -> None:
+    state = _StopState(initial_status)
+    service = ScriptMissionService(
+        runner=SimpleNamespace(
+            stop=lambda _root: _async_value(True),
+            trace_store=SimpleNamespace(get_trace=lambda _root: _async_value(None)),
+        ),
+        coordinator=SimpleNamespace(
+            task_store=SimpleNamespace(
+                load=lambda _root: _async_value(SimpleNamespace(operations={}))
+            )
+        ),
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(
+            get_by_build=lambda _build: _async_value(SimpleNamespace(root_trace_id="root"))
+        ),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    result = await service.stop(7, Principal("owner"))
+
+    assert result.status is BuildStatus.STOPPED
+    assert state.status is BuildStatus.STOPPED
+
+
+@pytest.mark.asyncio
+async def test_stop_timeout_keeps_intent_and_reports_active_operation_ids() -> None:
+    state = _StopState(BuildStatus.RUNNING)
+    operation = SimpleNamespace(
+        operation_id="operation-1",
+        status=OperationStatus.RUNNING,
+    )
+    service = ScriptMissionService(
+        runner=SimpleNamespace(
+            stop=lambda _root: _async_value(True),
+            trace_store=SimpleNamespace(get_trace=lambda _root: _async_value(None)),
+        ),
+        coordinator=SimpleNamespace(
+            task_store=SimpleNamespace(
+                load=lambda _root: _async_value(
+                    SimpleNamespace(operations={"operation-1": operation})
+                )
+            ),
+            stop_operation=lambda *_args, **_kwargs: _async_value(None),
+            get_operation=lambda *_args, **_kwargs: _async_value(operation),
+        ),
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(
+            get_by_build=lambda _build: _async_value(SimpleNamespace(root_trace_id="root"))
+        ),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+        stop_timeout_seconds=0.01,
+    )
+
+    result = await service.stop(7, Principal("owner"))
+
+    assert result.status is BuildStatus.STOPPING
+    assert result.stopped_operation_ids == ("operation-1",)
+    assert state.status is BuildStatus.STOPPING
+
+
+@pytest.mark.asyncio
+async def test_stop_with_missing_binding_records_terminal_failure() -> None:
+    class RecordingState(_StopState):
+        error_summary: str | None = None
+
+        async def set_status(
+            self,
+            _build: int,
+            status: BuildStatus,
+            *,
+            error_summary: str | None = None,
+        ) -> None:
+            self.status = status
+            self.error_summary = error_summary
+
+    state = RecordingState(BuildStatus.PARTIAL)
+
+    async def missing_binding(_build: int):
+        raise LookupError("missing")
+
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(get_by_build=missing_binding),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    with pytest.raises(LookupError, match="missing"):
+        await service.stop(7, Principal("owner"))
+
+    assert state.status is BuildStatus.FAILED
+    assert state.error_summary == "MISSION_BINDING_MISSING"
+
+
 async def _async_value(value):
     return value

+ 119 - 48
script_build_host/tests/test_phase_one_e2e.py

@@ -16,6 +16,7 @@ from script_build_host.domain.artifacts import (
     ArtifactVersion,
     EvidenceRecordV1,
 )
+from script_build_host.domain.errors import ArtifactNotFound
 from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
 from script_build_host.domain.records import (
     BuildStatus,
@@ -68,6 +69,52 @@ def _last_tool_json(messages: list[dict[str, object]]) -> dict[str, object] | No
     return None
 
 
+def _contract(
+    kind: str,
+    objective: str,
+    criterion_id: str,
+    description: str,
+    *,
+    scope: str,
+) -> dict[str, object]:
+    output_schema = {
+        "direction": "script-direction/v1",
+        "decode-retrieval": "evidence-record/v1",
+    }[kind]
+    return {
+        "schema_version": "script-task-contract/v1",
+        "task_kind": kind,
+        "scope_ref": scope,
+        "intent_class": "explore",
+        "objective": objective,
+        "input_decision_refs": [],
+        "base_artifact_ref": None,
+        "write_scope": [scope],
+        "gap_ref": None,
+        "output_schema": output_schema,
+        "criteria": [
+            {
+                "criterion_id": criterion_id,
+                "description": description,
+                "hard": True,
+            }
+        ],
+        "budget": {
+            "max_attempts": 4,
+            "max_tokens": 32000,
+            "max_seconds": 900,
+            "max_external_queries": 40,
+            "max_no_improvement": 3,
+        },
+        "supersedes_decision_ids": [],
+        "candidate_closure_decision_refs": [],
+        "adopted_decision_ids": [],
+        "held_or_rejected_decision_ids": [],
+        "compose_order": [],
+        "comparison_decision_refs": [],
+    }
+
+
 class _ScriptedLLM:
     def __init__(self, store: FileSystemTaskStore, root: str) -> None:
         self.store = store
@@ -92,7 +139,7 @@ class _ScriptedLLM:
 
     async def __call__(self, messages, tools, **_kwargs):
         names = {item["function"]["name"] for item in tools or []}
-        if "task_plan" in names:
+        if "plan_script_tasks" in names:
             return await self._planner()
         if "submit_validation" in names:
             return self._validator(messages)
@@ -113,21 +160,16 @@ class _ScriptedLLM:
         )
         if direction is None:
             return self.call(
-                "task_plan",
+                "plan_script_tasks",
                 {
-                    "operation": "create",
-                    "tasks": [
-                        {
-                            "objective": "build direction",
-                            "acceptance_criteria": [
-                                {
-                                    "criterion_id": "direction-ready",
-                                    "description": "direction is grounded",
-                                    "hard": True,
-                                }
-                            ],
-                            "context_refs": ["script-build://task-kinds/direction"],
-                        }
+                    "contracts": [
+                        _contract(
+                            "direction",
+                            "build direction",
+                            "direction-ready",
+                            "direction is grounded",
+                            scope="script-build://scopes/direction",
+                        )
                     ],
                 },
             )
@@ -141,22 +183,17 @@ class _ScriptedLLM:
         )
         if retrieval is None:
             return self.call(
-                "task_plan",
+                "plan_script_tasks",
                 {
-                    "operation": "create",
                     "parent_task_id": direction.task_id,
-                    "tasks": [
-                        {
-                            "objective": "decode evidence",
-                            "acceptance_criteria": [
-                                {
-                                    "criterion_id": "evidence-ready",
-                                    "description": "evidence is relevant",
-                                    "hard": True,
-                                }
-                            ],
-                            "context_refs": ["script-build://task-kinds/decode-retrieval"],
-                        }
+                    "contracts": [
+                        _contract(
+                            "decode-retrieval",
+                            "decode evidence",
+                            "evidence-ready",
+                            "evidence is relevant",
+                            scope="script-build://scopes/direction/decode",
+                        )
                     ],
                 },
             )
@@ -177,17 +214,23 @@ class _ScriptedLLM:
             validation = ledger.validations[retrieval.validation_ids[-1]]
             if validation.verdict == ValidationVerdict.FAILED:
                 return self.call(
-                    "task_decide",
+                    "decide_script_task",
                     {
                         "task_id": retrieval.task_id,
                         "validation_id": validation.validation_id,
                         "action": "revise",
                         "reason": "make decode query precise",
-                        "payload": {"objective": "decode evidence revised"},
+                        "replacement_contract": _contract(
+                            "decode-retrieval",
+                            "decode evidence revised",
+                            "evidence-ready",
+                            "evidence is relevant",
+                            scope="script-build://scopes/direction/decode",
+                        ),
                     },
                 )
             return self.call(
-                "task_decide",
+                "decide_script_task",
                 {
                     "task_id": retrieval.task_id,
                     "validation_id": validation.validation_id,
@@ -200,7 +243,7 @@ class _ScriptedLLM:
         if direction.status == TaskStatus.AWAITING_DECISION:
             validation = ledger.validations[direction.validation_ids[-1]]
             return self.call(
-                "task_decide",
+                "decide_script_task",
                 {
                     "task_id": direction.task_id,
                     "validation_id": validation.validation_id,
@@ -211,18 +254,27 @@ class _ScriptedLLM:
         assert direction.status == TaskStatus.COMPLETED
         if root.current_spec_version == 1:
             return self.call(
-                "task_decide",
+                "decide_script_task",
                 {
                     "task_id": root.task_id,
                     "action": "revise",
                     "reason": "record accepted direction in root objective",
-                    "payload": {
-                        "objective": root.current_spec.objective + " using accepted direction"
+                    "replacement_contract": {
+                        "objective": root.current_spec.objective + " using accepted direction",
+                        "acceptance_criteria": [
+                            {
+                                "criterion_id": item.criterion_id,
+                                "description": item.description,
+                                "hard": item.hard,
+                            }
+                            for item in root.current_spec.acceptance_criteria
+                        ],
+                        "context_refs": list(root.current_spec.context_refs),
                     },
                 },
             )
         return self.call(
-            "task_decide",
+            "decide_script_task",
             {
                 "task_id": root.task_id,
                 "action": "block",
@@ -243,11 +295,7 @@ class _ScriptedLLM:
                 )
             return self.call(
                 "submit_attempt",
-                {
-                    "summary": "decode evidence",
-                    "artifact_refs": [],
-                    "evidence_refs": [last["artifact_ref"]],
-                },
+                {},
             )
         if last is None or "artifact_ref" not in last:
             accepted = prompt["accepted_child_results"]
@@ -274,11 +322,7 @@ class _ScriptedLLM:
             )
         return self.call(
             "submit_attempt",
-            {
-                "summary": "direction candidate",
-                "artifact_refs": [last["artifact_ref"]],
-                "evidence_refs": [],
-            },
+            {},
         )
 
     def _validator(self, messages):
@@ -313,7 +357,22 @@ class _ScriptedLLM:
                     }
                     for item in spec["acceptance_criteria"]
                 ],
-                "summary": verdict,
+                "defects": (
+                    [
+                        {
+                            "defect_code": "DECODE_QUERY_TOO_BROAD",
+                            "criterion_id": spec["acceptance_criteria"][0]["criterion_id"],
+                            "scope_ref": "script-build://scopes/direction/decode",
+                            "observed_excerpt": "v1 decode evidence was too broad",
+                            "evidence_refs": [],
+                            "severity": "hard",
+                            "invalidated_inputs": [],
+                            "recommended_action_class": "revise",
+                        }
+                    ]
+                    if failed
+                    else []
+                ),
                 "recommendation": "revise" if failed else "accept",
             },
         )
@@ -388,6 +447,18 @@ class _Artifacts:
     async def read_by_ref(self, ref, **_kwargs):
         return self.values[int(ref.version)][0]
 
+    async def get_by_attempt(self, *, script_build_id, task_id, attempt_id):
+        try:
+            return next(
+                value
+                for value, _ref in self.values.values()
+                if value.script_build_id == script_build_id
+                and value.task_id == task_id
+                and value.attempt_id == attempt_id
+            )
+        except StopIteration as exc:
+            raise ArtifactNotFound() from exc
+
 
 class _Legacy:
     def __init__(self) -> None:

+ 467 - 0
script_build_host/tests/test_phase_two_agent_tools.py

@@ -0,0 +1,467 @@
+from __future__ import annotations
+
+import base64
+from collections.abc import Mapping
+from types import SimpleNamespace
+from typing import Any
+
+import pytest
+from agent import ToolRegistry
+from agent.orchestration import ArtifactRef, ValidationVerdict
+
+from script_build_host.domain.artifacts import ArtifactState
+from script_build_host.domain.errors import ArtifactNotFound, ProtocolViolation
+from script_build_host.tools.contracts import AttemptManifest
+from script_build_host.tools.gateway import LegacyScriptToolGateway, RetrievalResult
+from script_build_host.tools.registry import register_script_tools
+
+
+class _CandidateTools:
+    def __init__(self, artifact_ref: ArtifactRef) -> None:
+        self.artifact_ref = artifact_ref
+        self.last_payload: dict[str, Any] | None = None
+        self.last_context: dict[str, Any] | None = None
+        self.images: list[dict[str, Any]] = []
+
+    async def resolve_attempt_manifest(self, *, context: Mapping[str, Any]) -> AttemptManifest:
+        self.last_context = dict(context)
+        return AttemptManifest(
+            artifact_ref=self.artifact_ref,
+            scope_ref="script-build://scopes/body",
+        )
+
+    async def create_script_paragraph(
+        self, *, payload: Mapping[str, Any], context: Mapping[str, Any]
+    ) -> dict[str, int]:
+        self.last_payload = dict(payload)
+        self.last_context = dict(context)
+        return {"paragraph_id": 7}
+
+    async def view_frozen_images(
+        self,
+        *,
+        raw_artifact_refs: list[str],
+        context: Mapping[str, Any],
+    ) -> list[dict[str, Any]]:
+        self.last_payload = {"raw_artifact_refs": raw_artifact_refs}
+        self.last_context = dict(context)
+        return self.images
+
+
+class _TaskStore:
+    def __init__(self, task: Any) -> None:
+        self.task = task
+
+    async def load(self, _root_trace_id: str) -> Any:
+        return SimpleNamespace(
+            tasks={"task-1": self.task},
+            attempts={"attempt-a": SimpleNamespace(attempt_id="attempt-a", task_id="task-1")},
+        )
+
+
+class _ArtifactStore:
+    def __init__(self, ref: ArtifactRef) -> None:
+        self.ref = ref
+
+    async def get(self, _root_trace_id: str, _snapshot_id: str) -> Any:
+        return SimpleNamespace(attempt_id="attempt-a", artifact_refs=[self.ref], evidence_refs=[])
+
+
+class _Bindings:
+    async def get_by_root(self, _root_trace_id: str) -> Any:
+        return SimpleNamespace(script_build_id=9, input_snapshot_id=3)
+
+
+class _BudgetGuard:
+    def __init__(self, error: Exception | None = None) -> None:
+        self.error = error
+        self.calls: list[tuple[str, str, str]] = []
+
+    async def ensure_tool_budget(self, *, root_trace_id: str, task_id: str, entry: str) -> None:
+        self.calls.append((root_trace_id, task_id, entry))
+        if self.error is not None:
+            raise self.error
+
+
+class _BusinessArtifacts:
+    def __init__(self, ref: ArtifactRef) -> None:
+        self.ref = ref
+
+    async def read_by_ref(
+        self,
+        ref: ArtifactRef,
+        *,
+        script_build_id: int,
+        task_id: str,
+        attempt_id: str,
+    ) -> Any:
+        assert (ref, script_build_id, task_id, attempt_id) == (
+            self.ref,
+            9,
+            "task-1",
+            "attempt-a",
+        )
+        return SimpleNamespace(
+            artifact_type=SimpleNamespace(value=ref.kind),
+            canonical_sha256=ref.digest,
+            state=SimpleNamespace(value="frozen"),
+        )
+
+
+class _RetrievalArtifacts:
+    def __init__(self) -> None:
+        self.version: Any | None = None
+
+    async def get_by_attempt(self, **_owners: Any) -> Any:
+        if self.version is None:
+            raise ArtifactNotFound()
+        return self.version
+
+    async def freeze(self, *, artifact: Any, **owners: Any) -> tuple[Any, ArtifactRef]:
+        ref = ArtifactRef(
+            "script-build://artifact-versions/8",
+            "evidence",
+            "8",
+            "sha256:" + "b" * 64,
+        )
+        self.version = SimpleNamespace(
+            artifact_version_id=8,
+            state=ArtifactState.FROZEN,
+            artifact=artifact,
+        )
+        return self.version, ref
+
+
+class _CountingRetrievalAdapter:
+    def __init__(self) -> None:
+        self.calls = 0
+
+    async def retrieve(self, **_kwargs: Any) -> RetrievalResult:
+        self.calls += 1
+        return RetrievalResult(
+            source_refs=("decode://case/1",),
+            summary="one immutable result",
+        )
+
+
+class _Coordinator:
+    def __init__(self, ref: ArtifactRef) -> None:
+        criterion = SimpleNamespace(criterion_id="criterion-a", hard=True)
+        task = SimpleNamespace(
+            task_id="task-1",
+            validation_ids=(),
+            current_spec=SimpleNamespace(
+                acceptance_criteria=(criterion,),
+                context_refs=("script-build://task-kinds/paragraph",),
+            ),
+        )
+        self.task_store = _TaskStore(task)
+        self.artifact_store = _ArtifactStore(ref)
+        self.attempt_submission: tuple[dict[str, Any], Any] | None = None
+        self.validation_submission: tuple[dict[str, Any], tuple[Any, ...]] | None = None
+
+    async def submit_attempt(self, context: Mapping[str, Any], submission: Any) -> dict[str, Any]:
+        self.attempt_submission = (dict(context), submission)
+        return {
+            "attempt_id": context["attempt_id"],
+            "snapshot_id": "snapshot-a",
+            "status": "awaiting_validation",
+        }
+
+    async def submit_validation(self, context: Mapping[str, Any], *values: Any) -> dict[str, Any]:
+        self.validation_submission = (dict(context), values)
+        return {
+            "validation_id": context["validation_id"],
+            "verdict": values[0].value,
+            "status": "awaiting_decision",
+        }
+
+
+def _gateway(ref: ArtifactRef) -> tuple[LegacyScriptToolGateway, _Coordinator, _CandidateTools]:
+    coordinator = _Coordinator(ref)
+    candidates = _CandidateTools(ref)
+    gateway = LegacyScriptToolGateway(
+        bindings=_Bindings(),  # type: ignore[arg-type]
+        snapshots=SimpleNamespace(),
+        artifacts=_BusinessArtifacts(ref),  # type: ignore[arg-type]
+        coordinator=coordinator,
+        candidate_tools=candidates,  # type: ignore[arg-type]
+    )
+    return gateway, coordinator, candidates
+
+
+@pytest.mark.asyncio
+async def test_attempt_submission_is_derived_from_protected_attempt_artifact() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="paragraph",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, coordinator, candidates = _gateway(ref)
+    budget = _BudgetGuard()
+    gateway.budget_guard = budget
+    context = {
+        "root_trace_id": "root-a",
+        "task_id": "task-1",
+        "attempt_id": "attempt-a",
+        "spec_version": 1,
+    }
+    await gateway.submit_current_attempt(context)
+
+    assert coordinator.attempt_submission is not None
+    submitted_context, submission = coordinator.attempt_submission
+    assert submitted_context == context
+    assert submission.artifact_refs == [ref]
+    assert submission.evidence_refs == []
+    assert "paragraph" in submission.summary
+    assert "sha256:" in submission.summary
+    assert candidates.last_context == context
+    assert budget.calls == [("root-a", "task-1", "submit")]
+
+
+@pytest.mark.asyncio
+async def test_retrieval_budget_is_checked_before_adapter_or_snapshot_access() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="evidence",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, _, _ = _gateway(ref)
+    budget = _BudgetGuard(ProtocolViolation("retrieval budget exhausted"))
+    gateway.budget_guard = budget
+    with pytest.raises(ProtocolViolation, match="budget exhausted"):
+        await gateway.retrieve(
+            "decode",
+            "search_script_decode_case",
+            {"query": "opening"},
+            {
+                "root_trace_id": "root-a",
+                "task_id": "task-1",
+                "attempt_id": "attempt-a",
+            },
+        )
+    assert budget.calls == [("root-a", "task-1", "retrieval")]
+
+
+@pytest.mark.asyncio
+async def test_one_retrieval_attempt_never_calls_adapter_twice() -> None:
+    artifacts = _RetrievalArtifacts()
+    adapter = _CountingRetrievalAdapter()
+
+    class Snapshots:
+        async def get(self, snapshot_id: str, *, script_build_id: int) -> Any:
+            assert (snapshot_id, script_build_id) == ("3", 9)
+            return SimpleNamespace(account={})
+
+    gateway = LegacyScriptToolGateway(
+        bindings=_Bindings(),  # type: ignore[arg-type]
+        snapshots=Snapshots(),  # type: ignore[arg-type]
+        artifacts=artifacts,  # type: ignore[arg-type]
+        coordinator=SimpleNamespace(),
+        retrieval_adapters={"decode": adapter},
+    )
+    context = {
+        "root_trace_id": "root-a",
+        "task_id": "task-1",
+        "attempt_id": "attempt-a",
+        "spec_version": 1,
+    }
+    await gateway.retrieve("decode", "search_script_decode_case", {"query": "opening"}, context)
+    with pytest.raises(ProtocolViolation, match="only one Evidence"):
+        await gateway.retrieve(
+            "decode", "search_script_decode_case", {"query": "different"}, context
+        )
+    assert adapter.calls == 1
+
+
+@pytest.mark.asyncio
+async def test_candidate_write_receives_identity_only_from_protected_context() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="paragraph",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, _, candidates = _gateway(ref)
+    context = {
+        "root_trace_id": "root-a",
+        "task_id": "task-1",
+        "attempt_id": "attempt-a",
+        "spec_version": 1,
+    }
+    await gateway.candidate_command(
+        "create_script_paragraph",
+        {"paragraph_index": 1, "name": "opening", "content_range": {}},
+        context,
+    )
+    assert candidates.last_context == context
+    assert candidates.last_payload == {
+        "paragraph_index": 1,
+        "name": "opening",
+        "content_range": {},
+    }
+
+
+@pytest.mark.asyncio
+async def test_structured_validation_keeps_excerpt_out_of_framework_summary() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="paragraph",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, coordinator, _ = _gateway(ref)
+    context = {
+        "root_trace_id": "root-a",
+        "task_id": "task-1",
+        "attempt_id": "attempt-a",
+        "snapshot_id": "snapshot-a",
+        "validation_id": "validation-a",
+    }
+    excerpt = "this exact candidate passage must not enter the bounded summary"
+    await gateway.submit_structured_validation(
+        verdict="failed",
+        criterion_results=[
+            {"criterion_id": "criterion-a", "verdict": "failed", "reason": "not realized"}
+        ],
+        defects=[
+            {
+                "defect_code": "REALIZATION_PLACEHOLDER",
+                "criterion_id": "criterion-a",
+                "scope_ref": "script-build://scopes/body/end",
+                "observed_excerpt": excerpt,
+                "evidence_refs": [
+                    {
+                        "uri": ref.uri,
+                        "kind": ref.kind,
+                        "version": ref.version,
+                        "digest": ref.digest,
+                    }
+                ],
+                "severity": "hard",
+                "invalidated_inputs": ["decision-a"],
+                "recommended_action_class": "split",
+            }
+        ],
+        recommendation="split",
+        context=context,
+    )
+
+    assert coordinator.validation_submission is not None
+    _, values = coordinator.validation_submission
+    verdict, results, summary, evidence_refs, unverified, risks, recommendation = values
+    assert verdict is ValidationVerdict.FAILED
+    assert results[0].reason.endswith("defect_codes=REALIZATION_PLACEHOLDER")
+    assert excerpt not in summary
+    assert len(summary) <= 500
+    assert evidence_refs == (ref,)
+    assert unverified == ()
+    assert risks == ("REALIZATION_PLACEHOLDER",)
+    assert recommendation == "split"
+
+
+@pytest.mark.asyncio
+async def test_passed_validation_rejects_blocking_defect_and_unfrozen_evidence() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="paragraph",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, _, _ = _gateway(ref)
+    context = {
+        "root_trace_id": "root-a",
+        "task_id": "task-1",
+        "attempt_id": "attempt-a",
+        "snapshot_id": "snapshot-a",
+        "validation_id": "validation-a",
+    }
+    defect: dict[str, Any] = {
+        "defect_code": "WRITE_SCOPE_CONFLICT",
+        "criterion_id": "criterion-a",
+        "scope_ref": "script-build://scopes/body",
+        "observed_excerpt": "overlapping body",
+        "evidence_refs": [],
+        "severity": "hard",
+        "invalidated_inputs": [],
+        "recommended_action_class": "replace",
+    }
+    with pytest.raises(ProtocolViolation, match="cannot contain hard"):
+        await gateway.submit_structured_validation(
+            verdict="passed",
+            criterion_results=[
+                {"criterion_id": "criterion-a", "verdict": "passed", "reason": "ok"}
+            ],
+            defects=[defect],
+            recommendation="replace",
+            context=context,
+        )
+
+    defect["severity"] = "warning"
+    defect["evidence_refs"] = [
+        {
+            "uri": "script-build://artifact-versions/8",
+            "kind": "paragraph",
+            "version": "8",
+            "digest": "sha256:" + "b" * 64,
+        }
+    ]
+    with pytest.raises(ProtocolViolation, match="fixed validation snapshot"):
+        await gateway.submit_structured_validation(
+            verdict="passed",
+            criterion_results=[
+                {"criterion_id": "criterion-a", "verdict": "passed", "reason": "ok"}
+            ],
+            defects=[defect],
+            recommendation="accept",
+            context=context,
+        )
+
+
+@pytest.mark.asyncio
+async def test_view_frozen_images_returns_multimodal_without_url_or_data_in_text() -> None:
+    ref = ArtifactRef(
+        uri="script-build://artifact-versions/7",
+        kind="paragraph",
+        version="7",
+        digest="sha256:" + "a" * 64,
+    )
+    gateway, _, candidates = _gateway(ref)
+    encoded = base64.b64encode(b"safe-image-bytes").decode("ascii")
+    candidates.images = [
+        {
+            "type": "base64",
+            "media_type": "image/png",
+            "data": encoded,
+            "raw_artifact_ref": "script-build://raw-artifacts/sha256/abc",
+        }
+    ]
+    registry = ToolRegistry()
+    register_script_tools(registry, gateway)
+    result = await registry.execute(
+        "view_frozen_images",
+        {"raw_artifact_refs": ["script-build://raw-artifacts/sha256/abc"]},
+        context={"root_trace_id": "root-a", "task_id": "task-1"},
+    )
+    assert isinstance(result, dict)
+    assert result["images"] == [{"type": "base64", "media_type": "image/png", "data": encoded}]
+    assert encoded not in result["text"]
+    assert "http" not in result["text"]
+
+    candidates.images = [
+        {
+            "type": "base64",
+            "media_type": "image/png",
+            "data": encoded,
+            "url": "https://forbidden.invalid/image.png",
+        }
+    ]
+    rejected = await registry.execute(
+        "view_frozen_images",
+        {"raw_artifact_refs": ["script-build://raw-artifacts/sha256/abc"]},
+        context={"root_trace_id": "root-a", "task_id": "task-1"},
+    )
+    assert isinstance(rejected, str)
+    assert "must not contain network URLs" in rejected

+ 438 - 0
script_build_host/tests/test_phase_two_boundary.py

@@ -0,0 +1,438 @@
+from __future__ import annotations
+
+import hashlib
+import json
+from dataclasses import replace
+from datetime import UTC, datetime
+
+import pytest
+from agent.orchestration import (
+    AcceptanceCriterion,
+    ArtifactRef,
+    ArtifactSnapshot,
+    AttemptStatus,
+    AttemptSubmission,
+    DecisionAction,
+    PlannerDecision,
+    TaskAttempt,
+    TaskLedger,
+    TaskRecord,
+    TaskSpec,
+    TaskStatus,
+    ValidationReport,
+    ValidationRunStatus,
+    ValidationVerdict,
+)
+
+from script_build_host.application.phase_two_boundary import ScriptPhaseTwoBoundaryVerifier
+from script_build_host.application.phase_two_inputs import AcceptedInputResolver
+from script_build_host.domain.artifacts import ArtifactKind, ArtifactState, ArtifactVersion
+from script_build_host.domain.errors import PhaseTwoBoundaryNotReady
+from script_build_host.domain.phase_two_artifacts import (
+    CandidatePortfolioArtifactV1,
+    ScriptParagraphV1,
+    StructuredScriptArtifactV1,
+)
+from script_build_host.domain.records import MissionBinding
+from script_build_host.domain.task_contracts import (
+    AcceptedDecisionRef,
+    ScriptCriterion,
+    ScriptIntentClass,
+    ScriptTaskBudget,
+    ScriptTaskContractV1,
+    ScriptTaskKind,
+)
+
+SNAPSHOT_ID = "11"
+INPUT_DIGEST = "sha256:" + "c" * 64
+SCRIPT_DIGEST = "sha256:" + "a" * 64
+PORTFOLIO_DIGEST = "sha256:" + "b" * 64
+SCOPE = "script-build://scopes/full"
+WRITE = "script-build://writes/full"
+
+
+def _task(task_id: str, parent: str | None, kind: str, status: TaskStatus) -> TaskRecord:
+    return TaskRecord(
+        task_id=task_id,
+        goal_id=None,
+        parent_task_id=parent,
+        display_path=task_id,
+        specs=[
+            TaskSpec(
+                version=1,
+                objective="bounded",
+                acceptance_criteria=(AcceptanceCriterion("closed", "closed"),),
+                context_refs=(
+                    f"script-build://task-kinds/{kind}",
+                    f"script-build://inputs/{SNAPSHOT_ID}",
+                ),
+            )
+        ],
+        status=status,
+    )
+
+
+def _compose_contract() -> ScriptTaskContractV1:
+    return ScriptTaskContractV1(
+        ScriptTaskKind.COMPOSE,
+        SCOPE,
+        ScriptIntentClass.COMPOSE,
+        "compose one complete candidate",
+        (),
+        None,
+        (WRITE,),
+        None,
+        "structured-script/v1",
+        (ScriptCriterion("closed", "closed"),),
+        ScriptTaskBudget(),
+    )
+
+
+def _portfolio_contract(ref: ArtifactRef) -> ScriptTaskContractV1:
+    accepted = AcceptedDecisionRef("decision-compose", ref, SCOPE, ScriptTaskKind.COMPOSE)
+    return ScriptTaskContractV1(
+        ScriptTaskKind.CANDIDATE_PORTFOLIO,
+        SCOPE,
+        ScriptIntentClass.PORTFOLIO,
+        "govern the closed candidate set",
+        (accepted,),
+        None,
+        (WRITE,),
+        None,
+        "candidate-portfolio/v1",
+        (ScriptCriterion("closed", "closed"),),
+        ScriptTaskBudget(),
+        candidate_closure_decision_refs=(accepted,),
+        adopted_decision_ids=("decision-compose",),
+        compose_order=("decision-compose",),
+    )
+
+
+class _Store:
+    def __init__(self, value: object) -> None:
+        self.value = value
+
+    async def load(self, root: str) -> TaskLedger:
+        assert root == "root"
+        assert isinstance(self.value, TaskLedger)
+        return self.value
+
+    async def get(self, root: str, snapshot: str) -> ArtifactSnapshot:
+        assert root == "root"
+        return self.value[snapshot]  # type: ignore[index]
+
+
+class _Contracts:
+    def __init__(self, values: dict[str, ScriptTaskContractV1]) -> None:
+        self.values = values
+
+    async def read_for_task(
+        self, *, root_trace_id: str, task: TaskRecord, spec_version: int | None = None
+    ) -> ScriptTaskContractV1:
+        assert root_trace_id == "root" and spec_version == 1
+        return self.values[task.task_id]
+
+
+class _Artifacts:
+    def __init__(self, versions: dict[int, ArtifactVersion]) -> None:
+        self.versions = versions
+
+    async def read_by_ref(
+        self,
+        ref: ArtifactRef,
+        *,
+        script_build_id: int,
+        task_id: str | None = None,
+        attempt_id: str | None = None,
+    ) -> ArtifactVersion:
+        assert script_build_id == 7 and ref.version is not None
+        value = self.versions[int(ref.version)]
+        assert value.task_id == task_id and value.attempt_id == attempt_id
+        assert value.canonical_sha256 == ref.digest
+        return value
+
+    async def get_by_id(self, identifier: int, *, script_build_id: int) -> ArtifactVersion:
+        assert script_build_id == 7
+        return self.versions[identifier]
+
+
+class _Bindings:
+    async def get_by_build(self, script_build_id: int) -> MissionBinding:
+        assert script_build_id == 7
+        now = datetime.now(UTC)
+        return MissionBinding(1, 7, "root", 11, None, None, "test", "v1", now, now)
+
+
+def _version(
+    identifier: int,
+    task_id: str,
+    attempt_id: str,
+    kind: ArtifactKind,
+    digest: str,
+    artifact: object,
+) -> ArtifactVersion:
+    now = datetime.now(UTC)
+    return ArtifactVersion(
+        identifier,
+        7,
+        task_id,
+        attempt_id,
+        1,
+        kind,
+        digest,
+        ArtifactState.FROZEN,
+        artifact,  # type: ignore[arg-type]
+        now,
+        now,
+    )
+
+
+async def _fixture() -> tuple[ScriptPhaseTwoBoundaryVerifier, TaskLedger]:
+    script_ref = ArtifactRef(
+        "script-build://artifact-versions/1",
+        ArtifactKind.STRUCTURED_SCRIPT.value,
+        "1",
+        SCRIPT_DIGEST,
+    )
+    root = _task("root-task", None, "root", TaskStatus.BLOCKED)
+    portfolio_task = _task("portfolio", "root-task", "candidate-portfolio", TaskStatus.RUNNING)
+    compose_task = _task("compose", "portfolio", "compose", TaskStatus.COMPLETED)
+    root.child_task_ids.append("portfolio")
+    portfolio_task.child_task_ids.append("compose")
+
+    compose_attempt = TaskAttempt(
+        "attempt-compose",
+        "compose",
+        1,
+        "worker-compose",
+        "script_compose_worker",
+        "new",
+        (),
+        AttemptStatus.SUBMITTED,
+        snapshot_id="snapshot-compose",
+        submission=AttemptSubmission("closed", artifact_refs=[script_ref]),
+    )
+    compose_validation = ValidationReport(
+        "validation-compose",
+        "compose",
+        "attempt-compose",
+        1,
+        "snapshot-compose",
+        "validator-compose",
+        status=ValidationRunStatus.COMPLETED,
+        verdict=ValidationVerdict.PASSED,
+    )
+    compose_decision = PlannerDecision(
+        "decision-compose",
+        "compose",
+        DecisionAction.ACCEPT,
+        "passed",
+        TaskStatus.AWAITING_DECISION,
+        TaskStatus.COMPLETED,
+        "attempt-compose",
+        "validation-compose",
+    )
+    compose_task.attempt_ids.append("attempt-compose")
+    compose_task.validation_ids.append("validation-compose")
+    compose_task.decision_ids.append("decision-compose")
+
+    portfolio_attempt = TaskAttempt(
+        "attempt-portfolio",
+        "portfolio",
+        1,
+        "worker-portfolio",
+        "script_candidate_portfolio_worker",
+        "new",
+        ("decision-compose",),
+    )
+    ledger = TaskLedger(
+        "root",
+        "mission",
+        "root-task",
+        tasks={item.task_id: item for item in (root, portfolio_task, compose_task)},
+        attempts={
+            compose_attempt.attempt_id: compose_attempt,
+            portfolio_attempt.attempt_id: portfolio_attempt,
+        },
+        validations={compose_validation.validation_id: compose_validation},
+        decisions={compose_decision.decision_id: compose_decision},
+    )
+    compose_normalized = _normalized("closed", script_ref)
+    snapshots = {
+        "snapshot-compose": ArtifactSnapshot(
+            "snapshot-compose",
+            "attempt-compose",
+            compose_normalized,
+            _snapshot_sha(compose_normalized),
+            [script_ref],
+            [],
+        )
+    }
+    paragraph = ScriptParagraphV1(1, 1, 1, None, "opening", {}, description="done")
+    structured = StructuredScriptArtifactV1(
+        "script-build://artifact-versions/9",
+        INPUT_DIGEST,
+        (paragraph,),
+        (),
+        (),
+        ("script-build://artifact-versions/9",),
+        (),
+        (),
+        SCRIPT_DIGEST,
+    )
+    artifacts = _Artifacts(
+        {
+            1: _version(
+                1,
+                "compose",
+                "attempt-compose",
+                ArtifactKind.STRUCTURED_SCRIPT,
+                SCRIPT_DIGEST,
+                structured,
+            )
+        }
+    )
+    contracts = _Contracts(
+        {"compose": _compose_contract(), "portfolio": _portfolio_contract(script_ref)}
+    )
+    accepted = AcceptedInputResolver(
+        task_store=_Store(ledger),
+        artifact_store=_Store(snapshots),
+        artifacts=artifacts,
+        contracts=contracts,
+    )
+    bundle = await accepted.resolve(
+        root_trace_id="root",
+        script_build_id=7,
+        task_id="portfolio",
+        attempt_id="attempt-portfolio",
+        contract=contracts.values["portfolio"],
+        input_snapshot_id=SNAPSHOT_ID,
+    )
+    portfolio = CandidatePortfolioArtifactV1(
+        script_ref.uri,
+        (script_ref.uri,),
+        ("decision-compose",),
+        (),
+        (),
+        bundle.input_closure_digest,
+        (),
+        ("decision-compose",),
+        PORTFOLIO_DIGEST,
+    )
+    portfolio_ref = ArtifactRef(
+        "script-build://artifact-versions/2",
+        ArtifactKind.CANDIDATE_PORTFOLIO.value,
+        "2",
+        PORTFOLIO_DIGEST,
+    )
+    portfolio_attempt.status = AttemptStatus.SUBMITTED
+    portfolio_attempt.snapshot_id = "snapshot-portfolio"
+    portfolio_attempt.submission = AttemptSubmission("closed", artifact_refs=[portfolio_ref])
+    portfolio_validation = ValidationReport(
+        "validation-portfolio",
+        "portfolio",
+        "attempt-portfolio",
+        1,
+        "snapshot-portfolio",
+        "validator-portfolio",
+        status=ValidationRunStatus.COMPLETED,
+        verdict=ValidationVerdict.PASSED,
+    )
+    portfolio_decision = PlannerDecision(
+        "decision-portfolio",
+        "portfolio",
+        DecisionAction.ACCEPT,
+        "passed",
+        TaskStatus.AWAITING_DECISION,
+        TaskStatus.COMPLETED,
+        "attempt-portfolio",
+        "validation-portfolio",
+    )
+    portfolio_task.status = TaskStatus.COMPLETED
+    portfolio_task.attempt_ids.append("attempt-portfolio")
+    portfolio_task.validation_ids.append("validation-portfolio")
+    portfolio_task.decision_ids.append("decision-portfolio")
+    ledger.validations[portfolio_validation.validation_id] = portfolio_validation
+    ledger.decisions[portfolio_decision.decision_id] = portfolio_decision
+    portfolio_normalized = _normalized("closed", portfolio_ref)
+    snapshots["snapshot-portfolio"] = ArtifactSnapshot(
+        "snapshot-portfolio",
+        "attempt-portfolio",
+        portfolio_normalized,
+        _snapshot_sha(portfolio_normalized),
+        [portfolio_ref],
+        [],
+    )
+    artifacts.versions[2] = _version(
+        2,
+        "portfolio",
+        "attempt-portfolio",
+        ArtifactKind.CANDIDATE_PORTFOLIO,
+        PORTFOLIO_DIGEST,
+        portfolio,
+    )
+    verifier = ScriptPhaseTwoBoundaryVerifier(
+        task_store=_Store(ledger),
+        artifact_store=_Store(snapshots),
+        artifacts=artifacts,
+        bindings=_Bindings(),
+        accepted_inputs=accepted,
+    )
+    return verifier, ledger
+
+
+@pytest.mark.asyncio
+async def test_phase_two_boundary_requires_one_closed_adopted_portfolio() -> None:
+    verifier, _ = await _fixture()
+    await verifier.verify_checkpoint(
+        script_build_id=7, root_trace_id="root", input_snapshot_id=SNAPSHOT_ID
+    )
+
+
+@pytest.mark.asyncio
+async def test_phase_two_boundary_rejects_nonterminal_subtree_and_early_advance() -> None:
+    verifier, ledger = await _fixture()
+    ledger.tasks["compose"].status = TaskStatus.RUNNING
+    with pytest.raises(PhaseTwoBoundaryNotReady, match="subtree is not terminal"):
+        await verifier.verify_checkpoint(
+            script_build_id=7, root_trace_id="root", input_snapshot_id=SNAPSHOT_ID
+        )
+    with pytest.raises(PhaseTwoBoundaryNotReady, match="portfolio already exists"):
+        await verifier.verify_before_advance(
+            script_build_id=7, root_trace_id="root", input_snapshot_id=SNAPSHOT_ID
+        )
+
+
+@pytest.mark.asyncio
+async def test_phase_two_boundary_rejects_closure_digest_drift() -> None:
+    verifier, _ = await _fixture()
+    version = verifier._artifacts.versions[2]
+    artifact = replace(version.artifact, input_closure_digest="sha256:" + "f" * 64)
+    verifier._artifacts.versions[2] = replace(version, artifact=artifact)
+    with pytest.raises(PhaseTwoBoundaryNotReady, match="closure digest"):
+        await verifier.verify_checkpoint(
+            script_build_id=7, root_trace_id="root", input_snapshot_id=SNAPSHOT_ID
+        )
+
+
+def _normalized(summary: str, ref: ArtifactRef) -> dict[str, object]:
+    return {
+        "summary": summary,
+        "artifact_refs": [
+            {
+                "uri": ref.uri,
+                "kind": ref.kind,
+                "version": ref.version,
+                "digest": ref.digest,
+                "summary": ref.summary,
+                "metadata": ref.metadata,
+            }
+        ],
+        "evidence_refs": [],
+    }
+
+
+def _snapshot_sha(value: dict[str, object]) -> str:
+    canonical = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
+    return hashlib.sha256(canonical.encode()).hexdigest()

+ 932 - 0
script_build_host/tests/test_phase_two_lifecycle.py

@@ -0,0 +1,932 @@
+from __future__ import annotations
+
+import asyncio
+from dataclasses import replace
+from datetime import UTC, datetime
+from types import SimpleNamespace
+from unittest.mock import AsyncMock
+
+import pytest
+from agent.orchestration import DecisionAction, OperationStatus, TaskStatus
+from agent.trace.models import Message, Trace
+from agent.trace.store import FileSystemTraceStore
+
+from script_build_host.application.mission_factory import ScriptMissionFactory
+from script_build_host.application.mission_service import (
+    PHASE_ONE_CAPABILITY_BOUNDARY,
+    PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
+    ScriptMissionService,
+)
+from script_build_host.application.observation_views import mission_snapshot_view
+from script_build_host.domain.errors import (
+    MissionRecoveryRequired,
+    PhaseTwoBoundaryNotReady,
+    PlannerPolicyMigrationRequired,
+)
+from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
+from script_build_host.domain.records import (
+    BuildStatus,
+    MissionBinding,
+    Principal,
+    PublicationState,
+)
+from script_build_host.infrastructure.canonical_json import canonical_sha256
+
+
+class _TraceStore:
+    def __init__(self) -> None:
+        self.events: list[dict[str, object]] = []
+        self.messages: list[SimpleNamespace] = []
+
+    async def get_events(self, _root: str, _since: int):
+        return list(self.events)
+
+    async def append_event(self, _root: str, event: str, payload: dict[str, object]):
+        self.events.append({"event": event, **payload})
+        return len(self.events)
+
+    async def get_trace_messages(self, _root: str):
+        return list(self.messages)
+
+
+class _Runner:
+    def __init__(self, root: SimpleNamespace) -> None:
+        self.root = root
+        self.trace_store = _TraceStore()
+        self.tools = SimpleNamespace(
+            get_tool_names=lambda **_kwargs: ["plan_script_tasks", "decide_script_task"]
+        )
+        self.observed_config = None
+
+    async def run(self, *, messages, config):
+        self.observed_config = config
+        assert [item["role"] for item in messages] == ["system", "user"]
+        yield SimpleNamespace(trace_id="root")
+        yield SimpleNamespace(role="system", sequence=11)
+        yield SimpleNamespace(role="user", sequence=12)
+        assert self.root.status is TaskStatus.NEEDS_REPLAN
+        self.root.status = TaskStatus.BLOCKED
+        self.root.blocked_reason = PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+
+
+class _Coordinator:
+    def __init__(self, root: SimpleNamespace) -> None:
+        self.root = root
+        self.unblock_calls = 0
+        self.task_store = SimpleNamespace(
+            load=lambda _root: _async(
+                SimpleNamespace(root_task_id="root-task", tasks={"root-task": root})
+            )
+        )
+
+    async def decide_task(
+        self,
+        _root_trace_id,
+        task_id,
+        _validation_id,
+        action,
+        _payload,
+        idempotency_key,
+    ):
+        assert task_id == "root-task"
+        assert action is DecisionAction.UNBLOCK
+        assert idempotency_key == "phase-two-unblock:7"
+        self.unblock_calls += 1
+        self.root.status = TaskStatus.NEEDS_REPLAN
+        self.root.blocked_reason = None
+
+    async def root_completion(self, _root: str):
+        return {
+            "status": self.root.status.value,
+            "blocked_reason": self.root.blocked_reason,
+        }
+
+
+class _State:
+    def __init__(self) -> None:
+        self.status = BuildStatus.RUNNING
+        self.checkpoint = None
+
+    async def get_status(self, _build: int):
+        return self.status
+
+    async def set_status(self, _build: int, status: BuildStatus, **_kwargs):
+        self.status = status
+
+    async def set_checkpoint(self, _build: int, **values):
+        self.status = BuildStatus.PARTIAL
+        self.checkpoint = values
+
+
+def _snapshot() -> ScriptBuildInputSnapshotV1:
+    return ScriptBuildInputSnapshotV1(
+        snapshot_id="11",
+        script_build_id=7,
+        execution_id=1,
+        topic_build_id=2,
+        topic_id=3,
+        topic={"topic": {"id": 3, "result": "topic"}},
+        account={"account_name": "acct"},
+        persona_points=(),
+        section_patterns=(),
+        strategies=(),
+        prompt_manifest=(),
+        datasource_manifest={},
+        model_manifest={
+            "presets": {
+                "script_planner": {
+                    "model": "fake",
+                    "temperature": 0,
+                    "max_iterations": 10,
+                }
+            }
+        },
+        canonical_sha256="sha256:" + "1" * 64,
+        created_at=datetime.now(UTC),
+    )
+
+
+def _binding() -> MissionBinding:
+    now = datetime.now(UTC)
+    return MissionBinding(
+        binding_id=1,
+        script_build_id=7,
+        root_trace_id="root",
+        input_snapshot_id=11,
+        active_direction_artifact_version_id=3,
+        accepted_root_artifact_version_id=None,
+        engine_version="test",
+        schema_version="v1",
+        created_at=now,
+        updated_at=now,
+    )
+
+
+@pytest.mark.asyncio
+async def test_phase_two_policy_is_persisted_before_root_unblock_and_same_trace_continues() -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.BLOCKED,
+        blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
+    )
+    runner = _Runner(root)
+    coordinator = _Coordinator(root)
+    state = _State()
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=coordinator,
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+    await service.run_phase_two(
+        7,
+        binding=_binding(),
+        snapshot=_snapshot(),
+        direction_artifact_version_id=3,
+    )
+    assert runner.observed_config.trace_id == "root"
+    assert runner.observed_config.new_trace_id is None
+    assert runner.trace_store.events[0]["event"] == "planner_policy_migrated"
+    assert state.status is BuildStatus.PARTIAL
+    assert state.checkpoint["checkpoint_code"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+
+
+async def _async(value):
+    return value
+
+
+class _ReentryRunner(_Runner):
+    def __init__(self, root: SimpleNamespace) -> None:
+        super().__init__(root)
+        self.run_result_calls = 0
+
+    async def run_result(self, *, messages, config):
+        self.run_result_calls += 1
+        self.observed_config = config
+        assert messages == []
+        assert self.root.status is TaskStatus.NEEDS_REPLAN
+        self.root.status = TaskStatus.BLOCKED
+        self.root.blocked_reason = PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+        return SimpleNamespace()
+
+
+def _migration_event(runner: _Runner, snapshot: ScriptBuildInputSnapshotV1) -> dict[str, object]:
+    factory = ScriptMissionFactory()
+    policy = factory.build_phase_two_policy(snapshot)
+    continuation = factory.build_phase_two_message(
+        _binding(), snapshot, direction_artifact_version_id=3
+    )
+    runner.trace_store.messages = [
+        SimpleNamespace(sequence=11, role="system", content=policy, parent_sequence=10),
+        SimpleNamespace(sequence=12, role="user", content=continuation, parent_sequence=11),
+    ]
+    return {
+        "event": "planner_policy_migrated",
+        "policy_version": "script-build-phase-two/v1",
+        "policy_digest": canonical_sha256(policy).wire,
+        "toolset_digest": canonical_sha256(
+            sorted(runner.tools.get_tool_names(groups=["script_build"]))
+        ).wire,
+        "input_snapshot_id": snapshot.snapshot_id,
+        "system_sequence": 11,
+        "user_sequence": 12,
+    }
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("root_status", [TaskStatus.BLOCKED, TaskStatus.NEEDS_REPLAN])
+async def test_phase_two_policy_reentry_resumes_before_or_after_unblock_without_rewriting_policy(
+    root_status: TaskStatus,
+) -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=root_status,
+        blocked_reason=(
+            PHASE_ONE_CAPABILITY_BOUNDARY if root_status is TaskStatus.BLOCKED else None
+        ),
+    )
+    runner = _ReentryRunner(root)
+    runner.trace_store.events = [_migration_event(runner, _snapshot())]
+    coordinator = _Coordinator(root)
+    state = _State()
+    state.status = BuildStatus.PARTIAL
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=coordinator,
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    await service.run_phase_two(
+        7,
+        binding=_binding(),
+        snapshot=_snapshot(),
+        direction_artifact_version_id=3,
+    )
+
+    assert runner.run_result_calls == 1
+    assert len(runner.trace_store.events) == 1
+    assert coordinator.unblock_calls == (1 if root_status is TaskStatus.BLOCKED else 0)
+    assert state.status is BuildStatus.PARTIAL
+    assert state.checkpoint["checkpoint_code"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+
+
+@pytest.mark.asyncio
+async def test_phase_two_reentry_with_output_requires_recovery_without_terminal_failure() -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.NEEDS_REPLAN,
+        blocked_reason=None,
+    )
+    runner = _ReentryRunner(root)
+    runner.trace_store.events = [_migration_event(runner, _snapshot())]
+    runner.trace_store.messages.append(
+        SimpleNamespace(sequence=13, role="assistant", content="started")
+    )
+    state = _State()
+    state.status = BuildStatus.RUNNING
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=_Coordinator(root),
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    with pytest.raises(MissionRecoveryRequired):
+        await service.run_phase_two(
+            7,
+            binding=_binding(),
+            snapshot=_snapshot(),
+            direction_artifact_version_id=3,
+        )
+
+    assert state.status is BuildStatus.RUNNING
+    assert runner.run_result_calls == 0
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("root_status", "has_model_output"),
+    [
+        (TaskStatus.BLOCKED, False),
+        (TaskStatus.NEEDS_REPLAN, False),
+        (TaskStatus.NEEDS_REPLAN, True),
+    ],
+)
+async def test_historical_snapshot_v2_policy_reentry_survives_filestore_reload(
+    tmp_path,
+    root_status: TaskStatus,
+    has_model_output: bool,
+) -> None:
+    parent = _snapshot()
+    business_digest = canonical_sha256(parent.to_input().business_payload()).wire
+    current = replace(
+        parent,
+        snapshot_id="12",
+        parent_snapshot_id=parent.snapshot_id,
+        parent_snapshot_sha256=parent.canonical_sha256,
+        business_input_sha256=business_digest,
+    )
+    binding = replace(_binding(), input_snapshot_id=12)
+    policy = ScriptMissionFactory().build_phase_two_policy(current)
+    continuation = ScriptMissionFactory().build_phase_two_message(
+        binding,
+        current,
+        direction_artifact_version_id=3,
+    )
+    store_path = tmp_path / "phase-two-reentry"
+    writer = FileSystemTraceStore(str(store_path))
+    await writer.create_trace(
+        Trace(
+            trace_id="root",
+            mode="agent",
+            agent_type="script_planner",
+            agent_role="planner",
+        )
+    )
+    await writer.add_message(
+        Message.create(
+            trace_id="root",
+            role="system",
+            sequence=11,
+            parent_sequence=10,
+            content=policy,
+        )
+    )
+    await writer.add_message(
+        Message.create(
+            trace_id="root",
+            role="user",
+            sequence=12,
+            parent_sequence=11,
+            content=continuation,
+        )
+    )
+    toolset_digest = canonical_sha256(["decide_script_task", "plan_script_tasks"]).wire
+    await writer.append_event(
+        "root",
+        "planner_policy_migrated",
+        {
+            "policy_version": "script-build-phase-two/v1",
+            "policy_digest": canonical_sha256(policy).wire,
+            "toolset_digest": toolset_digest,
+            "input_snapshot_id": current.snapshot_id,
+            "system_sequence": 11,
+            "user_sequence": 12,
+        },
+    )
+    if has_model_output:
+        await writer.add_message(
+            Message.create(
+                trace_id="root",
+                role="assistant",
+                sequence=13,
+                parent_sequence=12,
+                content={"text": "phase two already started"},
+            )
+        )
+
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=root_status,
+        blocked_reason=(
+            PHASE_ONE_CAPABILITY_BOUNDARY if root_status is TaskStatus.BLOCKED else None
+        ),
+    )
+    runner = _ReentryRunner(root)
+    runner.trace_store = FileSystemTraceStore(str(store_path))
+    coordinator = _Coordinator(root)
+    state = _State()
+    state.status = BuildStatus.PARTIAL if root_status is TaskStatus.BLOCKED else BuildStatus.RUNNING
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=coordinator,
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    if has_model_output:
+        with pytest.raises(MissionRecoveryRequired):
+            await service.run_phase_two(
+                7,
+                binding=binding,
+                snapshot=current,
+                direction_artifact_version_id=3,
+            )
+        assert runner.run_result_calls == 0
+        assert state.status is BuildStatus.RUNNING
+        return
+
+    await service.run_phase_two(
+        7,
+        binding=binding,
+        snapshot=current,
+        direction_artifact_version_id=3,
+    )
+    reloaded = FileSystemTraceStore(str(store_path))
+    messages = await reloaded.get_trace_messages("root")
+    events = await reloaded.get_events("root", 0)
+
+    assert [(item.sequence, item.role) for item in messages if item.sequence >= 11] == [
+        (11, "system"),
+        (12, "user"),
+    ]
+    assert len([item for item in events if item.get("event") == "planner_policy_migrated"]) == 1
+    assert coordinator.unblock_calls == (1 if root_status is TaskStatus.BLOCKED else 0)
+    assert state.status is BuildStatus.PARTIAL
+
+
+@pytest.mark.asyncio
+async def test_phase_two_reentry_rejects_unbacked_policy_event() -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.BLOCKED,
+        blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
+    )
+    runner = _ReentryRunner(root)
+    runner.trace_store.events = [_migration_event(runner, _snapshot())]
+    runner.trace_store.messages[0].content = "tampered policy"
+    state = _State()
+    state.status = BuildStatus.PARTIAL
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=_Coordinator(root),
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    with pytest.raises(PlannerPolicyMigrationRequired):
+        await service.run_phase_two(
+            7,
+            binding=_binding(),
+            snapshot=_snapshot(),
+            direction_artifact_version_id=3,
+        )
+
+    assert runner.run_result_calls == 0
+    assert state.status is BuildStatus.FAILED
+
+
+@pytest.mark.asyncio
+async def test_phase_two_reentry_recovers_messages_written_before_migration_event() -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.BLOCKED,
+        blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
+    )
+    runner = _ReentryRunner(root)
+    event = _migration_event(runner, _snapshot())
+    runner.trace_store.events = []
+    runner.trace_store.messages[0].parent_sequence = 10
+    runner.trace_store.messages[1].parent_sequence = 11
+    state = _State()
+    state.status = BuildStatus.PARTIAL
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=_Coordinator(root),
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    await service.run_phase_two(
+        7,
+        binding=_binding(),
+        snapshot=_snapshot(),
+        direction_artifact_version_id=3,
+    )
+
+    assert runner.run_result_calls == 1
+    assert len(runner.trace_store.events) == 1
+    assert runner.trace_store.events[0]["policy_digest"] == event["policy_digest"]
+    assert state.status is BuildStatus.PARTIAL
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("failure", ["duplicate", "toolset"])
+async def test_phase_two_reentry_rejects_ambiguous_policy_identity(failure: str) -> None:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.BLOCKED,
+        blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
+    )
+    runner = _ReentryRunner(root)
+    event = _migration_event(runner, _snapshot())
+    if failure == "duplicate":
+        runner.trace_store.events = [event, dict(event)]
+    else:
+        runner.trace_store.events = [{**event, "toolset_digest": "sha256:" + "0" * 64}]
+    state = _State()
+    state.status = BuildStatus.PARTIAL
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=_Coordinator(root),
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+
+    with pytest.raises(PlannerPolicyMigrationRequired):
+        await service.run_phase_two(
+            7,
+            binding=_binding(),
+            snapshot=_snapshot(),
+            direction_artifact_version_id=3,
+        )
+
+    assert runner.run_result_calls == 0
+    assert state.status is BuildStatus.FAILED
+
+
+@pytest.mark.asyncio
+async def test_direction_accept_may_cross_only_one_verified_prompt_snapshot_lineage() -> None:
+    parent = _snapshot()
+    business_digest = canonical_sha256(parent.to_input().business_payload()).wire
+    current = replace(
+        parent,
+        snapshot_id="12",
+        parent_snapshot_id=parent.snapshot_id,
+        parent_snapshot_sha256=parent.canonical_sha256,
+        business_input_sha256=business_digest,
+    )
+    spec = SimpleNamespace(
+        version=1,
+        context_refs=(
+            "script-build://task-kinds/direction",
+            f"script-build://inputs/{parent.snapshot_id}",
+        ),
+    )
+    direction = SimpleNamespace(
+        status=TaskStatus.COMPLETED,
+        current_spec=spec,
+        specs=[spec],
+        decision_ids=["direction-accept"],
+    )
+    ledger = SimpleNamespace(
+        tasks={"direction": direction},
+        decisions={
+            "direction-accept": SimpleNamespace(
+                action=DecisionAction.ACCEPT,
+                attempt_id="direction-attempt",
+            )
+        },
+        attempts={"direction-attempt": SimpleNamespace(spec_version=1)},
+    )
+
+    class Snapshots:
+        async def get(self, snapshot_id: str, *, script_build_id: int):
+            assert (snapshot_id, script_build_id) == ("11", 7)
+            return parent
+
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=ScriptMissionFactory(),
+        input_snapshots=Snapshots(),
+        bindings=SimpleNamespace(),
+        legacy_state=SimpleNamespace(),
+        authorizer=SimpleNamespace(),
+        direction_reconciler=SimpleNamespace(),
+    )
+    await service._verify_direction_snapshot_lineage(
+        script_build_id=7,
+        ledger=ledger,
+        current_snapshot=current,
+    )
+
+    tampered = replace(current, business_input_sha256="sha256:" + "9" * 64)
+    with pytest.raises(PhaseTwoBoundaryNotReady, match="immutable lineage"):
+        await service._verify_direction_snapshot_lineage(
+            script_build_id=7,
+            ledger=ledger,
+            current_snapshot=tampered,
+        )
+
+
+class _AdvanceSnapshots:
+    def __init__(self) -> None:
+        self.extend_calls = 0
+
+    async def get(self, _snapshot_id: str, *, script_build_id: int):
+        assert script_build_id == 7
+        return _snapshot()
+
+    async def extend_prompt_lineage(self, *_args, **_kwargs):
+        self.extend_calls += 1
+        return _snapshot()
+
+
+class _AdvanceState(_State):
+    def __init__(self, status: BuildStatus) -> None:
+        super().__init__()
+        self.status = status
+
+
+def _advance_service(
+    *,
+    status: BuildStatus = BuildStatus.PARTIAL,
+    blocked_reason: str = PHASE_ONE_CAPABILITY_BOUNDARY,
+    operation_status: OperationStatus | None = None,
+    publication_closed: bool = True,
+) -> tuple[ScriptMissionService, _AdvanceSnapshots]:
+    root = SimpleNamespace(
+        task_id="root-task",
+        status=TaskStatus.BLOCKED,
+        blocked_reason=blocked_reason,
+        decision_ids=["root-decision"],
+        current_spec=SimpleNamespace(context_refs=()),
+    )
+    operation = (
+        {"operation": SimpleNamespace(status=operation_status)}
+        if operation_status is not None
+        else {}
+    )
+    ledger = SimpleNamespace(
+        root_task_id="root-task",
+        tasks={"root-task": root},
+        decisions={
+            "root-decision": SimpleNamespace(
+                action=DecisionAction.BLOCK,
+                reason=blocked_reason,
+            )
+        },
+        operations=operation,
+    )
+
+    async def root_completion(_root: str):
+        return {"status": root.status.value, "blocked_reason": root.blocked_reason}
+
+    coordinator = SimpleNamespace(
+        root_completion=root_completion,
+        task_store=SimpleNamespace(load=lambda _root: _async(ledger)),
+    )
+    snapshots = _AdvanceSnapshots()
+    publication = (
+        SimpleNamespace(
+            state=PublicationState.PUBLISHED,
+            artifact_version_id=3,
+        )
+        if publication_closed
+        else None
+    )
+    reconciler = SimpleNamespace(
+        reconcile=lambda _build, _root: _async(3),
+        publications=SimpleNamespace(get_by_build=lambda *_args, **_kwargs: _async(publication)),
+    )
+    service = ScriptMissionService(
+        runner=SimpleNamespace(
+            trace_store=SimpleNamespace(
+                get_trace=lambda _root: _async(SimpleNamespace(agent_type="script_planner"))
+            )
+        ),
+        coordinator=coordinator,
+        factory=ScriptMissionFactory(),
+        input_snapshots=snapshots,
+        bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
+        legacy_state=_AdvanceState(status),
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
+        direction_reconciler=reconciler,
+        phase_two_prompt_requests=(SimpleNamespace(),),
+        phase_two_required_presets=("missing-phase-two-preset",),
+    )
+    service.run_phase_two = AsyncMock()  # type: ignore[method-assign]
+    return service, snapshots
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("status", "blocked_reason", "operation_status", "publication_closed"),
+    [
+        (BuildStatus.STOPPED, PHASE_ONE_CAPABILITY_BOUNDARY, None, True),
+        (BuildStatus.STOPPING, PHASE_ONE_CAPABILITY_BOUNDARY, None, True),
+        (BuildStatus.PARTIAL, "WRONG_BOUNDARY", None, True),
+        (BuildStatus.PARTIAL, PHASE_ONE_CAPABILITY_BOUNDARY, OperationStatus.RUNNING, True),
+        (BuildStatus.PARTIAL, PHASE_ONE_CAPABILITY_BOUNDARY, None, False),
+    ],
+)
+async def test_phase_two_advance_rejects_invalid_entry_without_extending_snapshot(
+    status: BuildStatus,
+    blocked_reason: str,
+    operation_status: OperationStatus | None,
+    publication_closed: bool,
+) -> None:
+    service, snapshots = _advance_service(
+        status=status,
+        blocked_reason=blocked_reason,
+        operation_status=operation_status,
+        publication_closed=publication_closed,
+    )
+
+    with pytest.raises(PhaseTwoBoundaryNotReady):
+        await service.advance_to_phase_two(7, Principal("owner"))
+
+    assert snapshots.extend_calls == 0
+
+
+@pytest.mark.asyncio
+async def test_concurrent_phase_two_advance_creates_only_one_active_run() -> None:
+    state = _AdvanceState(BuildStatus.PARTIAL)
+    release_prepare = asyncio.Event()
+    release_run = asyncio.Event()
+    prepare_calls = 0
+
+    async def prepare(_build: int):
+        nonlocal prepare_calls
+        prepare_calls += 1
+        await release_prepare.wait()
+        return _binding(), _snapshot(), 3
+
+    async def run_phase_two(*_args, **_kwargs):
+        await release_run.wait()
+
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(
+            root_completion=lambda _root: _async(
+                {
+                    "status": TaskStatus.BLOCKED.value,
+                    "blocked_reason": PHASE_ONE_CAPABILITY_BOUNDARY,
+                }
+            )
+        ),
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+    service._prepare_phase_two_transition = prepare  # type: ignore[method-assign]
+    service.run_phase_two = run_phase_two  # type: ignore[method-assign]
+
+    first = asyncio.create_task(service.advance_to_phase_two(7, Principal("owner")))
+    second = asyncio.create_task(service.advance_to_phase_two(7, Principal("owner")))
+    await asyncio.sleep(0)
+    release_prepare.set()
+    first_result, second_result = await asyncio.gather(first, second)
+
+    assert prepare_calls == 1
+    assert first_result.status is BuildStatus.RUNNING
+    assert second_result.status is BuildStatus.RUNNING
+    assert first_result.root_trace_id == second_result.root_trace_id == "root"
+    release_run.set()
+    await asyncio.sleep(0)
+
+
+@pytest.mark.asyncio
+async def test_http_advance_uses_safe_reentry_after_policy_was_unblocked() -> None:
+    state = _AdvanceState(BuildStatus.PARTIAL)
+    release_run = asyncio.Event()
+    coordinator = SimpleNamespace(
+        root_completion=lambda _root: _async(
+            {"status": TaskStatus.NEEDS_REPLAN.value, "blocked_reason": None}
+        )
+    )
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=coordinator,
+        factory=ScriptMissionFactory(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+    service._prepare_phase_two_transition = AsyncMock(  # type: ignore[method-assign]
+        side_effect=AssertionError("normal Phase1 boundary path must not be used")
+    )
+    service._prepare_phase_two_reentry = AsyncMock(  # type: ignore[method-assign]
+        return_value=(_binding(), _snapshot(), 3)
+    )
+
+    async def run_phase_two(*_args, **_kwargs):
+        await release_run.wait()
+
+    service.run_phase_two = run_phase_two  # type: ignore[method-assign]
+    result = await service.advance_to_phase_two(7, Principal("owner"))
+    assert result.status is BuildStatus.RUNNING
+    service._prepare_phase_two_reentry.assert_awaited_once_with(7)  # type: ignore[attr-defined]
+    service._prepare_phase_two_transition.assert_not_awaited()  # type: ignore[attr-defined]
+    release_run.set()
+    await asyncio.sleep(0)
+
+
+def test_mission_observation_dto_bounds_text_and_omits_internal_payloads() -> None:
+    secret = "must-not-leak"
+    artifact_ref = SimpleNamespace(
+        uri="script-build://artifact-versions/1",
+        kind="paragraph",
+        version="1",
+        digest="sha256:" + "a" * 64,
+    )
+    task = SimpleNamespace(
+        task_id="task-1",
+        parent_task_id="root",
+        display_path="root/task-1",
+        status=TaskStatus.COMPLETED,
+        current_spec=SimpleNamespace(
+            version=1,
+            objective="candidate body " * 500,
+            acceptance_criteria=[],
+            context_refs=("protected://" + secret,),
+        ),
+        child_task_ids=[],
+        attempt_ids=["attempt-1"],
+        validation_ids=[],
+        decision_ids=["decision-1"],
+        blocked_reason=None,
+        superseded_by=None,
+        created_at="now",
+        updated_at="now",
+    )
+    attempt = SimpleNamespace(
+        attempt_id="attempt-1",
+        task_id="task-1",
+        spec_version=1,
+        worker_trace_id="worker",
+        worker_preset="script_paragraph_worker",
+        status="submitted",
+        operation_id="operation-1",
+        snapshot_id="snapshot-1",
+        accepted_child_decision_ids=[],
+        submission=SimpleNamespace(
+            summary="candidate body " * 500,
+            artifact_refs=[artifact_ref],
+            evidence_refs=[],
+        ),
+        error=None,
+        created_at="now",
+        updated_at="now",
+        protected_context={"token": secret},
+    )
+    decision = SimpleNamespace(
+        decision_id="decision-1",
+        task_id="task-1",
+        action=DecisionAction.ACCEPT,
+        reason="accepted",
+        from_status="awaiting_decision",
+        to_status="completed",
+        attempt_id="attempt-1",
+        validation_id=None,
+        created_at="now",
+        payload={"tool_arguments": secret},
+    )
+    operation = SimpleNamespace(
+        operation_id="operation-1",
+        kind="dispatch",
+        status=OperationStatus.COMPLETED,
+        task_ids=["task-1"],
+        attempt_ids=["attempt-1"],
+        validation_ids=[],
+        deadline_at=None,
+        error=None,
+        created_at="now",
+        updated_at="now",
+        request={"arguments": secret},
+    )
+    ledger = SimpleNamespace(
+        revision=1,
+        root_trace_id="root",
+        root_task_id="root",
+        root_objective="root",
+        focused_task_id=None,
+        tasks={"task-1": task},
+        attempts={"attempt-1": attempt},
+        validations={},
+        decisions={"decision-1": decision},
+        operations={"operation-1": operation},
+        protected_context={"token": secret},
+    )
+
+    view = mission_snapshot_view(ledger)
+    rendered = str(view)
+    assert secret not in rendered
+    assert view["tasks"][0]["current_spec"]["context_refs"] == ["[internal-ref]"]
+    assert len(view["tasks"][0]["current_spec"]["objective"]) == 2_000
+    assert len(view["attempts"][0]["submission"]["summary"]) == 500

+ 1063 - 0
script_build_host/tests/test_phase_two_sql_flow.py

@@ -0,0 +1,1063 @@
+from __future__ import annotations
+
+import asyncio
+import re
+from collections.abc import Mapping, Sequence
+from datetime import UTC, datetime
+from pathlib import Path
+from typing import Any, cast
+
+import pytest
+from agent import FileSystemArtifactStore, FileSystemTaskStore, FileSystemTraceStore
+from agent.orchestration import (
+    AgentRole,
+    ArtifactRef,
+    AttemptSubmission,
+    CriterionResult,
+    DecisionAction,
+    TaskCoordinator,
+    TaskStatus,
+    ValidationVerdict,
+)
+from agent.orchestration.protocols import ValidatorRunResult, WorkerRunResult
+from sqlalchemy import event, func, select
+
+from script_build_host.application.phase_two_boundary import ScriptPhaseTwoBoundaryVerifier
+from script_build_host.application.phase_two_candidates import PhaseTwoCandidateService
+from script_build_host.application.phase_two_inputs import (
+    AcceptedInputResolver,
+    ActiveFrontierResolver,
+    StoredTaskContractReader,
+)
+from script_build_host.application.phase_two_planning import (
+    PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
+    PhaseTwoPlanningService,
+)
+from script_build_host.domain.artifacts import (
+    ArtifactKind,
+    DirectionGoal,
+    EvidenceRecordV1,
+    ScriptDirectionArtifactV1,
+)
+from script_build_host.domain.errors import PhaseTwoBoundaryNotReady
+from script_build_host.domain.phase_two_artifacts import (
+    CandidatePortfolioArtifactV1,
+    StructuredScriptArtifactV1,
+)
+from script_build_host.domain.records import BuildStatus, MissionBinding
+from script_build_host.domain.task_contracts import ScriptTaskKind
+from script_build_host.infrastructure.legacy_tables import (
+    script_build_element,
+    script_build_paragraph,
+    script_build_paragraph_element,
+    script_build_record,
+)
+from script_build_host.infrastructure.tables import (
+    artifact_version_table,
+    publication_table,
+)
+from script_build_host.infrastructure.task_contract_store import FileScriptTaskContractStore
+from script_build_host.repositories.legacy_state import SqlAlchemyLegacyBuildStateRepository
+from script_build_host.repositories.sqlalchemy import (
+    SqlAlchemyMissionBindingRepository,
+    SqlAlchemyScriptBusinessArtifactRepository,
+)
+from script_build_host.repositories.workspace import SqlAlchemyCandidateWorkspaceRepository
+
+ROOT = "phase-two-sql-root"
+SNAPSHOT_ID = 11
+FULL_SCOPE = "script-build://scopes/full"
+FULL_WRITE = "script-build://writes/full"
+
+
+def _contract(
+    kind: ScriptTaskKind,
+    *,
+    scope: str = FULL_SCOPE,
+    write_scope: Sequence[str] = (FULL_WRITE,),
+    input_refs: Sequence[dict[str, Any]] = (),
+    base_ref: ArtifactRef | None = None,
+    supersedes: Sequence[str] = (),
+    closure: Sequence[dict[str, Any]] = (),
+    adopted: Sequence[str] = (),
+    held: Sequence[str] = (),
+    order: Sequence[str] = (),
+    comparison_refs: Sequence[dict[str, Any]] = (),
+    execution_ready: bool = True,
+) -> dict[str, Any]:
+    output_schema = {
+        ScriptTaskKind.DIRECTION: "script-direction/v1",
+        ScriptTaskKind.DECODE_RETRIEVAL: "evidence-record/v1",
+        ScriptTaskKind.STRUCTURE: "structure-artifact/v1",
+        ScriptTaskKind.PARAGRAPH: "paragraph-patch/v1" if base_ref else "paragraph-artifact/v1",
+        ScriptTaskKind.ELEMENT_SET: (
+            "element-set-patch/v1" if base_ref else "element-set-artifact/v1"
+        ),
+        ScriptTaskKind.COMPARE: "comparison-artifact/v1",
+        ScriptTaskKind.COMPOSE: "structured-script/v1",
+        ScriptTaskKind.CANDIDATE_PORTFOLIO: "candidate-portfolio/v1",
+    }[kind]
+    if kind in {ScriptTaskKind.COMPOSE, ScriptTaskKind.CANDIDATE_PORTFOLIO}:
+        intent = "compose" if kind is ScriptTaskKind.COMPOSE else "portfolio"
+    elif supersedes:
+        intent = "replace"
+    else:
+        intent = "explore"
+    payload: dict[str, Any] = {
+        "schema_version": "script-task-contract/v1",
+        "task_kind": kind.value,
+        "scope_ref": scope,
+        "intent_class": intent,
+        "objective": f"produce independently verifiable {kind.value} content",
+        "input_decision_refs": list(input_refs),
+        "base_artifact_ref": _ref_payload(base_ref) if base_ref else None,
+        "write_scope": list(write_scope),
+        "gap_ref": None,
+        "output_schema": output_schema,
+        "criteria": [
+            {
+                "criterion_id": "closed",
+                "description": "the frozen increment is concrete and independently verifiable",
+                "hard": True,
+            }
+        ],
+        "budget": {
+            "max_attempts": 4,
+            "max_tokens": 32_000,
+            "max_seconds": 900,
+            "max_external_queries": 40,
+            "max_no_improvement": 3,
+        },
+        "supersedes_decision_ids": list(supersedes),
+        "candidate_closure_decision_refs": list(closure),
+        "adopted_decision_ids": list(adopted),
+        "held_or_rejected_decision_ids": list(held),
+        "compose_order": list(order),
+        "comparison_decision_refs": list(comparison_refs),
+    }
+    if not execution_ready and kind in {
+        ScriptTaskKind.COMPOSE,
+        ScriptTaskKind.CANDIDATE_PORTFOLIO,
+    }:
+        payload.update(
+            candidate_closure_decision_refs=[],
+            adopted_decision_ids=[],
+            held_or_rejected_decision_ids=[],
+            compose_order=[],
+        )
+    return payload
+
+
+def _decision_ref(
+    decision_id: str,
+    ref: ArtifactRef,
+    *,
+    scope: str,
+    kind: ScriptTaskKind,
+) -> dict[str, Any]:
+    return {
+        "decision_id": decision_id,
+        "artifact_ref": _ref_payload(ref),
+        "scope_ref": scope,
+        "expected_task_kind": kind.value,
+    }
+
+
+def _ref_payload(ref: ArtifactRef) -> dict[str, Any]:
+    return {
+        "uri": ref.uri,
+        "kind": ref.kind,
+        "version": ref.version,
+        "digest": ref.digest,
+        "summary": ref.summary,
+        "metadata": ref.metadata,
+    }
+
+
+def _task_kind(context: Mapping[str, Any]) -> ScriptTaskKind:
+    refs = cast(dict[str, Any], context["task_spec"])["context_refs"]
+    values = [value.rsplit("/", 1)[-1] for value in refs if "/task-kinds/" in value]
+    assert len(values) == 1
+    return ScriptTaskKind(values[0])
+
+
+class _ScriptedExecutor:
+    """Drive real Coordinator stages through the same protected Host services."""
+
+    def __init__(
+        self,
+        *,
+        coordinator: TaskCoordinator,
+        candidates: PhaseTwoCandidateService,
+        artifacts: SqlAlchemyScriptBusinessArtifactRepository,
+        script_build_id: int,
+    ) -> None:
+        self.coordinator = coordinator
+        self.candidates = candidates
+        self.artifacts = artifacts
+        self.script_build_id = script_build_id
+
+    async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
+        kind = _task_kind(context)
+        candidate_context = {
+            "root_trace_id": context["root_trace_id"],
+            "task_id": context["task_id"],
+            "attempt_id": context["attempt_id"],
+            "spec_version": context["spec_version"],
+        }
+        if kind is ScriptTaskKind.DECODE_RETRIEVAL:
+            await self.artifacts.freeze(
+                script_build_id=self.script_build_id,
+                task_id=context["task_id"],
+                attempt_id=context["attempt_id"],
+                spec_version=context["spec_version"],
+                artifact=EvidenceRecordV1(
+                    evidence_id=f"evidence-{context['attempt_id']}",
+                    source_type="decode",
+                    tool_name="retrieve_decode",
+                    query={"query": "opening tension", "top_k": 3},
+                    source_refs=("decode-index://fixture/1",),
+                    raw_artifact_ref=None,
+                    summary="A concrete opening needs a visible tension and an observable detail.",
+                    supports=("direction",),
+                    confidence="high",
+                    limitations=(),
+                    content_sha256="",
+                    created_at=datetime.now(UTC),
+                ),
+            )
+        elif kind is ScriptTaskKind.DIRECTION:
+            child_results = cast(list[dict[str, Any]], context["accepted_child_results"])
+            assert len(child_results) == 1
+            evidence_refs = child_results[0]["submission"]["evidence_refs"]
+            await self.artifacts.freeze(
+                script_build_id=self.script_build_id,
+                task_id=context["task_id"],
+                attempt_id=context["attempt_id"],
+                spec_version=context["spec_version"],
+                artifact=ScriptDirectionArtifactV1(
+                    goals=(DirectionGoal("goal-1", "Use a concrete reversal to reveal the topic"),),
+                    evidence_refs=(str(evidence_refs[0]["uri"]),),
+                    legacy_markdown=(
+                        "# Direction\n\nReveal the topic through one concrete reversal."
+                    ),
+                ),
+            )
+        elif kind in {ScriptTaskKind.STRUCTURE, ScriptTaskKind.PARAGRAPH}:
+            workspace = await self.candidates.read_attempt_workspace(context=candidate_context)
+            if workspace["paragraphs"]:
+                paragraph_id = workspace["paragraphs"][0]["paragraph_id"]
+            else:
+                created = await self.candidates.create_script_paragraph(
+                    payload={
+                        "paragraph_index": 1,
+                        "name": f"{kind.value} opening",
+                        "content_range": {"scope": "opening"},
+                    },
+                    context=candidate_context,
+                )
+                paragraph_id = created["paragraph_id"]
+            await self.candidates.batch_update_script_paragraphs(
+                updates=(
+                    {
+                        "paragraph_id": paragraph_id,
+                        "theme": "a specific reversal",
+                        "form": "contrast",
+                        "function": "hook",
+                        "feeling": "curiosity",
+                        "description": "The observable detail changes the initial interpretation.",
+                        "full_description": (
+                            "The opening states a familiar assumption, then overturns it with "
+                            "one visible and source-grounded detail."
+                        ),
+                    },
+                ),
+                context=candidate_context,
+            )
+        elif kind is ScriptTaskKind.ELEMENT_SET:
+            created = await self.candidates.create_script_element(
+                payload={
+                    "name": "observable reversal detail",
+                    "dimension_primary": "实质",
+                    "dimension_secondary": "narrative evidence",
+                    "topic_support": {"scope": "opening"},
+                    "weight_score": {"score": 0.9},
+                },
+                context=candidate_context,
+            )
+            workspace = await self.candidates.read_attempt_workspace(context=candidate_context)
+            if workspace["paragraphs"]:
+                await self.candidates.batch_link_paragraph_elements(
+                    links=(
+                        {
+                            "paragraph_id": workspace["paragraphs"][0]["paragraph_id"],
+                            "element_ids": [created["element_id"]],
+                        },
+                    ),
+                    context=candidate_context,
+                )
+        elif kind is ScriptTaskKind.COMPARE:
+            pinned = await self.candidates.read_pinned_candidates(context=candidate_context)
+            candidate_refs = [str(item["artifact_ref"]["uri"]) for item in pinned]
+            await self.candidates.save_comparison_candidate(
+                payload={
+                    "criterion_results": [
+                        {
+                            "criterion_id": "closed",
+                            "candidate_results": [
+                                {"artifact_ref": ref, "reason": "immutable candidate reviewed"}
+                                for ref in candidate_refs
+                            ],
+                        }
+                    ],
+                    "conflicts": ["the replacement has the more concrete opening"],
+                    "recommendation": candidate_refs[-1],
+                },
+                context=candidate_context,
+            )
+        elif kind is ScriptTaskKind.COMPOSE:
+            await self.candidates.save_structured_script_candidate(
+                acceptance_notes=("all adopted increments are pinned by ACCEPT decisions",),
+                context=candidate_context,
+            )
+        elif kind is ScriptTaskKind.CANDIDATE_PORTFOLIO:
+            await self.candidates.save_candidate_portfolio(
+                payload={"unresolved_defects": []}, context=candidate_context
+            )
+        else:  # pragma: no cover - this executor is intentionally phase-two bounded
+            raise AssertionError(f"unexpected worker kind: {kind.value}")
+
+        manifest = await self.candidates.resolve_attempt_manifest(context=candidate_context)
+        safe_ref = ArtifactRef(
+            manifest.artifact_ref.uri,
+            manifest.artifact_ref.kind,
+            manifest.artifact_ref.version,
+            manifest.artifact_ref.digest,
+        )
+        await self.coordinator.submit_attempt(
+            {
+                "role": AgentRole.WORKER.value,
+                "root_trace_id": context["root_trace_id"],
+                "task_id": context["task_id"],
+                "attempt_id": context["attempt_id"],
+                "spec_version": context["spec_version"],
+                "trace_id": context["worker_trace_id"],
+                "tool_call_id": f"submit:{context['attempt_id']}",
+                "operation_id": context["operation_id"],
+                "execution_epoch": context["execution_epoch"],
+            },
+            AttemptSubmission(
+                summary=f"kind={safe_ref.kind};status=frozen",
+                artifact_refs=[] if safe_ref.kind == ArtifactKind.EVIDENCE.value else [safe_ref],
+                evidence_refs=[safe_ref] if safe_ref.kind == ArtifactKind.EVIDENCE.value else [],
+            ),
+        )
+        return WorkerRunResult(context["worker_trace_id"], "completed")
+
+    async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
+        criteria = cast(dict[str, Any], context["task_spec"])["acceptance_criteria"]
+        results = [
+            CriterionResult(
+                criterion_id=item["criterion_id"],
+                verdict=ValidationVerdict.PASSED,
+                reason="the immutable SQL artifact is concrete and owner/digest closed",
+            )
+            for item in criteria
+        ]
+        await self.coordinator.submit_validation(
+            {
+                "role": AgentRole.VALIDATOR.value,
+                "root_trace_id": context["root_trace_id"],
+                "task_id": context["task_id"],
+                "attempt_id": context["attempt_id"],
+                "validation_id": context["validation_id"],
+                "snapshot_id": context["snapshot_id"],
+                "trace_id": context["validator_trace_id"],
+                "tool_call_id": f"validate:{context['validation_id']}",
+                "operation_id": context["operation_id"],
+                "execution_epoch": context["execution_epoch"],
+            },
+            ValidationVerdict.PASSED,
+            results,
+            "all hard criteria passed against the immutable snapshot",
+            (),
+            (),
+            (),
+            "accept",
+        )
+        return ValidatorRunResult(context["validator_trace_id"], "completed")
+
+    async def stop(self, trace_id: str) -> bool:
+        del trace_id
+        return True
+
+
+async def _plan_one(
+    planning: PhaseTwoPlanningService,
+    payload: Mapping[str, Any],
+    *,
+    parent_task_id: str | None,
+    call_id: str,
+) -> str:
+    result = await planning.plan_script_tasks(
+        contract_payloads=(payload,),
+        parent_task_id=parent_task_id,
+        context={"root_trace_id": ROOT, "tool_call_id": call_id},
+    )
+    return cast(str, result["task_ids"][0])
+
+
+async def _dispatch_accept(
+    planning: PhaseTwoPlanningService, coordinator: TaskCoordinator, task_id: str
+) -> tuple[str, ArtifactRef]:
+    cycles = await planning.dispatch_script_tasks(
+        task_ids=(task_id,),
+        context={"root_trace_id": ROOT, "tool_call_id": f"dispatch:{task_id}"},
+    )
+    assert cycles[0]["error"] is None
+    validation_id = cast(str, cycles[0]["validation_id"])
+    accepted = await planning.decide_script_task(
+        task_id=task_id,
+        action=DecisionAction.ACCEPT.value,
+        reason="independent validation passed",
+        validation_id=validation_id,
+        replacement_contract=None,
+        child_contracts=(),
+        context={"root_trace_id": ROOT, "tool_call_id": f"accept:{task_id}"},
+    )
+    decision_id = cast(str, accepted["decision_id"])
+    ledger = await coordinator.task_store.load(ROOT)
+    decision = ledger.decisions[decision_id]
+    attempt = ledger.attempts[cast(str, decision.attempt_id)]
+    assert attempt.submission is not None
+    refs = [*attempt.submission.artifact_refs, *attempt.submission.evidence_refs]
+    assert len(refs) == 1
+    return decision_id, refs[0]
+
+
+@pytest.mark.asyncio
+async def test_sql_phase_two_dynamic_replacement_compose_portfolio_and_boundary(
+    database: Any, tmp_path: Path
+) -> None:
+    engine, sessions = database
+    mutation_statements: list[str] = []
+
+    def observe_sql(
+        _connection: Any,
+        _cursor: Any,
+        statement: str,
+        _parameters: Any,
+        _context: Any,
+        _executemany: bool,
+    ) -> None:
+        if statement.lstrip().upper().startswith(("INSERT ", "UPDATE ", "DELETE ")):
+            mutation_statements.append(statement)
+
+    event.listen(engine.sync_engine, "before_cursor_execute", observe_sql)
+    build_states = SqlAlchemyLegacyBuildStateRepository(sessions)
+    script_build_id = await build_states.create(
+        execution_id=101,
+        topic_build_id=202,
+        topic_id=303,
+        agent_type="script-planner",
+        agent_config={},
+        data_source_url=None,
+        strategies_config={},
+    )
+    bindings = SqlAlchemyMissionBindingRepository(sessions)
+    await bindings.create(
+        script_build_id=script_build_id,
+        root_trace_id=ROOT,
+        input_snapshot_id=SNAPSHOT_ID,
+        engine_version="phase-two-test",
+        schema_version="v1",
+    )
+    task_path = tmp_path / "task-ledger"
+    artifact_path = tmp_path / "framework-artifacts"
+    contract_path = tmp_path / "contract-data"
+    task_store = FileSystemTaskStore(str(task_path))
+    framework_artifacts = FileSystemArtifactStore(str(artifact_path))
+    coordinator = TaskCoordinator(
+        task_store,
+        framework_artifacts,
+        FileSystemTraceStore(str(tmp_path / "traces")),
+    )
+    await coordinator.ensure_ledger(
+        ROOT,
+        {
+            "objective": "produce a governed phase-two candidate portfolio",
+            "acceptance_criteria": [
+                {"criterion_id": "portfolio", "description": "portfolio closed", "hard": True}
+            ],
+            "context_refs": [f"script-build://inputs/{SNAPSHOT_ID}"],
+        },
+    )
+    contract_store = FileScriptTaskContractStore(contract_path)
+    contract_reader = StoredTaskContractReader(contract_store)
+    artifacts = SqlAlchemyScriptBusinessArtifactRepository(sessions)
+    accepted_inputs = AcceptedInputResolver(
+        task_store=task_store,
+        artifact_store=framework_artifacts,
+        artifacts=artifacts,
+        contracts=contract_reader,
+    )
+    workspaces = SqlAlchemyCandidateWorkspaceRepository(sessions, artifacts)
+    candidates = PhaseTwoCandidateService(
+        bindings=bindings,
+        task_store=task_store,
+        framework_artifact_store=framework_artifacts,
+        artifacts=artifacts,
+        accepted_inputs=accepted_inputs,
+        active_frontier=ActiveFrontierResolver(),
+        workspaces=workspaces,
+    )
+    boundary = ScriptPhaseTwoBoundaryVerifier(
+        task_store=task_store,
+        artifact_store=framework_artifacts,
+        artifacts=artifacts,
+        bindings=bindings,
+        accepted_inputs=accepted_inputs,
+    )
+    planning = PhaseTwoPlanningService(
+        coordinator=coordinator,
+        bindings=bindings,
+        contracts=contract_store,
+        closure_gate=boundary,
+    )
+    coordinator.set_executor(
+        _ScriptedExecutor(
+            coordinator=coordinator,
+            candidates=candidates,
+            artifacts=artifacts,
+            script_build_id=script_build_id,
+        )
+    )
+
+    direction = await _plan_one(
+        planning, _contract(ScriptTaskKind.DIRECTION), parent_task_id=None, call_id="direction"
+    )
+    retrieval = await _plan_one(
+        planning,
+        _contract(ScriptTaskKind.DECODE_RETRIEVAL, write_scope=()),
+        parent_task_id=direction,
+        call_id="retrieval",
+    )
+    await _dispatch_accept(planning, coordinator, retrieval)
+    direction_decision, direction_ref = await _dispatch_accept(planning, coordinator, direction)
+
+    portfolio = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.CANDIDATE_PORTFOLIO,
+            write_scope=("script-build://writes",),
+            execution_ready=False,
+        ),
+        parent_task_id=None,
+        call_id="portfolio",
+    )
+    compose = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.COMPOSE,
+            write_scope=("script-build://writes",),
+            execution_ready=False,
+        ),
+        parent_task_id=portfolio,
+        call_id="compose",
+    )
+
+    # Exploration order is intentionally Element -> Paragraph -> Structure.
+    old_element = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.ELEMENT_SET,
+            scope="script-build://scopes/full/elements",
+            write_scope=("script-build://writes/elements",),
+        ),
+        parent_task_id=compose,
+        call_id="element-first",
+    )
+    old_element_decision, old_element_ref = await _dispatch_accept(
+        planning, coordinator, old_element
+    )
+    old_paragraph = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.PARAGRAPH,
+            scope="script-build://scopes/full/opening",
+            write_scope=("script-build://writes/paragraphs/opening",),
+        ),
+        parent_task_id=compose,
+        call_id="paragraph-second",
+    )
+    old_paragraph_decision, old_paragraph_ref = await _dispatch_accept(
+        planning, coordinator, old_paragraph
+    )
+    structure = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.STRUCTURE,
+            scope="script-build://scopes/full/opening",
+            write_scope=("script-build://writes/paragraphs/opening",),
+            input_refs=(
+                _decision_ref(
+                    old_paragraph_decision,
+                    old_paragraph_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.PARAGRAPH,
+                ),
+            ),
+            base_ref=old_paragraph_ref,
+        ),
+        parent_task_id=compose,
+        call_id="structure-third",
+    )
+    structure_decision, structure_ref = await _dispatch_accept(planning, coordinator, structure)
+
+    paragraph_replacement = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.PARAGRAPH,
+            scope="script-build://scopes/full/opening",
+            write_scope=("script-build://writes/paragraphs/opening",),
+            input_refs=(
+                _decision_ref(
+                    structure_decision,
+                    structure_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.STRUCTURE,
+                ),
+                _decision_ref(
+                    old_paragraph_decision,
+                    old_paragraph_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.PARAGRAPH,
+                ),
+            ),
+            base_ref=structure_ref,
+            supersedes=(old_paragraph_decision,),
+        ),
+        parent_task_id=compose,
+        call_id="replace-paragraph",
+    )
+    paragraph_decision, paragraph_ref = await _dispatch_accept(
+        planning, coordinator, paragraph_replacement
+    )
+    element_replacement = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.ELEMENT_SET,
+            scope="script-build://scopes/full/elements",
+            write_scope=("script-build://writes/elements",),
+            input_refs=(
+                _decision_ref(
+                    old_element_decision,
+                    old_element_ref,
+                    scope="script-build://scopes/full/elements",
+                    kind=ScriptTaskKind.ELEMENT_SET,
+                ),
+                _decision_ref(
+                    paragraph_decision,
+                    paragraph_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.PARAGRAPH,
+                ),
+            ),
+            base_ref=paragraph_ref,
+            supersedes=(old_element_decision,),
+        ),
+        parent_task_id=compose,
+        call_id="replace-element",
+    )
+    element_decision, element_ref = await _dispatch_accept(
+        planning, coordinator, element_replacement
+    )
+
+    comparison = await _plan_one(
+        planning,
+        _contract(
+            ScriptTaskKind.COMPARE,
+            scope="script-build://scopes/full/opening",
+            write_scope=(),
+            comparison_refs=(
+                _decision_ref(
+                    old_paragraph_decision,
+                    old_paragraph_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.PARAGRAPH,
+                ),
+                _decision_ref(
+                    paragraph_decision,
+                    paragraph_ref,
+                    scope="script-build://scopes/full/opening",
+                    kind=ScriptTaskKind.PARAGRAPH,
+                ),
+            ),
+        ),
+        parent_task_id=compose,
+        call_id="compare-opening-candidates",
+    )
+    comparison_decision, comparison_ref = await _dispatch_accept(planning, coordinator, comparison)
+
+    with pytest.raises(PhaseTwoBoundaryNotReady, match="CandidatePortfolio"):
+        await planning.decide_script_task(
+            task_id=(await task_store.load(ROOT)).root_task_id,
+            action=DecisionAction.BLOCK.value,
+            reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
+            validation_id=None,
+            replacement_contract=None,
+            child_contracts=(),
+            context={"root_trace_id": ROOT, "tool_call_id": "early-boundary"},
+        )
+
+    local_refs = (
+        _decision_ref(
+            old_element_decision,
+            old_element_ref,
+            scope="script-build://scopes/full/elements",
+            kind=ScriptTaskKind.ELEMENT_SET,
+        ),
+        _decision_ref(
+            old_paragraph_decision,
+            old_paragraph_ref,
+            scope="script-build://scopes/full/opening",
+            kind=ScriptTaskKind.PARAGRAPH,
+        ),
+        _decision_ref(
+            structure_decision,
+            structure_ref,
+            scope="script-build://scopes/full/opening",
+            kind=ScriptTaskKind.STRUCTURE,
+        ),
+        _decision_ref(
+            paragraph_decision,
+            paragraph_ref,
+            scope="script-build://scopes/full/opening",
+            kind=ScriptTaskKind.PARAGRAPH,
+        ),
+        _decision_ref(
+            element_decision,
+            element_ref,
+            scope="script-build://scopes/full/elements",
+            kind=ScriptTaskKind.ELEMENT_SET,
+        ),
+    )
+    adopted = (structure_decision, paragraph_decision, element_decision)
+    held = (old_element_decision, old_paragraph_decision)
+    compose_contract = _contract(
+        ScriptTaskKind.COMPOSE,
+        write_scope=("script-build://writes",),
+        input_refs=(
+            _decision_ref(
+                direction_decision,
+                direction_ref,
+                scope=FULL_SCOPE,
+                kind=ScriptTaskKind.DIRECTION,
+            ),
+        ),
+        closure=local_refs,
+        adopted=adopted,
+        held=held,
+        order=adopted,
+        comparison_refs=(
+            _decision_ref(
+                comparison_decision,
+                comparison_ref,
+                scope="script-build://scopes/full/opening",
+                kind=ScriptTaskKind.COMPARE,
+            ),
+        ),
+    )
+    await planning.decide_script_task(
+        task_id=compose,
+        action=DecisionAction.REVISE.value,
+        reason="freeze the selected active frontier and formal compose order",
+        validation_id=None,
+        replacement_contract=compose_contract,
+        child_contracts=(),
+        context={"root_trace_id": ROOT, "tool_call_id": "close-compose"},
+    )
+    compose_decision, compose_ref = await _dispatch_accept(planning, coordinator, compose)
+
+    portfolio_contract = _contract(
+        ScriptTaskKind.CANDIDATE_PORTFOLIO,
+        write_scope=("script-build://writes",),
+        closure=(
+            _decision_ref(
+                compose_decision,
+                compose_ref,
+                scope=FULL_SCOPE,
+                kind=ScriptTaskKind.COMPOSE,
+            ),
+        ),
+        adopted=(compose_decision,),
+        order=(compose_decision,),
+    )
+    await planning.decide_script_task(
+        task_id=portfolio,
+        action=DecisionAction.REVISE.value,
+        reason="freeze the uniquely adopted StructuredScript",
+        validation_id=None,
+        replacement_contract=portfolio_contract,
+        child_contracts=(),
+        context={"root_trace_id": ROOT, "tool_call_id": "close-portfolio"},
+    )
+    portfolio_decision, _ = await _dispatch_accept(planning, coordinator, portfolio)
+
+    ledger = await task_store.load(ROOT)
+    compose_attempt = ledger.attempts[cast(str, ledger.decisions[compose_decision].attempt_id)]
+    frontier = await candidates.read_active_frontier(
+        context={
+            "root_trace_id": ROOT,
+            "task_id": compose,
+            "attempt_id": compose_attempt.attempt_id,
+            "spec_version": compose_attempt.spec_version,
+        }
+    )
+    assert [item["decision_id"] for item in frontier["inputs"]] == list(adopted)
+    assert old_element_decision not in {item["decision_id"] for item in frontier["inputs"]}
+    assert old_paragraph_decision not in {item["decision_id"] for item in frontier["inputs"]}
+
+    await planning.decide_script_task(
+        task_id=ledger.root_task_id,
+        action=DecisionAction.BLOCK.value,
+        reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
+        validation_id=None,
+        replacement_contract=None,
+        child_contracts=(),
+        context={"root_trace_id": ROOT, "tool_call_id": "phase-two-boundary"},
+    )
+    await build_states.set_checkpoint(
+        script_build_id,
+        checkpoint_code=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
+        summary="one validated CandidatePortfolio is ready for phase three",
+        root_trace_id=ROOT,
+    )
+
+    reloaded = await FileSystemTaskStore(str(task_path)).load(ROOT)
+    assert reloaded.tasks[reloaded.root_task_id].status is TaskStatus.BLOCKED
+    assert reloaded.tasks[reloaded.root_task_id].blocked_reason == (
+        PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+    )
+    compose_children = [reloaded.tasks[item] for item in reloaded.tasks[compose].child_task_ids]
+    stored_reader = StoredTaskContractReader(FileScriptTaskContractStore(contract_path))
+    assert [
+        (
+            await stored_reader.read_for_task(
+                root_trace_id=ROOT,
+                task=item,
+                spec_version=item.current_spec_version,
+            )
+        ).task_kind
+        for item in compose_children[:3]
+    ] == [
+        ScriptTaskKind.ELEMENT_SET,
+        ScriptTaskKind.PARAGRAPH,
+        ScriptTaskKind.STRUCTURE,
+    ]
+
+    compose_version = await artifacts.get_by_attempt(
+        script_build_id=script_build_id,
+        task_id=compose,
+        attempt_id=compose_attempt.attempt_id,
+    )
+    assert isinstance(compose_version.artifact, StructuredScriptArtifactV1)
+    assert len(compose_version.artifact.paragraphs) == 1
+    assert compose_version.artifact.elements
+    assert compose_version.artifact.paragraph_element_links
+    portfolio_attempt = reloaded.attempts[
+        cast(str, reloaded.decisions[portfolio_decision].attempt_id)
+    ]
+    portfolio_version = await artifacts.get_by_attempt(
+        script_build_id=script_build_id,
+        task_id=portfolio,
+        attempt_id=portfolio_attempt.attempt_id,
+    )
+    assert isinstance(portfolio_version.artifact, CandidatePortfolioArtifactV1)
+    assert portfolio_version.artifact.accepted_decision_ids == (compose_decision,)
+    assert portfolio_version.artifact.adopted_structured_script_ref == compose_ref.uri
+
+    async with sessions() as session:
+        build = (
+            (
+                await session.execute(
+                    select(script_build_record).where(script_build_record.c.id == script_build_id)
+                )
+            )
+            .mappings()
+            .one()
+        )
+        branch_values = [
+            *(
+                await session.scalars(
+                    select(script_build_paragraph.c.branch_id).where(
+                        script_build_paragraph.c.script_build_id == script_build_id
+                    )
+                )
+            ),
+            *(
+                await session.scalars(
+                    select(script_build_element.c.branch_id).where(
+                        script_build_element.c.script_build_id == script_build_id
+                    )
+                )
+            ),
+            *(
+                await session.scalars(
+                    select(script_build_paragraph_element.c.branch_id).where(
+                        script_build_paragraph_element.c.script_build_id == script_build_id
+                    )
+                )
+            ),
+        ]
+        final_publications = await session.scalar(
+            select(func.count())
+            .select_from(publication_table)
+            .where(
+                publication_table.c.script_build_id == script_build_id,
+                publication_table.c.publication_type == "final",
+            )
+        )
+        phase_two_rows = (
+            (
+                await session.execute(
+                    select(artifact_version_table).where(
+                        artifact_version_table.c.script_build_id == script_build_id,
+                        artifact_version_table.c.artifact_type.in_(
+                            [
+                                ArtifactKind.STRUCTURE.value,
+                                ArtifactKind.PARAGRAPH.value,
+                                ArtifactKind.ELEMENT_SET.value,
+                                ArtifactKind.COMPARISON.value,
+                                ArtifactKind.STRUCTURED_SCRIPT.value,
+                                ArtifactKind.CANDIDATE_PORTFOLIO.value,
+                            ]
+                        ),
+                    )
+                )
+            )
+            .mappings()
+            .all()
+        )
+    mutation_targets = {
+        match.group(1).strip('`"').lower()
+        for statement in mutation_statements
+        if (
+            match := re.match(
+                r"\s*(?:INSERT\s+INTO|UPDATE|DELETE\s+FROM)\s+([`\"\w]+)",
+                statement,
+                re.IGNORECASE,
+            )
+        )
+    }
+    assert mutation_targets <= {
+        "script_build_record",
+        "script_build_mission_binding",
+        "script_build_artifact_version",
+        "script_build_publication",
+        "script_build_paragraph",
+        "script_build_element",
+        "script_build_paragraph_element",
+        "script_build_task_plan_step",
+    }
+    assert not mutation_targets & {
+        "script_build_branch",
+        "script_build_round",
+        "script_build_multipath",
+    }
+    assert BuildStatus(str(build["status"])) is BuildStatus.PARTIAL
+    assert build["error_message"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
+    assert build["reson_trace_id"] == ROOT and build["end_time"] is not None
+    assert branch_values and min(branch_values) > 0 and 0 not in branch_values
+    assert final_publications == 0
+    assert (
+        sum(
+            row["artifact_type"] == ArtifactKind.CANDIDATE_PORTFOLIO.value for row in phase_two_rows
+        )
+        == 1
+    )
+    for row in phase_two_rows:
+        if row["artifact_type"] in {
+            ArtifactKind.STRUCTURE.value,
+            ArtifactKind.PARAGRAPH.value,
+            ArtifactKind.ELEMENT_SET.value,
+            ArtifactKind.STRUCTURED_SCRIPT.value,
+        }:
+            assert row["legacy_branch_id"] == row["id"] > 0
+        else:
+            assert row["legacy_branch_id"] is None
+
+
+class _SingleBinding:
+    def __init__(self, root: str) -> None:
+        now = datetime.now(UTC)
+        self.value = MissionBinding(1, 1, root, 11, None, None, "test", "v1", now, now)
+
+    async def get_by_root(self, root_trace_id: str) -> MissionBinding:
+        assert root_trace_id == self.value.root_trace_id
+        return self.value
+
+
+class _BlockingExecutor:
+    def __init__(self) -> None:
+        self.started = asyncio.Event()
+        self.release = asyncio.Event()
+
+    async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
+        self.started.set()
+        await self.release.wait()
+        return WorkerRunResult(context["worker_trace_id"], "failed", error="test release")
+
+    async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
+        raise AssertionError(f"unexpected validation: {context}")
+
+    async def stop(self, trace_id: str) -> bool:
+        del trace_id
+        self.release.set()
+        return True
+
+
+@pytest.mark.asyncio
+async def test_same_task_cannot_reserve_parallel_attempts(tmp_path: Path) -> None:
+    root = "parallel-attempt-root"
+    task_store = FileSystemTaskStore(str(tmp_path / "ledger"))
+    executor = _BlockingExecutor()
+    coordinator = TaskCoordinator(
+        task_store,
+        FileSystemArtifactStore(str(tmp_path / "artifacts")),
+        FileSystemTraceStore(str(tmp_path / "traces")),
+        executor=executor,
+    )
+    await coordinator.ensure_ledger(
+        root,
+        {
+            "objective": "parallel attempt guard",
+            "acceptance_criteria": [
+                {"criterion_id": "guard", "description": "guard", "hard": True}
+            ],
+            "context_refs": ["script-build://inputs/11"],
+        },
+    )
+    planning = PhaseTwoPlanningService(
+        coordinator=coordinator,
+        bindings=cast(Any, _SingleBinding(root)),
+        contracts=FileScriptTaskContractStore(tmp_path / "contracts"),
+    )
+    task = await planning.plan_script_tasks(
+        contract_payloads=(_contract(ScriptTaskKind.DIRECTION),),
+        parent_task_id=None,
+        context={"root_trace_id": root},
+    )
+    task_id = task["task_ids"][0]
+    first = asyncio.create_task(
+        planning.dispatch_script_tasks(task_ids=(task_id,), context={"root_trace_id": root})
+    )
+    await asyncio.wait_for(executor.started.wait(), timeout=2)
+    second = await planning.dispatch_script_tasks(
+        task_ids=(task_id,), context={"root_trace_id": root}
+    )
+    assert "Dispatch conflict" in cast(str, second[0]["error"])
+    executor.release.set()
+    await first
+    ledger = await task_store.load(root)
+    assert len(ledger.tasks[task_id].attempt_ids) == 1

+ 7 - 1
script_build_host/tests/test_production_composition.py

@@ -68,7 +68,7 @@ def _settings(tmp_path: Path, **overrides: object) -> ScriptBuildSettings:
 async def test_production_composition_wires_real_stores_and_closes_resources(
     tmp_path: Path,
 ) -> None:
-    host = compose_production_host(_settings(tmp_path), _runtime(tmp_path))
+    host = compose_production_host(_settings(tmp_path, phase_two_max_tasks=7), _runtime(tmp_path))
     await host.validate_startup()
     manifest = await host.runtime_manifest.load()
     assert manifest["decode_index"]["file_count"] == 1
@@ -80,6 +80,7 @@ async def test_production_composition_wires_real_stores_and_closes_resources(
         host.composition.coordinator.task_store.base_path
         == (tmp_path / "agent-data" / "task-ledger").resolve()
     )
+    assert host.composition.phase_two_planning.limits.max_tasks == 7
     await host.close()
     await host.close()
     assert host.http_client.is_closed
@@ -99,3 +100,8 @@ def test_production_requires_decode_service_and_durable_trace_store(tmp_path: Pa
     runtime.runner.trace_store = None
     with pytest.raises(RuntimeError, match="durable trace store"):
         compose_production_host(_settings(tmp_path), runtime)
+
+    external = _runtime(tmp_path)
+    external.runner.trace_store = object()
+    with pytest.raises(RuntimeError, match="must provide a DurableTraceStoreVerifier"):
+        compose_production_host(_settings(tmp_path), external)