Ver Fonte

feat(script-build): recover phase three and preserve legacy api

SamLee há 1 dia atrás
pai
commit
fcf7466629

+ 4 - 4
script_build_host/src/script_build_host/adapters/uploaded_topic.py

@@ -104,7 +104,7 @@ class SqlUploadedTopicGateway:
         resolved_account = (account_name or parsed.get("account_name") or "").strip() or None
         now = datetime.now(UTC).replace(tzinfo=None)
         async with self._sessions() as session, session.begin():
-            build_result = cast(
+            build_result = cast(  # type: ignore[redundant-cast]
                 CursorResult[Any],
                 await session.execute(
                     insert(topic_build_record).values(
@@ -128,7 +128,7 @@ class SqlUploadedTopicGateway:
                 ),
             )
             topic_build_id = _inserted_id(build_result)
-            topic_result = cast(
+            topic_result = cast(  # type: ignore[redundant-cast]
                 CursorResult[Any],
                 await session.execute(
                     insert(topic_build_topic).values(
@@ -143,7 +143,7 @@ class SqlUploadedTopicGateway:
             topic_id = _inserted_id(topic_result)
             sort_order = 0
             for point_data in parsed["points"]:
-                point_result = cast(
+                point_result = cast(  # type: ignore[redundant-cast]
                     CursorResult[Any],
                     await session.execute(
                         insert(topic_build_point).values(
@@ -160,7 +160,7 @@ class SqlUploadedTopicGateway:
                 )
                 point_id = _inserted_id(point_result)
                 for item in point_data["items"]:
-                    item_result = cast(
+                    item_result = cast(  # type: ignore[redundant-cast]
                         CursorResult[Any],
                         await session.execute(
                             insert(topic_build_composition_item).values(

+ 46 - 1
script_build_host/src/script_build_host/api/app.py

@@ -4,7 +4,8 @@ from __future__ import annotations
 
 from typing import Any
 
-from fastapi import FastAPI, Request
+from fastapi import FastAPI, HTTPException, Request
+from fastapi.exceptions import RequestValidationError
 from fastapi.responses import JSONResponse
 
 from script_build_host.domain.errors import ScriptBuildError
@@ -36,6 +37,50 @@ def create_app(**router_dependencies: Any) -> FastAPI:
             },
         )
 
+    @app.exception_handler(HTTPException)
+    async def http_error(_: Request, exc: HTTPException) -> JSONResponse:
+        detail: dict[str, Any] = exc.detail if isinstance(exc.detail, dict) else {}
+        code = str(detail.get("error_code") or f"HTTP_{exc.status_code}")
+        message = str(detail.get("message") or code.replace("_", " ").lower())
+        return JSONResponse(
+            status_code=exc.status_code,
+            content={
+                "success": False,
+                "error_code": code,
+                "message": message,
+                # Compatibility overlay for callers that previously consumed
+                # FastAPI's nested HTTPException body.
+                "detail": {**detail, "error_code": code},
+            },
+            headers=exc.headers,
+        )
+
+    @app.exception_handler(RequestValidationError)
+    async def validation_error(_: Request, exc: RequestValidationError) -> JSONResponse:
+        issues = [
+            {"location": list(item["loc"]), "type": item["type"]} for item in exc.errors()[:20]
+        ]
+        return JSONResponse(
+            status_code=422,
+            content={
+                "success": False,
+                "error_code": "REQUEST_VALIDATION_FAILED",
+                "message": "request validation failed",
+                "issues": issues,
+            },
+        )
+
+    @app.exception_handler(Exception)
+    async def internal_error(_: Request, __: Exception) -> JSONResponse:
+        return JSONResponse(
+            status_code=500,
+            content={
+                "success": False,
+                "error_code": "INTERNAL_SERVER_ERROR",
+                "message": "an internal error occurred",
+            },
+        )
+
     app.include_router(create_script_build_router(**router_dependencies))
     return app
 

+ 465 - 58
script_build_host/src/script_build_host/api/routes.py

@@ -3,12 +3,14 @@
 import asyncio
 import json
 import re
-from collections.abc import Mapping
-from dataclasses import asdict
+from collections.abc import Awaitable, Callable, Mapping
+from dataclasses import asdict, replace
 from typing import Annotated, Any, Protocol, cast
 from urllib.parse import urlsplit, urlunsplit
+from uuid import uuid4
 
-from fastapi import APIRouter, Depends, Query, Request, WebSocket, WebSocketDisconnect
+from fastapi import APIRouter, Depends, Header, Query, Request, WebSocket, WebSocketDisconnect
+from fastapi.responses import HTMLResponse
 
 from script_build_host.agents.prompt_catalog import script_prompt_requests
 from script_build_host.application.mission_service import (
@@ -19,6 +21,7 @@ from script_build_host.application.observation_views import (
     mission_snapshot_view,
     task_detail_view,
 )
+from script_build_host.domain.errors import BuildNotFound
 from script_build_host.domain.ports import (
     MissionBindingRepository,
     PublicationRepository,
@@ -58,6 +61,12 @@ def create_script_build_router(
     outbound_policy: OutboundPolicy | None = None,
     default_model_manifest: Mapping[str, object] | None = None,
     default_datasource_manifest: Mapping[str, object] | None = None,
+    phase_three: Any | None = None,
+    finalization: Any | None = None,
+    recovery: Any | None = None,
+    legacy_projection: Any | None = None,
+    legacy_api: Any | None = None,
+    http_command_journal: Any | None = None,
 ) -> APIRouter:
     router = APIRouter(tags=["script-build"])
 
@@ -70,10 +79,33 @@ def create_script_build_router(
         await security.authorize(script_build_id, actor)
         return await bindings.get_by_build(script_build_id)
 
+    async def execute_mutation(
+        *,
+        actor: Principal,
+        route_family: str,
+        idempotency_key: str | None,
+        request_payload: Mapping[str, Any],
+        command: Callable[[], Awaitable[tuple[int, dict[str, Any]]]],
+        resource_id: int | None = None,
+    ) -> dict[str, Any]:
+        if idempotency_key and http_command_journal is not None:
+            _, payload = await http_command_journal.execute(
+                principal=actor,
+                route_family=route_family,
+                idempotency_key=idempotency_key,
+                request_payload=request_payload,
+                command=command,
+                resource_id=resource_id,
+            )
+            return cast(dict[str, Any], payload)
+        _, payload = await command()
+        return payload
+
     @router.post("/api/pattern/script_builds", response_model=ScriptBuildStartResponse)
     async def start_script_build(
         body: ScriptBuildRequest,
         actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
     ) -> ScriptBuildStartResponse:
         if body.agent_type != "AigcAgent":
             raise problem(400, "UNSUPPORTED_AGENT_TYPE", "script build supports AigcAgent only")
@@ -90,11 +122,84 @@ def create_script_build_router(
             default_model_manifest=default_model_manifest,
             default_datasource_manifest=default_datasource_manifest,
         )
-        result = await mission_service.start(command)
-        return ScriptBuildStartResponse(
-            script_build_id=result.script_build_id,
-            status=result.status.value,
-            root_trace_id=result.root_trace_id,
+        preallocated_root_trace_id = str(uuid4()) if idempotency_key else None
+        if preallocated_root_trace_id is not None:
+            command = replace(command, root_trace_id=preallocated_root_trace_id)
+
+        async def execute_with_root(root_trace_id: str | None) -> tuple[int, dict[str, Any]]:
+            result = await mission_service.start(
+                replace(command, root_trace_id=root_trace_id or command.root_trace_id)
+            )
+            return (
+                200,
+                {
+                    "script_build_id": result.script_build_id,
+                    "status": result.status.value,
+                    "root_trace_id": result.root_trace_id,
+                },
+            )
+
+        async def execute() -> tuple[int, dict[str, Any]]:
+            return await execute_with_root(preallocated_root_trace_id)
+
+        async def resume_reserved(
+            root_trace_id: str | None, _resource_id: int | None
+        ) -> tuple[int, dict[str, Any]]:
+            if root_trace_id is None:
+                raise problem(409, "IDEMPOTENCY_RECOVERY_REQUIRED", "reserved start has no Root")
+            try:
+                result = await mission_service.start_result_by_root(root_trace_id)
+            except BuildNotFound:
+                return await execute_with_root(root_trace_id)
+            return 200, {
+                "script_build_id": result.script_build_id,
+                "status": result.status.value,
+                "root_trace_id": result.root_trace_id,
+            }
+
+        if idempotency_key and http_command_journal is not None:
+            _, payload = await http_command_journal.execute(
+                principal=actor,
+                route_family="script-build:start",
+                idempotency_key=idempotency_key,
+                request_payload=body.model_dump(mode="json"),
+                command=execute,
+                resume_reserved=resume_reserved,
+                preallocated_root_trace_id=preallocated_root_trace_id,
+            )
+        else:
+            _, payload = await execute()
+        return ScriptBuildStartResponse(**payload)
+
+    @router.get("/api/pattern/script_builds")
+    async def list_script_builds(
+        actor: PrincipalDependency,
+        page: int = Query(1, ge=1),
+        page_size: int = Query(20, ge=1, le=100),
+        status: str | None = None,
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "legacy list is unavailable")
+        return cast(
+            dict[str, Any],
+            await legacy_api.list_authorized(actor, page=page, page_size=page_size, status=status),
+        )
+
+    @router.get("/api/pattern/script_builds/overview")
+    async def script_build_overview(actor: PrincipalDependency) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "overview is unavailable")
+        return cast(dict[str, Any], await legacy_api.overview(actor))
+
+    @router.get("/api/pattern/script_builds/overview/{topic_build_id}")
+    async def script_build_topic_overview(
+        topic_build_id: int, actor: PrincipalDependency
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "overview is unavailable")
+        return cast(
+            dict[str, Any],
+            await legacy_api.overview(actor, topic_build_id=topic_build_id),
         )
 
     @router.post("/api/pattern/script_builds/parse_topic_json")
@@ -132,6 +237,7 @@ def create_script_build_router(
     async def start_from_topic_json(
         body: ScriptBuildFromTopicJsonRequest,
         actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
     ) -> dict[str, Any]:
         if uploaded_topics is None:
             raise problem(
@@ -143,68 +249,331 @@ def create_script_build_router(
             topic_build_id=0,
             topic_id=0,
         )
-        parsed = await uploaded_topics.parse(body.json_data)
-        created = await uploaded_topics.create(parsed, account_name=body.account_name)
-        request = ScriptBuildRequest(
-            execution_id=0,
-            topic_build_id=int(created["topic_build_id"]),
-            topic_id=int(created["topic_id"]),
-            agent_config=body.agent_config,
-            data_source_url=body.data_source_url,
-            strategies_always_on=body.strategies_always_on,
-            strategies_on_demand=body.strategies_on_demand,
-        )
-        await security.authorize_source(
-            actor,
-            execution_id=request.execution_id,
-            topic_build_id=request.topic_build_id,
-            topic_id=request.topic_id,
-        )
-        result = await mission_service.start(
-            await _start_command(
-                request,
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            parsed = await uploaded_topics.parse(body.json_data)
+            created = await uploaded_topics.create(parsed, account_name=body.account_name)
+            request = ScriptBuildRequest(
+                execution_id=0,
+                topic_build_id=int(created["topic_build_id"]),
+                topic_id=int(created["topic_id"]),
+                agent_config=body.agent_config,
+                data_source_url=body.data_source_url,
+                strategies_always_on=body.strategies_always_on,
+                strategies_on_demand=body.strategies_on_demand,
+            )
+            await security.authorize_source(
                 actor,
-                outbound_policy,
-                default_model_manifest=default_model_manifest,
-                default_datasource_manifest=default_datasource_manifest,
+                execution_id=request.execution_id,
+                topic_build_id=request.topic_build_id,
+                topic_id=request.topic_id,
             )
+            result = await mission_service.start(
+                await _start_command(
+                    request,
+                    actor,
+                    outbound_policy,
+                    default_model_manifest=default_model_manifest,
+                    default_datasource_manifest=default_datasource_manifest,
+                )
+            )
+            return 200, {
+                "success": True,
+                "message": "脚本构建任务已提交",
+                "script_build_id": result.script_build_id,
+                "topic_id": created["topic_id"],
+                "topic_build_id": created["topic_build_id"],
+                "account_name": created.get("account_name"),
+                "item_count": created.get("item_count", 0),
+                "point_count": created.get("point_count", 0),
+                "status": result.status.value,
+                "root_trace_id": result.root_trace_id,
+            }
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:from-topic-json",
+            idempotency_key=idempotency_key,
+            request_payload=body.model_dump(mode="json"),
+            command=command,
         )
-        return {
-            "success": True,
-            "message": "脚本构建任务已提交",
-            "script_build_id": result.script_build_id,
-            "topic_id": created["topic_id"],
-            "topic_build_id": created["topic_build_id"],
-            "account_name": created.get("account_name"),
-            "item_count": created.get("item_count", 0),
-            "point_count": created.get("point_count", 0),
-            "status": result.status.value,
-            "root_trace_id": result.root_trace_id,
-        }
 
     @router.post("/api/pattern/script_builds/{script_build_id}/stop")
     async def stop_script_build(
         script_build_id: int,
         actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
     ) -> dict[str, Any]:
         await security.authorize(script_build_id, actor)
-        result = await mission_service.stop(script_build_id, actor)
-        return {"success": True, **asdict(result), "status": result.status.value}
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await mission_service.stop(script_build_id, actor)
+            return 200, {"success": True, **asdict(result), "status": result.status.value}
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:stop",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.get("/api/pattern/script_builds/{script_build_id}/prefill")
+    async def script_build_prefill(
+        script_build_id: int, actor: PrincipalDependency
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "prefill is unavailable")
+        await security.authorize(script_build_id, actor)
+        return cast(dict[str, Any], await legacy_api.prefill(script_build_id, actor))
+
+    @router.get("/api/pattern/script_builds/{script_build_id}/log")
+    async def script_build_log(script_build_id: int, actor: PrincipalDependency) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "build log is unavailable")
+        await security.authorize(script_build_id, actor)
+        return cast(dict[str, Any], await legacy_api.log(script_build_id, actor))
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/favorite")
+    @router.patch("/api/pattern/script_builds/{script_build_id}/favorite")
+    async def favorite_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        is_favorited: bool = True,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "favorite is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            return 200, cast(
+                dict[str, Any],
+                await legacy_api.set_favorite(script_build_id, actor, value=is_favorited),
+            )
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:favorite",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id, "is_favorited": is_favorited},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.delete("/api/pattern/script_builds/{script_build_id}")
+    async def delete_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "delete is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            return 200, cast(dict[str, Any], await legacy_api.delete(script_build_id, actor))
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:delete",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/retry")
+    async def retry_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "retry is unavailable")
+        await security.authorize(script_build_id, actor)
+        prefill = await legacy_api.prefill(script_build_id, actor)
+        strategies = prefill.get("strategies_config") or {}
+        request = ScriptBuildRequest(
+            execution_id=int(prefill["execution_id"]),
+            topic_build_id=int(prefill["topic_build_id"]),
+            topic_id=int(prefill["topic_id"]),
+            agent_type=str(prefill.get("agent_type") or "AigcAgent"),
+            agent_config=prefill.get("agent_config"),
+            data_source_url=prefill.get("data_source_url"),
+            strategies_always_on=list(strategies.get("always_on", [])),
+            strategies_on_demand=list(strategies.get("on_demand", [])),
+        )
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await mission_service.start(
+                await _start_command(
+                    request,
+                    actor,
+                    outbound_policy,
+                    default_model_manifest=default_model_manifest,
+                    default_datasource_manifest=default_datasource_manifest,
+                )
+            )
+            return 200, {
+                "success": True,
+                "script_build_id": result.script_build_id,
+                "status": result.status.value,
+                "root_trace_id": result.root_trace_id,
+                "retried_from": script_build_id,
+            }
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:retry",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
 
     @router.post("/api/pattern/script_builds/{script_build_id}/phase-two/advance")
     async def advance_to_phase_two(
         script_build_id: int,
         actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
     ) -> dict[str, Any]:
         await security.authorize(script_build_id, actor)
-        result = await mission_service.advance_to_phase_two(script_build_id, actor)
-        return {
-            "success": True,
-            "script_build_id": result.script_build_id,
-            "status": result.status.value,
-            "root_trace_id": result.root_trace_id,
-            "input_snapshot_id": result.input_snapshot_id,
-        }
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await mission_service.advance_to_phase_two(script_build_id, actor)
+            return 200, {
+                "success": True,
+                "script_build_id": result.script_build_id,
+                "status": result.status.value,
+                "root_trace_id": result.root_trace_id,
+                "input_snapshot_id": result.input_snapshot_id,
+            }
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:phase-two-advance",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/phase-three/advance")
+    async def advance_to_phase_three(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if phase_three is None:
+            raise problem(503, "PHASE_THREE_NOT_CONFIGURED", "phase three is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await phase_three.advance(script_build_id, actor)
+            return 200, {
+                "success": True,
+                "script_build_id": result.script_build_id,
+                "status": result.status.value,
+                "root_trace_id": result.root_trace_id,
+                "input_snapshot_id": result.input_snapshot_id,
+            }
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:phase-three-advance",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/resume")
+    async def resume_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if recovery is None:
+            raise problem(503, "RECOVERY_NOT_CONFIGURED", "recovery is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await recovery.resume(script_build_id, actor)
+            return 200, {"success": True, **_json_values(asdict(result))}
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:resume",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/finalize")
+    async def finalize_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if finalization is None:
+            raise problem(503, "FINAL_PUBLISHER_NOT_CONFIGURED", "publisher is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await finalization.finalize(script_build_id, actor)
+            return 200, {"success": True, **_json_values(asdict(result))}
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:finalize",
+            idempotency_key=idempotency_key,
+            request_payload={"script_build_id": script_build_id},
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.post("/api/pattern/script_builds/{script_build_id}/reconcile")
+    async def reconcile_script_build(
+        script_build_id: int,
+        actor: PrincipalDependency,
+        retry_finalize: bool = False,
+        idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
+    ) -> dict[str, Any]:
+        if recovery is None:
+            raise problem(503, "RECOVERY_NOT_CONFIGURED", "recovery is unavailable")
+        await security.authorize(script_build_id, actor)
+
+        async def command() -> tuple[int, dict[str, Any]]:
+            result = await recovery.reconcile(script_build_id, actor, retry_finalize=retry_finalize)
+            return 200, {"success": True, **_json_values(asdict(result))}
+
+        return await execute_mutation(
+            actor=actor,
+            route_family="script-build:reconcile",
+            idempotency_key=idempotency_key,
+            request_payload={
+                "script_build_id": script_build_id,
+                "retry_finalize": retry_finalize,
+            },
+            command=command,
+            resource_id=script_build_id,
+        )
+
+    @router.get("/api/pattern/script_builds/{script_build_id}")
+    async def legacy_detail(
+        script_build_id: int,
+        actor: PrincipalDependency,
+    ) -> dict[str, Any]:
+        if legacy_projection is None:
+            raise problem(503, "LEGACY_PROJECTION_NOT_CONFIGURED", "detail is unavailable")
+        await security.authorize(script_build_id, actor)
+        return cast(
+            dict[str, Any],
+            await legacy_projection.project_authorized(script_build_id, principal=actor),
+        )
 
     @router.get("/api/pattern/script_builds/{script_build_id}/mission")
     async def mission_snapshot(
@@ -260,11 +629,18 @@ def create_script_build_router(
     async def publication_detail(
         script_build_id: int,
         actor: PrincipalDependency,
+        publication_type: str = Query("direction", alias="type"),
     ) -> dict[str, Any]:
         await scope(script_build_id, actor)
-        value = await publications.get_by_build(
-            script_build_id, publication_type=PublicationType.DIRECTION
-        )
+        try:
+            selected_type = PublicationType(publication_type)
+        except ValueError as exc:
+            raise problem(
+                422,
+                "INVALID_PUBLICATION_TYPE",
+                "publication type must be direction or final",
+            ) from exc
+        value = await publications.get_by_build(script_build_id, publication_type=selected_type)
         if value is None:
             raise problem(404, "PUBLICATION_NOT_FOUND", "publication does not exist")
         return _json_values(asdict(value))
@@ -347,6 +723,30 @@ def create_script_build_router(
         messages = await trace_store.get_trace_messages(trace_id)
         return {"messages": [item.to_dict() for item in messages]}
 
+    @router.get("/api/pattern/script_builds/{script_build_id}/trace_messages")
+    async def legacy_trace_messages(
+        script_build_id: int, actor: PrincipalDependency
+    ) -> dict[str, Any]:
+        if legacy_api is None:
+            raise problem(503, "LEGACY_API_NOT_CONFIGURED", "trace messages are unavailable")
+        await security.authorize(script_build_id, actor)
+        trace_id = await legacy_api.root_trace_id(script_build_id, actor)
+        messages = await trace_store.get_trace_messages(trace_id) if trace_id else []
+        return {
+            "success": True,
+            "script_build_id": script_build_id,
+            "trace_id": trace_id,
+            "messages": [item.to_dict() for item in messages],
+        }
+
+    @router.get("/script_build_detail/{script_build_id}", response_class=HTMLResponse)
+    async def legacy_detail_entry(script_build_id: int, actor: PrincipalDependency) -> HTMLResponse:
+        await security.authorize(script_build_id, actor)
+        return HTMLResponse(
+            '<!doctype html><html><body><div id="script-build-detail" '
+            f'data-script-build-id="{script_build_id}"></div></body></html>'
+        )
+
     @router.websocket("/api/pattern/script_builds/{script_build_id}/traces/{trace_id}/watch")
     async def watch_trace(
         websocket: WebSocket,
@@ -365,9 +765,16 @@ def create_script_build_router(
         since = 0
         try:
             while True:
-                events = await trace_store.get_events(trace_id, since)
+                try:
+                    actor = await security.principal(websocket)
+                    binding = await scope(script_build_id, actor)
+                    await _owned_trace(trace_store, binding.root_trace_id, trace_id)
+                except Exception:
+                    await websocket.close(code=4404)
+                    return
+                events = (await trace_store.get_events(trace_id, since))[:100]
                 for event in events:
-                    await websocket.send_json(event)
+                    await asyncio.wait_for(websocket.send_json(event), timeout=2.0)
                     if isinstance(event.get("event_id"), int):
                         since = max(since, event["event_id"])
                 try:

+ 195 - 0
script_build_host/src/script_build_host/application/legacy_api.py

@@ -0,0 +1,195 @@
+"""Bounded legacy list/overview/mutation projections."""
+
+from __future__ import annotations
+
+from collections import Counter
+from typing import Any
+
+from sqlalchemy import select, update
+from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
+
+from script_build_host.domain.errors import BuildNotFound
+from script_build_host.domain.ports import BuildAuthorizer
+from script_build_host.domain.records import Principal
+from script_build_host.infrastructure.legacy_tables import (
+    script_build_log,
+    script_build_record,
+)
+
+
+class LegacyScriptBuildApiService:
+    def __init__(
+        self,
+        read_sessions: async_sessionmaker[AsyncSession],
+        write_sessions: async_sessionmaker[AsyncSession],
+        authorizer: BuildAuthorizer,
+    ) -> None:
+        self._read = read_sessions
+        self._write = write_sessions
+        self._authorizer = authorizer
+
+    async def list_authorized(
+        self,
+        principal: Principal,
+        *,
+        page: int,
+        page_size: int,
+        status: str | None = None,
+        topic_build_id: int | None = None,
+    ) -> dict[str, Any]:
+        authorized = await self._authorized_rows(
+            principal, status=status, topic_build_id=topic_build_id
+        )
+        total = len(authorized)
+        start = (page - 1) * page_size
+        items = [_list_item(row) for row in authorized[start : start + page_size]]
+        return {
+            "success": True,
+            "items": items,
+            "total": total,
+            "page": page,
+            "page_size": page_size,
+        }
+
+    async def _authorized_rows(
+        self,
+        principal: Principal,
+        *,
+        status: str | None = None,
+        topic_build_id: int | None = None,
+    ) -> list[Any]:
+        async with self._read() as session:
+            statement = select(script_build_record).where(
+                script_build_record.c.is_deleted.is_(False)
+            )
+            if status is not None:
+                statement = statement.where(script_build_record.c.status == status)
+            if topic_build_id is not None:
+                statement = statement.where(script_build_record.c.topic_build_id == topic_build_id)
+            rows = list((await session.execute(statement)).mappings().all())
+        authorized: list[Any] = []
+        for row in rows:
+            try:
+                await self._authorizer.require_access(principal, int(row["id"]))
+            except Exception:
+                continue
+            authorized.append(row)
+        authorized.sort(key=lambda row: int(row["id"]), reverse=True)
+        return authorized
+
+    async def overview(
+        self, principal: Principal, *, topic_build_id: int | None = None
+    ) -> dict[str, Any]:
+        rows = await self._authorized_rows(principal, topic_build_id=topic_build_id)
+        counts = Counter(str(item["status"]) for item in rows)
+        return {"success": True, "total": len(rows), "status_counts": dict(counts)}
+
+    async def prefill(self, script_build_id: int, principal: Principal) -> dict[str, Any]:
+        row = await self._row(script_build_id, principal)
+        return {
+            "success": True,
+            "execution_id": int(row["execution_id"]),
+            "topic_build_id": int(row["topic_build_id"]),
+            "topic_id": int(row["topic_id"]),
+            "agent_type": row["agent_type"],
+            "agent_config": row["agent_config"],
+            "data_source_url": row["data_source_url"],
+            "strategies_config": row["strategies_config"],
+        }
+
+    async def set_favorite(
+        self, script_build_id: int, principal: Principal, *, value: bool
+    ) -> dict[str, Any]:
+        await self._authorizer.require_access(principal, script_build_id)
+        async with self._write() as session, session.begin():
+            result = await session.execute(
+                update(script_build_record)
+                .where(
+                    script_build_record.c.id == script_build_id,
+                    script_build_record.c.is_deleted.is_(False),
+                )
+                .values(is_favorited=value)
+            )
+            if result.rowcount != 1:
+                raise BuildNotFound()
+        return {"success": True, "script_build_id": script_build_id, "is_favorited": value}
+
+    async def delete(self, script_build_id: int, principal: Principal) -> dict[str, Any]:
+        await self._authorizer.require_access(principal, script_build_id)
+        async with self._write() as session, session.begin():
+            result = await session.execute(
+                update(script_build_record)
+                .where(
+                    script_build_record.c.id == script_build_id,
+                    script_build_record.c.is_deleted.is_(False),
+                )
+                .values(is_deleted=True)
+            )
+            if result.rowcount != 1:
+                raise BuildNotFound()
+        return {"success": True, "script_build_id": script_build_id}
+
+    async def log(self, script_build_id: int, principal: Principal) -> dict[str, Any]:
+        row = await self._row(script_build_id, principal)
+        async with self._read() as session:
+            log_rows = list(
+                (
+                    await session.execute(
+                        select(script_build_log.c.log_content)
+                        .where(script_build_log.c.script_build_id == script_build_id)
+                        .order_by(script_build_log.c.id)
+                    )
+                ).scalars()
+            )
+        content = "\n".join(str(item) for item in log_rows if item is not None)
+        if not content:
+            content = str(row["summary"] or row["error_message"] or "")
+        return {
+            "success": True,
+            "log_content": content,
+            "live": row["status"] in {"running", "stopping"},
+            "status": row["status"],
+        }
+
+    async def root_trace_id(self, script_build_id: int, principal: Principal) -> str | None:
+        row = await self._row(script_build_id, principal)
+        value = row["reson_trace_id"]
+        return str(value) if value else None
+
+    async def _row(self, script_build_id: int, principal: Principal) -> Any:
+        await self._authorizer.require_access(principal, script_build_id)
+        async with self._read() as session:
+            row = (
+                (
+                    await session.execute(
+                        select(script_build_record).where(
+                            script_build_record.c.id == script_build_id,
+                            script_build_record.c.is_deleted.is_(False),
+                        )
+                    )
+                )
+                .mappings()
+                .one_or_none()
+            )
+        if row is None:
+            raise BuildNotFound()
+        return row
+
+
+def _list_item(row: Any) -> dict[str, Any]:
+    return {
+        "id": int(row["id"]),
+        "execution_id": int(row["execution_id"]),
+        "topic_build_id": int(row["topic_build_id"]),
+        "topic_id": int(row["topic_id"]),
+        "status": row["status"],
+        "summary": row["summary"],
+        "paragraph_count": row["paragraph_count"],
+        "element_count": row["element_count"],
+        "is_favorited": bool(row["is_favorited"]),
+        "start_time": row["start_time"].isoformat() if row["start_time"] else None,
+        "end_time": row["end_time"].isoformat() if row["end_time"] else None,
+    }
+
+
+__all__ = ["LegacyScriptBuildApiService"]

+ 224 - 0
script_build_host/src/script_build_host/application/recovery.py

@@ -0,0 +1,224 @@
+"""Explicit per-build phase-three finalization and recovery classification."""
+
+from __future__ import annotations
+
+from collections.abc import Callable
+from typing import Any
+
+from agent.orchestration import OperationStatus, TaskStatus
+
+from script_build_host.application.phase_three import PhaseThreeContinuationService
+from script_build_host.application.publisher import ScriptPublisher
+from script_build_host.domain.artifacts import ArtifactKind, ArtifactState
+from script_build_host.domain.errors import ProtocolViolation, PublicationStateInconsistent
+from script_build_host.domain.phase_three_artifacts import RootDeliveryManifestV1
+from script_build_host.domain.ports import (
+    BuildAuthorizer,
+    LegacyBuildStateRepository,
+    MissionBindingRepository,
+    PublicationRepository,
+    ScriptBusinessArtifactRepository,
+)
+from script_build_host.domain.records import (
+    BuildStatus,
+    Principal,
+    PublicationResult,
+    PublicationState,
+    PublicationType,
+    RecoveryClassification,
+    ResumeResult,
+)
+from script_build_host.infrastructure.ownership import OwnerLease
+
+
+class PhaseThreeFinalizationService:
+    def __init__(
+        self,
+        *,
+        coordinator: Any,
+        bindings: MissionBindingRepository,
+        authorizer: BuildAuthorizer,
+        publisher: ScriptPublisher,
+        owner_lease_factory: Callable[[int], OwnerLease],
+    ) -> None:
+        self.coordinator = coordinator
+        self.bindings = bindings
+        self.authorizer = authorizer
+        self.publisher = publisher
+        self.owner_lease_factory = owner_lease_factory
+
+    async def finalize(self, script_build_id: int, principal: Principal) -> PublicationResult:
+        await self.authorizer.require_access(principal, script_build_id)
+        binding = await self.bindings.get_by_build(script_build_id)
+        ledger = await self.coordinator.task_store.load(binding.root_trace_id)
+        root = ledger.tasks[ledger.root_task_id]
+        decisions = [
+            ledger.decisions[item]
+            for item in root.decision_ids
+            if ledger.decisions[item].action.value == "accept"
+        ]
+        if root.status is not TaskStatus.COMPLETED or len(decisions) != 1:
+            raise ProtocolViolation("finalize requires one completed Root ACCEPT")
+        decision = decisions[0]
+        if not decision.attempt_id:
+            raise ProtocolViolation("Root ACCEPT has no Attempt")
+        attempt = ledger.attempts[decision.attempt_id]
+        if attempt.submission is None:
+            raise ProtocolViolation("Root ACCEPT Attempt has no submission")
+        refs = [
+            item
+            for item in attempt.submission.artifact_refs
+            if item.kind == ArtifactKind.ROOT_DELIVERY_MANIFEST.value
+        ]
+        if len(refs) != 1:
+            raise ProtocolViolation("Root ACCEPT requires one delivery manifest")
+        lease = self.owner_lease_factory(script_build_id)
+        token = await lease.acquire(script_build_id)
+        try:
+            return await self.publisher.publish(
+                binding.root_trace_id,
+                decision.decision_id,
+                refs[0],
+                owner_token=token,
+            )
+        finally:
+            await lease.release()
+
+
+class PhaseThreeRecoveryService:
+    def __init__(
+        self,
+        *,
+        coordinator: Any,
+        bindings: MissionBindingRepository,
+        publications: PublicationRepository,
+        legacy_state: LegacyBuildStateRepository,
+        authorizer: BuildAuthorizer,
+        continuation: PhaseThreeContinuationService,
+        finalization: PhaseThreeFinalizationService,
+        artifacts: ScriptBusinessArtifactRepository | None = None,
+        legacy_projection: Any | None = None,
+    ) -> None:
+        self.coordinator = coordinator
+        self.bindings = bindings
+        self.publications = publications
+        self.legacy_state = legacy_state
+        self.authorizer = authorizer
+        self.continuation = continuation
+        self.finalization = finalization
+        self.artifacts = artifacts
+        self.legacy_projection = legacy_projection
+
+    async def resume(self, script_build_id: int, principal: Principal) -> ResumeResult:
+        await self.authorizer.require_access(principal, script_build_id)
+        binding = await self.bindings.get_by_build(script_build_id)
+        status = await self.legacy_state.get_status(script_build_id)
+        publication = await self.publications.get_by_build(
+            script_build_id, publication_type=PublicationType.FINAL
+        )
+        if status is BuildStatus.SUCCESS:
+            if publication is None or publication.state is not PublicationState.PUBLISHED:
+                raise PublicationStateInconsistent()
+            if binding.accepted_root_artifact_version_id != publication.artifact_version_id:
+                raise PublicationStateInconsistent()
+            if self.artifacts is not None:
+                manifest_version = await self.artifacts.get_by_id(
+                    publication.artifact_version_id, script_build_id=script_build_id
+                )
+                if (
+                    manifest_version.state is not ArtifactState.PUBLISHED
+                    or manifest_version.canonical_sha256 != publication.expected_sha256
+                    or not isinstance(manifest_version.artifact, RootDeliveryManifestV1)
+                ):
+                    raise PublicationStateInconsistent()
+                structured_id = int(
+                    manifest_version.artifact.structured_script_ref.rsplit("/", 1)[-1]
+                )
+                structured = await self.artifacts.get_by_id(
+                    structured_id, script_build_id=script_build_id
+                )
+                if structured.state is not ArtifactState.PUBLISHED:
+                    raise PublicationStateInconsistent()
+                if self.legacy_projection is not None:
+                    detail = await self.legacy_projection.project_authorized(
+                        script_build_id, principal=principal
+                    )
+                    canonical = self.legacy_projection.to_canonical(detail)
+                    if (
+                        canonical.canonical_sha256
+                        != manifest_version.artifact.legacy_projection_digest
+                    ):
+                        raise PublicationStateInconsistent()
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.ALREADY_COMPLETE,
+                BuildStatus.SUCCESS,
+            )
+        ledger = await self.coordinator.task_store.load(binding.root_trace_id)
+        root = ledger.tasks[ledger.root_task_id]
+        if root.status is TaskStatus.COMPLETED:
+            if publication is not None and publication.state is PublicationState.PUBLISHING:
+                raise PublicationStateInconsistent()
+            await self.finalization.finalize(script_build_id, principal)
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.FINALIZE_ONLY,
+                BuildStatus.SUCCESS,
+            )
+        active = [
+            item
+            for item in ledger.operations.values()
+            if item.status
+            in {OperationStatus.PENDING, OperationStatus.RUNNING, OperationStatus.STOP_REQUESTED}
+        ]
+        if any(item.status is OperationStatus.PENDING for item in active):
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.RESUME_PENDING_OPERATION,
+                status,
+                "pending durable Operation requires the framework recovery hook",
+            )
+        if active:
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.REPLAN_REQUIRED,
+                status,
+                "started execution must normalize before a new Attempt",
+            )
+        if root.status is TaskStatus.BLOCKED:
+            await self.continuation.advance(script_build_id, principal)
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.REPLAN_REQUIRED,
+                BuildStatus.RUNNING,
+                "phase-three continuation started",
+            )
+        return ResumeResult(
+            script_build_id,
+            RecoveryClassification.MANUAL_RECONCILIATION,
+            status,
+            "durable state is not safe for automatic replay",
+        )
+
+    async def reconcile(
+        self, script_build_id: int, principal: Principal, *, retry_finalize: bool = False
+    ) -> ResumeResult:
+        await self.authorizer.require_access(principal, script_build_id)
+        if retry_finalize:
+            await self.finalization.finalize(script_build_id, principal)
+            return ResumeResult(
+                script_build_id,
+                RecoveryClassification.FINALIZE_ONLY,
+                BuildStatus.SUCCESS,
+                "same immutable final identity reconciled",
+            )
+        status = await self.legacy_state.get_status(script_build_id)
+        return ResumeResult(
+            script_build_id,
+            RecoveryClassification.MANUAL_RECONCILIATION,
+            status,
+            "inspection only; no Artifact, digest, or Decision was changed",
+        )
+
+
+__all__ = ["PhaseThreeFinalizationService", "PhaseThreeRecoveryService"]

+ 85 - 0
script_build_host/src/script_build_host/composition.py

@@ -21,6 +21,7 @@ from script_build_host.agents import (
 from script_build_host.agents.model_resolver import SnapshotModelManifestSource
 from script_build_host.agents.prompt_catalog import (
     PHASE_TWO_PRESETS,
+    phase_three_prompt_requests,
     script_prompt_requests,
 )
 from script_build_host.api import ApiSecurity, UploadedTopicGateway, create_app
@@ -31,6 +32,7 @@ from script_build_host.application.mission_service import (
     DirectionReconciler,
     ScriptMissionService,
 )
+from script_build_host.application.phase_three import PhaseThreeContinuationService
 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 (
@@ -39,6 +41,12 @@ from script_build_host.application.phase_two_inputs import (
     StoredTaskContractReader,
 )
 from script_build_host.application.phase_two_planning import PhaseTwoPlanningService
+from script_build_host.application.publisher import ScriptPublisher
+from script_build_host.application.recovery import (
+    PhaseThreeFinalizationService,
+    PhaseThreeRecoveryService,
+)
+from script_build_host.application.root_delivery import RootDeliveryService
 from script_build_host.domain.ports import (
     BuildAuthorizer,
     InputSnapshotRepository,
@@ -89,6 +97,13 @@ class HostDependencies:
     max_image_bytes: int = 10 * 1024 * 1024
     max_total_image_bytes: int = 25 * 1024 * 1024
     phase_two_limits: PhaseTwoLimits = field(default_factory=PhaseTwoLimits)
+    owner_lease_factory: Any = None
+    fenced_command_gate: Any = None
+    final_publication_uow: Any = None
+    final_publications: PublicationRepository | None = None
+    legacy_projection: Any = None
+    legacy_api: Any = None
+    http_command_journal: Any = None
 
 
 @dataclass(frozen=True)
@@ -99,6 +114,9 @@ class HostComposition:
     tool_gateway: LegacyScriptToolGateway
     phase_two_planning: PhaseTwoPlanningService
     phase_two_candidates: PhaseTwoCandidateService
+    phase_three: PhaseThreeContinuationService | None = None
+    finalization: PhaseThreeFinalizationService | None = None
+    recovery: PhaseThreeRecoveryService | None = None
 
 
 def compose_host(dependencies: HostDependencies) -> HostComposition:
@@ -180,6 +198,11 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         max_total_image_bytes=dependencies.max_total_image_bytes,
     )
     planning_service.workspace_lifecycle = candidate_service
+    root_delivery = RootDeliveryService(
+        bindings=dependencies.bindings,
+        task_store=dependencies.task_store,
+        artifacts=dependencies.business_artifacts,
+    )
     transition_gate = BuildTransitionGate()
     direction_reconciler = DirectionReconciler(
         coordinator=coordinator,
@@ -207,6 +230,8 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         phase_two_required_presets=tuple(sorted(PHASE_TWO_PRESETS)),
         phase_two_model_manifest=dict(dependencies.default_model_manifest),
         phase_two_boundary_verifier=boundary_verifier,
+        owner_lease_factory=dependencies.owner_lease_factory,
+        fenced_command_gate=dependencies.fenced_command_gate,
     )
     planning_service.dispatch_guard = mission_service
     gateway = LegacyScriptToolGateway(
@@ -220,6 +245,7 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         planner_tools=planning_service,
         candidate_tools=candidate_service,
         budget_guard=planning_service,
+        root_tools=root_delivery,
     )
     register_script_tools(dependencies.runner.tools, gateway)
     security = ApiSecurity(
@@ -228,6 +254,56 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         hide_forbidden=dependencies.hide_forbidden,
         websocket_allowed_origins=dependencies.websocket_allowed_origins,
     )
+    phase_three: PhaseThreeContinuationService | None = None
+    finalization: PhaseThreeFinalizationService | None = None
+    recovery: PhaseThreeRecoveryService | None = None
+    if (
+        dependencies.owner_lease_factory is not None
+        and dependencies.final_publication_uow is not None
+    ):
+        final_publications = dependencies.final_publications or dependencies.publications
+        phase_three = PhaseThreeContinuationService(
+            runner=dependencies.runner,
+            coordinator=coordinator,
+            factory=ScriptMissionFactory(),
+            input_snapshots=dependencies.input_snapshot_service,
+            bindings=dependencies.bindings,
+            artifacts=dependencies.business_artifacts,
+            legacy_state=dependencies.legacy_state,
+            authorizer=dependencies.build_authorizer,
+            contracts=contract_store,
+            owner_lease_factory=dependencies.owner_lease_factory,
+            phase_two_boundary_verifier=boundary_verifier,
+            fenced_command_gate=dependencies.fenced_command_gate,
+            owner_registry=mission_service,
+            prompt_requests=phase_three_prompt_requests(),
+            model_manifest=dict(dependencies.default_model_manifest),
+        )
+        publisher = ScriptPublisher(
+            coordinator=coordinator,
+            bindings=dependencies.bindings,
+            artifacts=dependencies.business_artifacts,
+            publications=final_publications,
+            final_uow=dependencies.final_publication_uow,
+        )
+        finalization = PhaseThreeFinalizationService(
+            coordinator=coordinator,
+            bindings=dependencies.bindings,
+            authorizer=dependencies.build_authorizer,
+            publisher=publisher,
+            owner_lease_factory=dependencies.owner_lease_factory,
+        )
+        recovery = PhaseThreeRecoveryService(
+            coordinator=coordinator,
+            bindings=dependencies.bindings,
+            publications=final_publications,
+            legacy_state=dependencies.legacy_state,
+            authorizer=dependencies.build_authorizer,
+            continuation=phase_three,
+            finalization=finalization,
+            artifacts=dependencies.business_artifacts,
+            legacy_projection=dependencies.legacy_projection,
+        )
     app = create_app(
         mission_service=mission_service,
         security=security,
@@ -240,6 +316,12 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         outbound_policy=dependencies.outbound_policy,
         default_model_manifest=dependencies.default_model_manifest,
         default_datasource_manifest=dependencies.default_datasource_manifest,
+        phase_three=phase_three,
+        finalization=finalization,
+        recovery=recovery,
+        legacy_projection=dependencies.legacy_projection,
+        legacy_api=dependencies.legacy_api,
+        http_command_journal=dependencies.http_command_journal,
     )
     return HostComposition(
         app=app,
@@ -248,6 +330,9 @@ def compose_host(dependencies: HostDependencies) -> HostComposition:
         tool_gateway=gateway,
         phase_two_planning=planning_service,
         phase_two_candidates=candidate_service,
+        phase_three=phase_three,
+        finalization=finalization,
+        recovery=recovery,
     )
 
 

+ 246 - 0
script_build_host/src/script_build_host/infrastructure/http_commands.py

@@ -0,0 +1,246 @@
+"""SQL command journal for optional HTTP Idempotency-Key semantics."""
+
+from __future__ import annotations
+
+import asyncio
+from collections.abc import Awaitable, Callable, Mapping
+from datetime import UTC, datetime
+from typing import Any, TypeVar
+
+from sqlalchemy import insert, select, update
+from sqlalchemy.exc import IntegrityError
+from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
+
+from script_build_host.domain.canonical_json import canonical_json_bytes, canonical_sha256
+from script_build_host.domain.errors import HttpIdempotencyConflict, ProtocolViolation
+from script_build_host.domain.records import Principal
+from script_build_host.infrastructure.tables import http_command_table
+
+T = TypeVar("T")
+MAX_IDEMPOTENCY_RESPONSE_BYTES = 64 * 1024
+
+
+class HttpCommandJournal:
+    def __init__(self, sessions: async_sessionmaker[AsyncSession]) -> None:
+        self._sessions = sessions
+
+    async def execute(
+        self,
+        *,
+        principal: Principal,
+        route_family: str,
+        idempotency_key: str,
+        request_payload: Mapping[str, Any],
+        command: Callable[[], Awaitable[tuple[int, dict[str, Any]]]],
+        resume_reserved: Callable[[str | None, int | None], Awaitable[tuple[int, dict[str, Any]]]]
+        | None = None,
+        preallocated_root_trace_id: str | None = None,
+        resource_id: int | None = None,
+    ) -> tuple[int, dict[str, Any]]:
+        if not idempotency_key.strip() or len(idempotency_key) > 128:
+            raise ProtocolViolation("Idempotency-Key must contain between 1 and 128 characters")
+        if not route_family or len(route_family) > 64:
+            raise ProtocolViolation("HTTP command route family exceeds 64 characters")
+        scope = canonical_sha256(
+            {"subject": principal.subject, "tenant_id": principal.tenant_id}
+        ).hex_value
+        fingerprint = canonical_sha256(
+            {"route_family": route_family, "request": dict(request_payload)}
+        ).hex_value
+        existing = await self._reserve(
+            scope=scope,
+            route_family=route_family,
+            key=idempotency_key,
+            fingerprint=fingerprint,
+            preallocated_root_trace_id=preallocated_root_trace_id,
+            resource_id=resource_id,
+        )
+        if existing is not None:
+            if existing["request_fingerprint"] != fingerprint:
+                raise HttpIdempotencyConflict()
+            if existing["state"] == "completed":
+                response = existing["response_json"]
+                if not isinstance(response, dict) or existing["response_status"] is None:
+                    raise ProtocolViolation("completed HTTP command has no bounded response")
+                return int(existing["response_status"]), dict(response)
+            if existing["state"] == "reserved":
+                if resume_reserved is not None:
+                    status, response = await resume_reserved(
+                        existing["preallocated_root_trace_id"], existing["resource_id"]
+                    )
+                    await self._complete(
+                        scope=scope,
+                        route_family=route_family,
+                        key=idempotency_key,
+                        fingerprint=fingerprint,
+                        status=status,
+                        response=response,
+                        resource_id=(
+                            int(existing["resource_id"])
+                            if existing["resource_id"] is not None
+                            else resource_id
+                        ),
+                    )
+                    return status, response
+                return await self._await_completed(
+                    scope=scope,
+                    route_family=route_family,
+                    key=idempotency_key,
+                    fingerprint=fingerprint,
+                )
+            raise ProtocolViolation("the idempotent HTTP command requires reconciliation")
+        try:
+            status, response = await command()
+            if len(canonical_json_bytes(response)) > MAX_IDEMPOTENCY_RESPONSE_BYTES:
+                raise ProtocolViolation("idempotent HTTP response exceeds 64 KiB")
+        except BaseException:
+            async with self._sessions() as session, session.begin():
+                await session.execute(
+                    update(http_command_table)
+                    .where(
+                        http_command_table.c.principal_scope_sha256 == scope,
+                        http_command_table.c.route_family == route_family,
+                        http_command_table.c.idempotency_key == idempotency_key,
+                        http_command_table.c.request_fingerprint == fingerprint,
+                        http_command_table.c.state == "reserved",
+                    )
+                    .values(state="failed", updated_at=datetime.now(UTC))
+                )
+            raise
+        await self._complete(
+            scope=scope,
+            route_family=route_family,
+            key=idempotency_key,
+            fingerprint=fingerprint,
+            status=status,
+            response=response,
+            resource_id=resource_id,
+        )
+        return status, response
+
+    async def _complete(
+        self,
+        *,
+        scope: str,
+        route_family: str,
+        key: str,
+        fingerprint: str,
+        status: int,
+        response: dict[str, Any],
+        resource_id: int | None,
+    ) -> None:
+        if len(canonical_json_bytes(response)) > MAX_IDEMPOTENCY_RESPONSE_BYTES:
+            raise ProtocolViolation("idempotent HTTP response exceeds 64 KiB")
+        async with self._sessions() as session, session.begin():
+            result = await session.execute(
+                update(http_command_table)
+                .where(
+                    http_command_table.c.principal_scope_sha256 == scope,
+                    http_command_table.c.route_family == route_family,
+                    http_command_table.c.idempotency_key == key,
+                    http_command_table.c.request_fingerprint == fingerprint,
+                    http_command_table.c.state == "reserved",
+                )
+                .values(
+                    state="completed",
+                    response_status=status,
+                    response_json=response,
+                    resource_id=resource_id,
+                    updated_at=datetime.now(UTC),
+                )
+            )
+            if result.rowcount != 1:
+                raise ProtocolViolation("idempotent HTTP command completion lost its reservation")
+
+    async def _await_completed(
+        self, *, scope: str, route_family: str, key: str, fingerprint: str
+    ) -> tuple[int, dict[str, Any]]:
+        for _ in range(500):
+            await asyncio.sleep(0.01)
+            async with self._sessions() as session:
+                row = (
+                    (
+                        await session.execute(
+                            select(http_command_table).where(
+                                http_command_table.c.principal_scope_sha256 == scope,
+                                http_command_table.c.route_family == route_family,
+                                http_command_table.c.idempotency_key == key,
+                            )
+                        )
+                    )
+                    .mappings()
+                    .one()
+                )
+            if row["request_fingerprint"] != fingerprint:
+                raise HttpIdempotencyConflict()
+            if row["state"] == "failed":
+                raise ProtocolViolation("the idempotent HTTP command requires reconciliation")
+            if row["state"] == "completed":
+                response = row["response_json"]
+                if not isinstance(response, dict) or row["response_status"] is None:
+                    raise ProtocolViolation("completed HTTP command has no bounded response")
+                return int(row["response_status"]), dict(response)
+        raise ProtocolViolation("the idempotent HTTP command is still reserved")
+
+    async def _reserve(
+        self,
+        *,
+        scope: str,
+        route_family: str,
+        key: str,
+        fingerprint: str,
+        preallocated_root_trace_id: str | None,
+        resource_id: int | None,
+    ) -> Any | None:
+        now = datetime.now(UTC)
+        try:
+            async with self._sessions() as session, session.begin():
+                existing = (
+                    (
+                        await session.execute(
+                            select(http_command_table)
+                            .where(
+                                http_command_table.c.principal_scope_sha256 == scope,
+                                http_command_table.c.route_family == route_family,
+                                http_command_table.c.idempotency_key == key,
+                            )
+                            .with_for_update()
+                        )
+                    )
+                    .mappings()
+                    .one_or_none()
+                )
+                if existing is not None:
+                    return existing
+                await session.execute(
+                    insert(http_command_table).values(
+                        principal_scope_sha256=scope,
+                        route_family=route_family,
+                        idempotency_key=key,
+                        request_fingerprint=fingerprint,
+                        preallocated_root_trace_id=preallocated_root_trace_id,
+                        resource_id=resource_id,
+                        state="reserved",
+                        created_at=now,
+                        updated_at=now,
+                    )
+                )
+                return None
+        except IntegrityError:
+            async with self._sessions() as session:
+                return (
+                    (
+                        await session.execute(
+                            select(http_command_table).where(
+                                http_command_table.c.principal_scope_sha256 == scope,
+                                http_command_table.c.route_family == route_family,
+                                http_command_table.c.idempotency_key == key,
+                            )
+                        )
+                    )
+                    .mappings()
+                    .one()
+                )
+
+
+__all__ = ["MAX_IDEMPOTENCY_RESPONSE_BYTES", "HttpCommandJournal"]

+ 37 - 1
script_build_host/src/script_build_host/production.py

@@ -23,6 +23,8 @@ from script_build_host.adapters import (
     SqlUploadedTopicGateway,
 )
 from script_build_host.application import ScriptInputSnapshotService
+from script_build_host.application.legacy_api import LegacyScriptBuildApiService
+from script_build_host.application.legacy_projection import LegacyDetailProjectionService
 from script_build_host.composition import HostComposition, HostDependencies, compose_host
 from script_build_host.domain.ports import (
     BuildAuthorizer,
@@ -31,8 +33,13 @@ from script_build_host.domain.ports import (
 )
 from script_build_host.infrastructure.config import ScriptBuildSettings
 from script_build_host.infrastructure.db import DatabaseSessions, create_database_sessions
+from script_build_host.infrastructure.final_publication import (
+    SqlAlchemyFinalPublicationUnitOfWork,
+)
+from script_build_host.infrastructure.http_commands import HttpCommandJournal
 from script_build_host.infrastructure.manifests import SettingsRuntimeManifestProvider
 from script_build_host.infrastructure.outbound import OutboundPolicy
+from script_build_host.infrastructure.ownership import FencedCommandGate, OwnerLease
 from script_build_host.infrastructure.raw_artifacts import FileRawArtifactStore
 from script_build_host.infrastructure.trace_verifier import FileSystemTraceStoreVerifier
 from script_build_host.repositories import (
@@ -156,9 +163,11 @@ def compose_production_host(
     bindings = SqlAlchemyMissionBindingRepository(database.write)
     artifacts = SqlAlchemyScriptBusinessArtifactRepository(database.write)
     publications = SqlAlchemyPublicationRepository(database.write)
+    final_publications = SqlAlchemyPublicationRepository(database.final)
     legacy_state = SqlAlchemyLegacyBuildStateRepository(database.write)
+    legacy_input = LegacySqlAlchemyInputReader(database.read)
     input_service = ScriptInputSnapshotService(
-        legacy_input=LegacySqlAlchemyInputReader(database.read),
+        legacy_input=legacy_input,
         persona_source=FilePersonaSource(
             persona_root=settings.persona_root,
             section_pattern_root=settings.section_pattern_root,
@@ -180,6 +189,24 @@ def compose_production_host(
     data_root = settings.agent_data_root.resolve()
     raw_artifacts = FileRawArtifactStore(data_root)
     candidate_workspaces = SqlAlchemyCandidateWorkspaceRepository(database.write, artifacts)
+    legacy_projection = LegacyDetailProjectionService(
+        database.read,
+        runtime.build_authorizer,
+        topic_reader=legacy_input,
+        bindings=bindings,
+        snapshots=snapshots,
+    )
+    final_uow = SqlAlchemyFinalPublicationUnitOfWork(
+        database.final,
+        legacy_projection,
+        FencedCommandGate(database.final),
+    )
+    runtime_fencing = FencedCommandGate(database.write)
+    lease_root = data_root / "mission-owner-leases"
+
+    def owner_lease_factory(_script_build_id: int) -> OwnerLease:
+        return OwnerLease(database.write, lease_root)
+
     composition = compose_host(
         HostDependencies(
             runner=runtime.runner,
@@ -224,6 +251,15 @@ def compose_production_host(
             max_image_bytes=settings.max_image_bytes,
             max_total_image_bytes=settings.max_total_image_bytes,
             phase_two_limits=settings.phase_two_limits(),
+            owner_lease_factory=owner_lease_factory,
+            fenced_command_gate=runtime_fencing,
+            final_publication_uow=final_uow,
+            final_publications=final_publications,
+            legacy_projection=legacy_projection,
+            legacy_api=LegacyScriptBuildApiService(
+                database.read, database.write, runtime.build_authorizer
+            ),
+            http_command_journal=HttpCommandJournal(database.write),
         )
     )
     host = ProductionHost(