Explorar o código

feat(script-build): publish legacy projection atomically

SamLee hai 1 día
pai
achega
e0ed388e63

+ 380 - 0
script_build_host/src/script_build_host/application/legacy_projection.py

@@ -0,0 +1,380 @@
+"""One legacy detail projection shared by HTTP reads and final publication."""
+
+from __future__ import annotations
+
+from collections import defaultdict
+from collections.abc import Mapping
+from datetime import date, datetime
+from typing import Any, cast
+
+from sqlalchemy import or_, select
+from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
+
+from script_build_host.domain.errors import BuildNotFound, LegacyCanonicalIdentityConflict
+from script_build_host.domain.phase_three_artifacts import (
+    LegacyProjectionCanonicalV1,
+    canonicalize_legacy_projection,
+)
+from script_build_host.domain.phase_two_artifacts import (
+    ScriptElementV1,
+    ScriptParagraphElementLinkV1,
+    ScriptParagraphV1,
+    StructuredScriptArtifactV1,
+)
+from script_build_host.domain.ports import (
+    BuildAuthorizer,
+    InputSnapshotRepository,
+    MissionBindingRepository,
+)
+from script_build_host.domain.records import Principal
+from script_build_host.infrastructure.legacy_tables import (
+    script_build_decision_trace,
+    script_build_element,
+    script_build_paragraph,
+    script_build_paragraph_element,
+    script_build_record,
+)
+
+
+class LegacyDetailProjectionService:
+    def __init__(
+        self,
+        sessions: async_sessionmaker[AsyncSession],
+        authorizer: BuildAuthorizer,
+        *,
+        topic_reader: Any | None = None,
+        bindings: MissionBindingRepository | None = None,
+        snapshots: InputSnapshotRepository | None = None,
+    ) -> None:
+        self._sessions = sessions
+        self._authorizer = authorizer
+        self._topic_reader = topic_reader
+        self._bindings = bindings
+        self._snapshots = snapshots
+
+    async def project_authorized(
+        self, script_build_id: int, *, principal: Principal
+    ) -> dict[str, Any]:
+        await self._authorizer.require_access(principal, script_build_id)
+        async with self._sessions() as session:
+            detail = await self.project_in_session(script_build_id, session=session)
+        if self._bindings is not None and self._snapshots is not None:
+            try:
+                binding = await self._bindings.get_by_build(script_build_id)
+            except BuildNotFound:
+                binding = None
+            if binding is not None:
+                snapshot = await self._snapshots.get(
+                    str(binding.input_snapshot_id), script_build_id=script_build_id
+                )
+                detail["topic"] = snapshot.topic
+                return detail
+        if self._topic_reader is not None:
+            build = detail["build"]
+            detail["topic"] = await self._topic_reader.read_topic_graph(
+                execution_id=int(build["execution_id"]),
+                topic_build_id=int(build["topic_build_id"]),
+                topic_id=int(build["topic_id"]),
+            )
+        return detail
+
+    async def project_in_session(
+        self, script_build_id: int, *, session: AsyncSession
+    ) -> dict[str, Any]:
+        build = (
+            (
+                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 build is None:
+            raise BuildNotFound()
+        paragraphs: list[Any] = list(
+            (
+                await session.execute(
+                    select(script_build_paragraph)
+                    .where(
+                        script_build_paragraph.c.script_build_id == script_build_id,
+                        script_build_paragraph.c.branch_id == 0,
+                        or_(
+                            script_build_paragraph.c.is_active.is_(True),
+                            script_build_paragraph.c.is_active.is_(None),
+                        ),
+                    )
+                    .order_by(
+                        script_build_paragraph.c.level,
+                        script_build_paragraph.c.paragraph_index,
+                        script_build_paragraph.c.id,
+                    )
+                )
+            )
+            .mappings()
+            .all()
+        )
+        elements: list[Any] = list(
+            (
+                await session.execute(
+                    select(script_build_element)
+                    .where(
+                        script_build_element.c.script_build_id == script_build_id,
+                        script_build_element.c.branch_id == 0,
+                        or_(
+                            script_build_element.c.is_active.is_(True),
+                            script_build_element.c.is_active.is_(None),
+                        ),
+                    )
+                    .order_by(script_build_element.c.id)
+                )
+            )
+            .mappings()
+            .all()
+        )
+        links: list[Any] = list(
+            (
+                await session.execute(
+                    select(script_build_paragraph_element)
+                    .where(
+                        script_build_paragraph_element.c.script_build_id == script_build_id,
+                        script_build_paragraph_element.c.branch_id == 0,
+                    )
+                    .order_by(script_build_paragraph_element.c.id)
+                )
+            )
+            .mappings()
+            .all()
+        )
+        decision_traces: list[Any] = list(
+            (
+                await session.execute(
+                    select(script_build_decision_trace)
+                    .where(script_build_decision_trace.c.script_build_id == script_build_id)
+                    .order_by(script_build_decision_trace.c.id)
+                )
+            )
+            .mappings()
+            .all()
+        )
+        paragraph_ids = {int(item["id"]) for item in paragraphs}
+        element_ids = {int(item["id"]) for item in elements}
+        owned_links = [
+            item
+            for item in links
+            if int(item["paragraph_id"]) in paragraph_ids and int(item["element_id"]) in element_ids
+        ]
+        links_by_paragraph: dict[int, list[Mapping[str, Any]]] = defaultdict(list)
+        for item in owned_links:
+            links_by_paragraph[int(item["paragraph_id"])].append(item)
+        children: dict[int, list[dict[str, Any]]] = defaultdict(list)
+        projected: dict[int, dict[str, Any]] = {}
+        for item in paragraphs:
+            identifier = int(item["id"])
+            paragraph = _legacy_paragraph(item)
+            paragraph_links = links_by_paragraph[identifier]
+            paragraph["linked_element_ids"] = [int(link["element_id"]) for link in paragraph_links]
+            paragraph["linked_elements"] = [
+                {"element_id": int(link["element_id"]), "link_id": int(link["id"])}
+                for link in paragraph_links
+            ]
+            paragraph["sub_paragraphs"] = []
+            projected[identifier] = paragraph
+        roots: list[dict[str, Any]] = []
+        for item in paragraphs:
+            identifier = int(item["id"])
+            parent = item["parent_id"]
+            if parent is None:
+                roots.append(projected[identifier])
+            elif int(parent) in projected:
+                children[int(parent)].append(projected[identifier])
+            else:
+                raise LegacyCanonicalIdentityConflict("legacy paragraph parent is dangling")
+        _reject_physical_parent_cycles(paragraphs)
+        for parent_id, values in children.items():
+            values.sort(key=lambda item: (item["paragraph_index"], item["id"]))
+            projected[parent_id]["sub_paragraphs"] = values
+
+        return {
+            "success": True,
+            "topic": None,
+            "build": _legacy_build(build),
+            "paragraphs": roots,
+            "elements": [_legacy_element(item) for item in elements],
+            "decision_traces": [
+                {
+                    "id": int(item["id"]),
+                    "trace_type": item["trace_type"],
+                    "content": item["content"],
+                    "summary": item["summary"],
+                    "created_at": _json_value(item["created_at"]),
+                }
+                for item in decision_traces
+            ],
+            "post_datas": {},
+            "script_datas": {},
+        }
+
+    def to_canonical(self, detail: Mapping[str, Any]) -> LegacyProjectionCanonicalV1:
+        build = cast(Mapping[str, Any], detail.get("build") or {})
+        paragraphs = _flatten_paragraphs(detail.get("paragraphs"))
+        elements = detail.get("elements") or []
+        if not isinstance(elements, list):
+            raise LegacyCanonicalIdentityConflict("legacy elements must be an array")
+        candidate_paragraphs = tuple(ScriptParagraphV1.from_payload(item) for item in paragraphs)
+        candidate_elements = tuple(
+            ScriptElementV1.from_payload({**item, "element_id": int(item["id"])})
+            for item in elements
+        )
+        links: list[ScriptParagraphElementLinkV1] = []
+        for item in paragraphs:
+            linked = item.get("linked_elements") or []
+            if not isinstance(linked, list):
+                raise LegacyCanonicalIdentityConflict("linked_elements must be an array")
+            links.extend(
+                ScriptParagraphElementLinkV1(
+                    paragraph_id=int(item["paragraph_id"]), element_id=int(link["element_id"])
+                )
+                for link in linked
+            )
+        script = StructuredScriptArtifactV1(
+            direction_ref="script-build://artifact-versions/1",
+            input_closure_digest="sha256:" + "0" * 64,
+            paragraphs=candidate_paragraphs,
+            elements=candidate_elements,
+            paragraph_element_links=tuple(links),
+            source_artifact_refs=(),
+            evidence_refs=(),
+            acceptance_notes=(),
+        )
+        return canonicalize_legacy_projection(
+            script,
+            direction=_optional_text(build.get("script_direction")),
+            summary=_optional_text(build.get("summary")),
+        )
+
+
+def _legacy_build(row: Any) -> dict[str, Any]:
+    fields = (
+        "id",
+        "execution_id",
+        "topic_build_id",
+        "topic_id",
+        "agent_type",
+        "agent_config",
+        "data_source_url",
+        "status",
+        "summary",
+        "paragraph_count",
+        "element_count",
+        "input_tokens",
+        "output_tokens",
+        "cost_usd",
+        "error_message",
+        "strategies_config",
+        "script_direction",
+        "build_workflow",
+        "start_time",
+        "end_time",
+    )
+    output = {name: _json_value(row[name]) for name in fields}
+    config = output.get("agent_config")
+    output["model_name"] = config.get("model_name") if isinstance(config, dict) else None
+    return output
+
+
+def _legacy_paragraph(row: Any) -> dict[str, Any]:
+    fields = (
+        "id",
+        "paragraph_index",
+        "level",
+        "parent_id",
+        "name",
+        "content_range",
+        "theme",
+        "form",
+        "function",
+        "feeling",
+        "theme_elements",
+        "form_elements",
+        "function_elements",
+        "feeling_elements",
+        "description",
+        "full_description",
+        "is_active",
+        "created_at",
+        "updated_at",
+    )
+    return {name: _json_value(row[name]) for name in fields}
+
+
+def _legacy_element(row: Any) -> dict[str, Any]:
+    fields = (
+        "id",
+        "name",
+        "dimension_primary",
+        "dimension_secondary",
+        "commonality_analysis",
+        "topic_support",
+        "weight_score",
+        "support_elements",
+        "is_active",
+        "created_at",
+        "updated_at",
+    )
+    return {name: _json_value(row[name]) for name in fields}
+
+
+def _flatten_paragraphs(value: Any) -> list[dict[str, Any]]:
+    if not isinstance(value, list):
+        raise LegacyCanonicalIdentityConflict("legacy paragraphs must be an array")
+    output: list[dict[str, Any]] = []
+
+    def visit(item: Any, parent_id: int | None) -> None:
+        if not isinstance(item, Mapping):
+            raise LegacyCanonicalIdentityConflict("legacy paragraph must be an object")
+        current = dict(item)
+        children = current.pop("sub_paragraphs", []) or []
+        current["paragraph_id"] = int(current.pop("id"))
+        current["parent_id"] = parent_id
+        current.pop("linked_element_ids", None)
+        output.append(current)
+        if not isinstance(children, list):
+            raise LegacyCanonicalIdentityConflict("sub_paragraphs must be an array")
+        for child in children:
+            visit(child, int(current["paragraph_id"]))
+
+    for root in value:
+        visit(root, None)
+    return output
+
+
+def _reject_physical_parent_cycles(paragraphs: list[Any]) -> None:
+    parents = {int(item["id"]): item["parent_id"] for item in paragraphs}
+    for start in parents:
+        seen: set[int] = set()
+        current: int | None = start
+        while current is not None:
+            if current in seen:
+                raise LegacyCanonicalIdentityConflict(
+                    "legacy paragraph parent graph contains a cycle"
+                )
+            seen.add(current)
+            raw = parents.get(current)
+            current = int(raw) if raw is not None else None
+
+
+def _json_value(value: Any) -> Any:
+    if isinstance(value, (datetime, date)):
+        return value.isoformat()
+    return value
+
+
+def _optional_text(value: Any) -> str | None:
+    return str(value) if value is not None else None
+
+
+__all__ = ["LegacyDetailProjectionService"]

+ 292 - 0
script_build_host/src/script_build_host/application/publisher.py

@@ -0,0 +1,292 @@
+"""Root ACCEPT closure verification and deterministic final publication."""
+
+from __future__ import annotations
+
+from dataclasses import asdict
+from typing import Protocol
+
+from agent.orchestration import (
+    ArtifactRef,
+    DecisionAction,
+    TaskStatus,
+    ValidationVerdict,
+)
+
+from script_build_host.domain.artifacts import (
+    ArtifactKind,
+    ArtifactState,
+    ArtifactVersion,
+    ScriptDirectionArtifactV1,
+)
+from script_build_host.domain.errors import (
+    ProtocolViolation,
+    PublicationStateInconsistent,
+)
+from script_build_host.domain.phase_three_artifacts import (
+    RootDeliveryManifestV1,
+    canonicalize_legacy_projection,
+    root_delivery_input_closure_digest,
+)
+from script_build_host.domain.phase_two_artifacts import (
+    CandidatePortfolioArtifactV1,
+    StructuredScriptArtifactV1,
+)
+from script_build_host.domain.ports import (
+    MissionBindingRepository,
+    PublicationRepository,
+    ScriptBusinessArtifactRepository,
+)
+from script_build_host.domain.records import (
+    MissionOwnerToken,
+    Publication,
+    PublicationResult,
+    PublicationState,
+    PublicationType,
+)
+
+
+class FinalPublicationPort(Protocol):
+    async def reserve(
+        self,
+        *,
+        script_build_id: int,
+        accept_decision_id: str,
+        artifact_version_id: int,
+        expected_sha256: str,
+        owner_token: MissionOwnerToken,
+    ) -> Publication: ...
+
+    async def mark_failed(
+        self,
+        publication_id: int,
+        *,
+        error_code: str,
+        error_summary: str,
+        owner_token: MissionOwnerToken,
+    ) -> None: ...
+
+    async def publish(
+        self,
+        *,
+        script_build_id: int,
+        publication_id: int,
+        manifest_version: ArtifactVersion,
+        structured_version: ArtifactVersion,
+        manifest: RootDeliveryManifestV1,
+        structured_script: StructuredScriptArtifactV1,
+        direction_markdown: str,
+        owner_token: MissionOwnerToken,
+    ) -> PublicationResult: ...
+
+
+class ScriptPublisher:
+    def __init__(
+        self,
+        *,
+        coordinator: object,
+        bindings: MissionBindingRepository,
+        artifacts: ScriptBusinessArtifactRepository,
+        publications: PublicationRepository,
+        final_uow: FinalPublicationPort,
+    ) -> None:
+        self.coordinator = coordinator
+        self.bindings = bindings
+        self.artifacts = artifacts
+        self.publications = publications
+        self.final_uow = final_uow
+
+    async def publish(
+        self,
+        root_trace_id: str,
+        accept_decision_id: str,
+        root_delivery_manifest_ref: ArtifactRef,
+        *,
+        owner_token: MissionOwnerToken,
+    ) -> PublicationResult:
+        closure = await self._load_closure(
+            root_trace_id, accept_decision_id, root_delivery_manifest_ref, owner_token
+        )
+        manifest_version, manifest, structured_version, structured, direction = closure
+        publication = await self.final_uow.reserve(
+            script_build_id=owner_token.script_build_id,
+            accept_decision_id=accept_decision_id,
+            artifact_version_id=manifest_version.artifact_version_id,
+            expected_sha256=manifest_version.canonical_sha256,
+            owner_token=owner_token,
+        )
+        if publication.state is PublicationState.PUBLISHED:
+            return await self.final_uow.publish(
+                script_build_id=owner_token.script_build_id,
+                publication_id=publication.publication_id,
+                manifest_version=manifest_version,
+                structured_version=structured_version,
+                manifest=manifest,
+                structured_script=structured,
+                direction_markdown=direction.legacy_markdown,
+                owner_token=owner_token,
+            )
+        try:
+            return await self.final_uow.publish(
+                script_build_id=owner_token.script_build_id,
+                publication_id=publication.publication_id,
+                manifest_version=manifest_version,
+                structured_version=structured_version,
+                manifest=manifest,
+                structured_script=structured,
+                direction_markdown=direction.legacy_markdown,
+                owner_token=owner_token,
+            )
+        except Exception as exc:
+            current = await self.publications.get_by_build(
+                owner_token.script_build_id, publication_type=PublicationType.FINAL
+            )
+            if current is not None and current.state is PublicationState.PUBLISHED:
+                # Lost commit ACK: use the immutable final identity, never replay writes.
+                return await self.final_uow.publish(
+                    script_build_id=owner_token.script_build_id,
+                    publication_id=current.publication_id,
+                    manifest_version=manifest_version,
+                    structured_version=structured_version,
+                    manifest=manifest,
+                    structured_script=structured,
+                    direction_markdown=direction.legacy_markdown,
+                    owner_token=owner_token,
+                )
+            if (
+                current is None
+                or current.artifact_version_id != manifest_version.artifact_version_id
+            ):
+                raise PublicationStateInconsistent() from exc
+            await self.final_uow.mark_failed(
+                current.publication_id,
+                error_code=getattr(exc, "code", "FINAL_PUBLICATION_FAILED"),
+                error_summary=type(exc).__name__,
+                owner_token=owner_token,
+            )
+            raise
+
+    async def _load_closure(
+        self,
+        root_trace_id: str,
+        accept_decision_id: str,
+        manifest_ref: ArtifactRef,
+        owner_token: MissionOwnerToken,
+    ) -> tuple[
+        ArtifactVersion,
+        RootDeliveryManifestV1,
+        ArtifactVersion,
+        StructuredScriptArtifactV1,
+        ScriptDirectionArtifactV1,
+    ]:
+        if (
+            owner_token.root_trace_id != root_trace_id
+            or manifest_ref.kind != ArtifactKind.ROOT_DELIVERY_MANIFEST.value
+        ):
+            raise ProtocolViolation("final publication owner or manifest reference differs")
+        binding = await self.bindings.get_by_root(root_trace_id)
+        if binding.script_build_id != owner_token.script_build_id:
+            raise ProtocolViolation("final publication binding differs from its owner token")
+        ledger = await self.coordinator.task_store.load(root_trace_id)  # type: ignore[attr-defined]
+        root = ledger.tasks[ledger.root_task_id]
+        if root.status is not TaskStatus.COMPLETED:
+            raise ProtocolViolation("Root must be completed by an explicit ACCEPT")
+        decision = ledger.decisions.get(accept_decision_id)
+        if (
+            decision is None
+            or decision.action is not DecisionAction.ACCEPT
+            or decision.task_id != root.task_id
+            or not decision.attempt_id
+            or not decision.validation_id
+        ):
+            raise ProtocolViolation("final publication requires the Root ACCEPT decision")
+        attempt = ledger.attempts.get(decision.attempt_id)
+        validation = ledger.validations.get(decision.validation_id)
+        if (
+            attempt is None
+            or validation is None
+            or attempt.submission is None
+            or validation.attempt_id != attempt.attempt_id
+            or validation.snapshot_id != attempt.snapshot_id
+            or validation.verdict is not ValidationVerdict.PASSED
+        ):
+            raise ProtocolViolation("Root ACCEPT is not closed over one passed snapshot")
+        submitted = [
+            item
+            for item in attempt.submission.artifact_refs
+            if item.kind == ArtifactKind.ROOT_DELIVERY_MANIFEST.value
+        ]
+        if len(submitted) != 1 or asdict(submitted[0]) != asdict(manifest_ref):
+            raise ProtocolViolation("Root Attempt does not submit the requested manifest")
+        manifest_version = await self.artifacts.read_by_ref(
+            manifest_ref,
+            script_build_id=binding.script_build_id,
+            task_id=root.task_id,
+            attempt_id=attempt.attempt_id,
+        )
+        if manifest_version.state not in {
+            ArtifactState.FROZEN,
+            ArtifactState.PUBLISHED,
+        } or not isinstance(manifest_version.artifact, RootDeliveryManifestV1):
+            raise ProtocolViolation("Root delivery manifest is not frozen")
+        manifest = manifest_version.artifact
+        manifest.require_publishable()
+        direction_version = await self._version_from_uri(
+            manifest.direction_ref, binding.script_build_id
+        )
+        portfolio_version = await self._version_from_uri(
+            manifest.candidate_portfolio_ref, binding.script_build_id
+        )
+        structured_version = await self._version_from_uri(
+            manifest.structured_script_ref, binding.script_build_id
+        )
+        if (
+            not isinstance(direction_version.artifact, ScriptDirectionArtifactV1)
+            or not isinstance(portfolio_version.artifact, CandidatePortfolioArtifactV1)
+            or not isinstance(structured_version.artifact, StructuredScriptArtifactV1)
+        ):
+            raise ProtocolViolation("Root delivery manifest references the wrong artifact kinds")
+        portfolio = portfolio_version.artifact
+        structured = structured_version.artifact
+        if portfolio.adopted_structured_script_ref != manifest.structured_script_ref:
+            raise ProtocolViolation("Root delivery differs from the adopted StructuredScript")
+        expected_closure = root_delivery_input_closure_digest(
+            direction_ref=manifest.direction_ref,
+            direction_digest=direction_version.canonical_sha256,
+            candidate_portfolio_ref=manifest.candidate_portfolio_ref,
+            candidate_portfolio_digest=portfolio_version.canonical_sha256,
+            structured_script_ref=manifest.structured_script_ref,
+            structured_script_digest=structured_version.canonical_sha256,
+        )
+        if (
+            manifest.input_closure_digest != expected_closure
+            or structured.direction_ref != manifest.direction_ref
+        ):
+            raise ProtocolViolation("Root delivery immutable input closure differs")
+        canonical = canonicalize_legacy_projection(
+            structured,
+            direction=direction_version.artifact.legacy_markdown,
+            summary=manifest.build_summary,
+        )
+        if canonical.canonical_sha256 != manifest.legacy_projection_digest:
+            raise ProtocolViolation(
+                "Root delivery dry-run digest differs from its StructuredScript"
+            )
+        return (
+            manifest_version,
+            manifest,
+            structured_version,
+            structured,
+            direction_version.artifact,
+        )
+
+    async def _version_from_uri(self, uri: str, script_build_id: int) -> ArtifactVersion:
+        value = uri.removeprefix("script-build://artifact-versions/")
+        if not value.isdigit():
+            raise ProtocolViolation("delivery artifact URI is invalid")
+        version = await self.artifacts.get_by_id(int(value), script_build_id=script_build_id)
+        if version.state not in {ArtifactState.FROZEN, ArtifactState.PUBLISHED}:
+            raise ProtocolViolation("delivery artifact is not frozen")
+        return version
+
+
+__all__ = ["ScriptPublisher"]

+ 478 - 0
script_build_host/src/script_build_host/infrastructure/final_publication.py

@@ -0,0 +1,478 @@
+"""SQLAlchemy final-publication unit of work.
+
+All branch-zero, readback, pointer, artifact, publication and build mutations use
+one AsyncSession and one commit.
+"""
+
+from __future__ import annotations
+
+from datetime import UTC, datetime
+from typing import Any, cast
+
+from sqlalchemy import CursorResult, delete, insert, select, update
+from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
+
+from script_build_host.application.legacy_projection import LegacyDetailProjectionService
+from script_build_host.domain.artifacts import ArtifactState, ArtifactVersion
+from script_build_host.domain.errors import (
+    BuildNotFound,
+    MissionFencingTokenStale,
+    ProtocolViolation,
+    PublicationReadbackMismatch,
+)
+from script_build_host.domain.phase_three_artifacts import RootDeliveryManifestV1
+from script_build_host.domain.phase_two_artifacts import StructuredScriptArtifactV1
+from script_build_host.domain.records import (
+    MissionOwnerToken,
+    Publication,
+    PublicationResult,
+    PublicationState,
+    PublicationType,
+)
+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.ownership import FencedCommandGate
+from script_build_host.infrastructure.tables import (
+    artifact_version_table,
+    mission_binding_table,
+    publication_table,
+)
+
+
+class SqlAlchemyFinalPublicationUnitOfWork:
+    def __init__(
+        self,
+        sessions: async_sessionmaker[AsyncSession],
+        projection: LegacyDetailProjectionService,
+        fencing: FencedCommandGate,
+    ) -> None:
+        self._sessions = sessions
+        self._projection = projection
+        self._fencing = fencing
+
+    async def reserve(
+        self,
+        *,
+        script_build_id: int,
+        accept_decision_id: str,
+        artifact_version_id: int,
+        expected_sha256: str,
+        owner_token: MissionOwnerToken,
+    ) -> Publication:
+        now = datetime.now(UTC)
+        expected = expected_sha256.removeprefix("sha256:")
+        async with self._sessions() as session, session.begin():
+            await self._fencing.verify_in_session(session, owner_token)
+            row = (
+                (
+                    await session.execute(
+                        select(publication_table)
+                        .where(
+                            publication_table.c.script_build_id == script_build_id,
+                            publication_table.c.publication_type == PublicationType.FINAL.value,
+                        )
+                        .with_for_update()
+                    )
+                )
+                .mappings()
+                .one_or_none()
+            )
+            if row is None:
+                result = await session.execute(
+                    insert(publication_table).values(
+                        script_build_id=script_build_id,
+                        publication_type=PublicationType.FINAL.value,
+                        accept_decision_id=accept_decision_id,
+                        artifact_version_id=artifact_version_id,
+                        expected_sha256=expected,
+                        state=PublicationState.PENDING.value,
+                        attempt_count=0,
+                        publication_revision=0,
+                        created_at=now,
+                        updated_at=now,
+                    )
+                )
+                publication_id = int(cast(Any, result).inserted_primary_key[0])
+                row = (
+                    (
+                        await session.execute(
+                            select(publication_table).where(
+                                publication_table.c.id == publication_id
+                            )
+                        )
+                    )
+                    .mappings()
+                    .one()
+                )
+            if (
+                int(row["script_build_id"]) != script_build_id
+                or row["accept_decision_id"] != accept_decision_id
+                or int(row["artifact_version_id"]) != artifact_version_id
+                or row["expected_sha256"] != expected
+            ):
+                raise ProtocolViolation("final publication identity already differs")
+            return _publication(row)
+
+    async def mark_failed(
+        self,
+        publication_id: int,
+        *,
+        error_code: str,
+        error_summary: str,
+        owner_token: MissionOwnerToken,
+    ) -> None:
+        async with self._sessions() as session, session.begin():
+            await self._fencing.verify_in_session(session, owner_token)
+            result = await session.execute(
+                update(publication_table)
+                .where(
+                    publication_table.c.id == publication_id,
+                    publication_table.c.script_build_id == owner_token.script_build_id,
+                    publication_table.c.state != PublicationState.PUBLISHED.value,
+                )
+                .values(
+                    state=PublicationState.FAILED.value,
+                    last_error_code=error_code[:64],
+                    last_error_summary=error_summary[:1000],
+                    updated_at=datetime.now(UTC),
+                )
+            )
+            if result.rowcount != 1:
+                raise ProtocolViolation("final publication failure CAS did not match")
+
+    async def publish(
+        self,
+        *,
+        script_build_id: int,
+        publication_id: int,
+        manifest_version: ArtifactVersion,
+        structured_version: ArtifactVersion,
+        manifest: RootDeliveryManifestV1,
+        structured_script: StructuredScriptArtifactV1,
+        direction_markdown: str,
+        owner_token: MissionOwnerToken,
+    ) -> PublicationResult:
+        now = datetime.now(UTC)
+        async with self._sessions() as session:
+            async with session.begin():
+                # A published identity may be read back after a lost commit ACK,
+                # even though the build is already success. No writes occur on
+                # that path.
+                await self._fencing.verify_in_session(session, owner_token, allow_success=True)
+                publication = (
+                    (
+                        await session.execute(
+                            select(publication_table)
+                            .where(publication_table.c.id == publication_id)
+                            .with_for_update()
+                        )
+                    )
+                    .mappings()
+                    .one_or_none()
+                )
+                if publication is None:
+                    raise BuildNotFound()
+                if (
+                    int(publication["script_build_id"]) != script_build_id
+                    or publication["publication_type"] != PublicationType.FINAL.value
+                    or int(publication["artifact_version_id"])
+                    != manifest_version.artifact_version_id
+                    or publication["expected_sha256"]
+                    != manifest_version.canonical_sha256.removeprefix("sha256:")
+                ):
+                    raise ProtocolViolation("final publication identity differs from its manifest")
+                if publication["state"] == PublicationState.PUBLISHED.value:
+                    detail = await self._projection.project_in_session(
+                        script_build_id, session=session
+                    )
+                    readback = self._projection.to_canonical(detail)
+                    if readback.canonical_sha256 != manifest.legacy_projection_digest:
+                        raise PublicationReadbackMismatch()
+                    return PublicationResult(
+                        script_build_id=script_build_id,
+                        publication_id=publication_id,
+                        publication_type=PublicationType.FINAL,
+                        state=PublicationState.PUBLISHED,
+                        legacy_projection_digest=readback.canonical_sha256,
+                        committed=True,
+                    )
+                await self._fencing.verify_in_session(session, owner_token)
+                await session.execute(
+                    update(publication_table)
+                    .where(publication_table.c.id == publication_id)
+                    .values(
+                        state=PublicationState.PUBLISHING.value,
+                        attempt_count=publication_table.c.attempt_count + 1,
+                        publication_revision=publication_table.c.publication_revision + 1,
+                        updated_at=now,
+                        last_error_code=None,
+                        last_error_summary=None,
+                    )
+                )
+                await self._lock_branch_zero(session, script_build_id)
+                await session.execute(
+                    delete(script_build_paragraph_element).where(
+                        script_build_paragraph_element.c.script_build_id == script_build_id,
+                        script_build_paragraph_element.c.branch_id == 0,
+                    )
+                )
+                await session.execute(
+                    delete(script_build_paragraph).where(
+                        script_build_paragraph.c.script_build_id == script_build_id,
+                        script_build_paragraph.c.branch_id == 0,
+                    )
+                )
+                await session.execute(
+                    delete(script_build_element).where(
+                        script_build_element.c.script_build_id == script_build_id,
+                        script_build_element.c.branch_id == 0,
+                    )
+                )
+                paragraph_ids = await self._insert_paragraphs(
+                    session, script_build_id, structured_script, now
+                )
+                element_ids = await self._insert_elements(
+                    session, script_build_id, structured_script, now
+                )
+                await self._insert_links(
+                    session,
+                    script_build_id,
+                    structured_script,
+                    paragraph_ids,
+                    element_ids,
+                )
+                await session.execute(
+                    update(script_build_record)
+                    .where(script_build_record.c.id == script_build_id)
+                    .values(
+                        script_direction=direction_markdown,
+                        summary=manifest.build_summary,
+                        paragraph_count=len(
+                            [p for p in structured_script.paragraphs if p.is_active]
+                        ),
+                        element_count=len([e for e in structured_script.elements if e.is_active]),
+                    )
+                )
+                detail = await self._projection.project_in_session(script_build_id, session=session)
+                readback = self._projection.to_canonical(detail)
+                if readback.canonical_sha256 != manifest.legacy_projection_digest:
+                    raise PublicationReadbackMismatch()
+                await self._fencing.verify_in_session(session, owner_token)
+                artifact_ids = sorted(
+                    {manifest_version.artifact_version_id, structured_version.artifact_version_id}
+                )
+                locked_artifacts = list(
+                    (
+                        await session.execute(
+                            select(artifact_version_table.c.id)
+                            .where(artifact_version_table.c.id.in_(artifact_ids))
+                            .order_by(artifact_version_table.c.id)
+                            .with_for_update()
+                        )
+                    )
+                    .scalars()
+                    .all()
+                )
+                if locked_artifacts != artifact_ids:
+                    raise ProtocolViolation("final publication artifacts are missing")
+                pointer_result = await session.execute(
+                    update(mission_binding_table)
+                    .where(
+                        mission_binding_table.c.script_build_id == script_build_id,
+                        mission_binding_table.c.accepted_root_artifact_version_id.is_(None),
+                    )
+                    .values(
+                        accepted_root_artifact_version_id=manifest_version.artifact_version_id,
+                        updated_at=now,
+                    )
+                )
+                if pointer_result.rowcount != 1:
+                    pointer = await session.scalar(
+                        select(mission_binding_table.c.accepted_root_artifact_version_id).where(
+                            mission_binding_table.c.script_build_id == script_build_id
+                        )
+                    )
+                    if pointer != manifest_version.artifact_version_id:
+                        raise MissionFencingTokenStale()
+                artifact_result = await session.execute(
+                    update(artifact_version_table)
+                    .where(
+                        artifact_version_table.c.id.in_(artifact_ids),
+                        artifact_version_table.c.state.in_(
+                            [ArtifactState.FROZEN.value, ArtifactState.PUBLISHED.value]
+                        ),
+                    )
+                    .values(state=ArtifactState.PUBLISHED.value, published_at=now)
+                )
+                if artifact_result.rowcount != len(artifact_ids):
+                    raise ProtocolViolation("final publication artifacts are not publishable")
+                await session.execute(
+                    update(publication_table)
+                    .where(publication_table.c.id == publication_id)
+                    .values(
+                        state=PublicationState.PUBLISHED.value,
+                        updated_at=now,
+                        published_at=now,
+                    )
+                )
+                build_result = await session.execute(
+                    update(script_build_record)
+                    .where(script_build_record.c.id == script_build_id)
+                    .values(
+                        status="success",
+                        error_message=None,
+                        end_time=now,
+                    )
+                )
+                if build_result.rowcount != 1:
+                    raise BuildNotFound()
+            # The context above performs the only commit in the final UoW.
+        return PublicationResult(
+            script_build_id=script_build_id,
+            publication_id=publication_id,
+            publication_type=PublicationType.FINAL,
+            state=PublicationState.PUBLISHED,
+            legacy_projection_digest=manifest.legacy_projection_digest,
+            committed=True,
+        )
+
+    async def _lock_branch_zero(self, session: AsyncSession, script_build_id: int) -> None:
+        for table in (
+            script_build_paragraph_element,
+            script_build_paragraph,
+            script_build_element,
+        ):
+            await session.execute(
+                select(table.c.id)
+                .where(table.c.script_build_id == script_build_id, table.c.branch_id == 0)
+                .order_by(table.c.id)
+                .with_for_update()
+            )
+
+    async def _insert_paragraphs(
+        self,
+        session: AsyncSession,
+        script_build_id: int,
+        script: StructuredScriptArtifactV1,
+        now: datetime,
+    ) -> dict[int, int]:
+        active = sorted(
+            (item for item in script.paragraphs if item.is_active),
+            key=lambda item: (item.level, item.paragraph_index, item.paragraph_id),
+        )
+        identifiers: dict[int, int] = {}
+        for item in active:
+            parent_id = identifiers.get(item.parent_id) if item.parent_id is not None else None
+            if item.parent_id is not None and parent_id is None:
+                raise ProtocolViolation("final paragraph parent is outside the accepted script")
+            result = await session.execute(
+                insert(script_build_paragraph).values(
+                    script_build_id=script_build_id,
+                    branch_id=0,
+                    base_ref_id=None,
+                    paragraph_index=item.paragraph_index,
+                    level=item.level,
+                    parent_id=parent_id,
+                    name=item.name,
+                    content_range=item.content_range,
+                    theme=item.theme,
+                    form=item.form,
+                    function=item.function,
+                    feeling=item.feeling,
+                    theme_elements=list(item.theme_elements),
+                    form_elements=list(item.form_elements),
+                    function_elements=list(item.function_elements),
+                    feeling_elements=list(item.feeling_elements),
+                    description=item.description,
+                    full_description=item.full_description,
+                    is_active=True,
+                    created_at=now,
+                    updated_at=now,
+                )
+            )
+            identifiers[item.paragraph_id] = _inserted_id(result)
+        return identifiers
+
+    async def _insert_elements(
+        self,
+        session: AsyncSession,
+        script_build_id: int,
+        script: StructuredScriptArtifactV1,
+        now: datetime,
+    ) -> dict[int, int]:
+        identifiers: dict[int, int] = {}
+        for item in (value for value in script.elements if value.is_active):
+            result = await session.execute(
+                insert(script_build_element).values(
+                    script_build_id=script_build_id,
+                    branch_id=0,
+                    base_ref_id=None,
+                    name=item.name,
+                    dimension_primary=item.dimension_primary,
+                    dimension_secondary=item.dimension_secondary,
+                    commonality_analysis=item.commonality_analysis,
+                    topic_support=item.topic_support,
+                    weight_score=item.weight_score,
+                    support_elements=list(item.support_elements),
+                    is_active=True,
+                    created_at=now,
+                    updated_at=now,
+                )
+            )
+            identifiers[item.element_id] = _inserted_id(result)
+        return identifiers
+
+    async def _insert_links(
+        self,
+        session: AsyncSession,
+        script_build_id: int,
+        script: StructuredScriptArtifactV1,
+        paragraph_ids: dict[int, int],
+        element_ids: dict[int, int],
+    ) -> None:
+        for item in script.paragraph_element_links:
+            paragraph_id = paragraph_ids.get(item.paragraph_id)
+            element_id = element_ids.get(item.element_id)
+            if paragraph_id is None or element_id is None:
+                raise ProtocolViolation("final paragraph-element link is dangling")
+            await session.execute(
+                insert(script_build_paragraph_element).values(
+                    script_build_id=script_build_id,
+                    branch_id=0,
+                    paragraph_id=paragraph_id,
+                    element_id=element_id,
+                )
+            )
+
+
+def _publication(row: Any) -> Publication:
+    return Publication(
+        publication_id=int(row["id"]),
+        script_build_id=int(row["script_build_id"]),
+        publication_type=PublicationType(str(row["publication_type"])),
+        accept_decision_id=str(row["accept_decision_id"]),
+        artifact_version_id=int(row["artifact_version_id"]),
+        expected_sha256="sha256:" + str(row["expected_sha256"]),
+        state=PublicationState(str(row["state"])),
+        attempt_count=int(row["attempt_count"]),
+        publication_revision=int(row["publication_revision"]),
+        last_error_code=row["last_error_code"],
+        last_error_summary=row["last_error_summary"],
+        created_at=cast(datetime, row["created_at"]),
+        updated_at=cast(datetime, row["updated_at"]),
+        published_at=row["published_at"],
+    )
+
+
+def _inserted_id(result: Any) -> int:
+    primary_key = cast(CursorResult[Any], result).inserted_primary_key
+    if not primary_key or primary_key[0] is None:
+        raise ProtocolViolation("database did not return a final projection identity")
+    return int(primary_key[0])
+
+
+__all__ = ["SqlAlchemyFinalPublicationUnitOfWork"]