from __future__ import annotations import os from dataclasses import replace from datetime import UTC, datetime from pathlib import Path from typing import Any from uuid import uuid4 import pytest from sqlalchemy import insert, select, text from sqlalchemy.ext.asyncio import ( AsyncEngine, AsyncSession, async_sessionmaker, create_async_engine, ) from script_build_host.application.legacy_projection import LegacyDetailProjectionService from script_build_host.domain.artifacts import ( ArtifactState, DirectionArtifact, DirectionGoal, ) from script_build_host.domain.errors import PublicationReadbackMismatch from script_build_host.domain.phase_three_artifacts import ( LegacyProjectionCanonicalV1, RootDeliveryManifestV1, canonicalize_legacy_projection, ) from script_build_host.domain.phase_two_artifacts import ( GoalCoverage, ScriptElementV1, ScriptParagraphElementLinkV1, ScriptParagraphV1, StructuredScriptArtifactV1, ) from script_build_host.domain.records import BuildStatus, PublicationState, PublicationType from script_build_host.infrastructure.final_publication import ( SqlAlchemyFinalPublicationUnitOfWork, ) from script_build_host.infrastructure.legacy_tables import ( script_build_element, script_build_paragraph, script_build_record, ) from script_build_host.infrastructure.ownership import FencedCommandGate, OwnerLease from script_build_host.infrastructure.tables import ( artifact_version_table, mission_binding_table, publication_table, ) from script_build_host.repositories.legacy_state import SqlAlchemyLegacyBuildStateRepository from script_build_host.repositories.sqlalchemy import ( SqlAlchemyMissionBindingRepository, SqlAlchemyPublicationRepository, SqlAlchemyScriptBusinessArtifactRepository, ) _CLOSURE = "sha256:" + "b" * 64 class _Authorizer: async def require_access(self, principal: Any, script_build_id: int) -> None: return None class _MismatchingProjection(LegacyDetailProjectionService): def to_canonical(self, detail: Any) -> LegacyProjectionCanonicalV1: return replace(super().to_canonical(detail), summary="tampered") def _structured(direction_ref: str) -> StructuredScriptArtifactV1: return StructuredScriptArtifactV1( direction_ref=direction_ref, input_closure_digest=_CLOSURE, paragraphs=( ScriptParagraphV1(10, 1, 1, None, "opening", {"start": 0, "end": 4}), ScriptParagraphV1(20, 2, 2, 10, "detail", {"start": 4, "end": 8}), ), elements=(ScriptElementV1(30, "hook", "形式", "opening"),), paragraph_element_links=(ScriptParagraphElementLinkV1(20, 30),), source_artifact_refs=("script-build://artifact-versions/2",), goal_coverage=( GoalCoverage("goal-1", ("script-build://artifact-versions/2",)), ), evidence_refs=(), acceptance_notes=(), ) async def _publication_fixture( sessions: async_sessionmaker[AsyncSession], tmp_path: Path, *, publication_sessions: async_sessionmaker[AsyncSession] | None = None, branch_sessions: async_sessionmaker[AsyncSession] | None = None, prepare_publication: bool = True, ) -> dict[str, Any]: state = SqlAlchemyLegacyBuildStateRepository(sessions) root_trace_id = f"root-publication-{uuid4()}" build_id = await state.create( execution_id=1, topic_build_id=2, topic_id=3, agent_type="AigcAgent", agent_config={}, data_source_url=None, strategies_config={}, root_trace_id=root_trace_id, ) await SqlAlchemyMissionBindingRepository(sessions).create( script_build_id=build_id, root_trace_id=root_trace_id, input_snapshot_id=1, engine_version="test", schema_version="test/v1", ) now = datetime.now(UTC) async with (branch_sessions or sessions)() as session, session.begin(): await session.execute( insert(script_build_paragraph).values( script_build_id=build_id, branch_id=0, paragraph_index=99, level=1, name="old", content_range={}, is_active=True, created_at=now, updated_at=now, ) ) await session.execute( insert(script_build_element).values( script_build_id=build_id, branch_id=0, name="old", dimension_primary="形式", dimension_secondary="old", is_active=True, created_at=now, updated_at=now, ) ) artifacts = SqlAlchemyScriptBusinessArtifactRepository(sessions) direction, _ = await artifacts.freeze( script_build_id=build_id, task_id="direction", attempt_id=f"direction-{build_id}", spec_version=1, artifact=DirectionArtifact( goals=( DirectionGoal( "goal", "deliver", "source", success_criteria=("the delivery closes",), ), ), evidence_refs=("script-build://artifact-versions/777",), ), ) script = _structured(f"script-build://artifact-versions/{direction.artifact_version_id}") structured, _ = await artifacts.freeze( script_build_id=build_id, task_id="structured", attempt_id=f"structured-{build_id}", spec_version=1, artifact=script, ) canonical = canonicalize_legacy_projection( script, direction="# Direction", summary="final summary" ) manifest_value = RootDeliveryManifestV1( direction_ref=f"script-build://artifact-versions/{direction.artifact_version_id}", candidate_portfolio_ref="script-build://artifact-versions/999", structured_script_ref=( f"script-build://artifact-versions/{structured.artifact_version_id}" ), input_closure_digest=_CLOSURE, legacy_projection_digest=canonical.canonical_sha256, build_summary="final summary", ) manifest, _ = await artifacts.freeze( script_build_id=build_id, task_id="root", attempt_id=f"root-{build_id}", spec_version=1, artifact=manifest_value, ) publication = None if prepare_publication: publication = await SqlAlchemyPublicationRepository( publication_sessions or sessions ).prepare( script_build_id=build_id, publication_type=PublicationType.FINAL, accept_decision_id=f"accept-{build_id}", artifact_version_id=manifest.artifact_version_id, expected_sha256=manifest.canonical_sha256, ) lease = OwnerLease(sessions, tmp_path / "leases", owner_instance_id=f"owner-{build_id}") token = await lease.acquire(build_id) return { "build_id": build_id, "script": script, "structured": structured, "manifest": manifest, "manifest_value": manifest.artifact, "publication": publication, "lease": lease, "token": token, } @pytest.mark.asyncio async def test_final_publication_commits_projection_pointer_artifacts_and_success( database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path ) -> None: _, sessions = database values = await _publication_fixture(sessions, tmp_path) projection = LegacyDetailProjectionService(sessions, _Authorizer()) # type: ignore[arg-type] uow = SqlAlchemyFinalPublicationUnitOfWork(sessions, projection, FencedCommandGate(sessions)) try: result = await uow.publish( script_build_id=values["build_id"], publication_id=values["publication"].publication_id, manifest_version=values["manifest"], structured_version=values["structured"], manifest=values["manifest_value"], structured_script=values["script"], direction_markdown="# Direction", owner_token=values["token"], ) assert result.committed is True replay = await uow.publish( script_build_id=values["build_id"], publication_id=values["publication"].publication_id, manifest_version=values["manifest"], structured_version=values["structured"], manifest=values["manifest_value"], structured_script=values["script"], direction_markdown="# Direction", owner_token=values["token"], ) assert replay.committed is True async with sessions() as session: assert ( await session.scalar( select(script_build_record.c.status).where( script_build_record.c.id == values["build_id"] ) ) == BuildStatus.SUCCESS.value ) assert ( await session.scalar( select(mission_binding_table.c.accepted_root_artifact_version_id).where( mission_binding_table.c.script_build_id == values["build_id"] ) ) == values["manifest"].artifact_version_id ) assert set( ( await session.execute( select(artifact_version_table.c.state).where( artifact_version_table.c.id.in_( [ values["manifest"].artifact_version_id, values["structured"].artifact_version_id, ] ) ) ) ).scalars() ) == {ArtifactState.PUBLISHED.value} finally: await values["lease"].release() @pytest.mark.asyncio async def test_final_publication_readback_failure_rolls_back_every_final_write( database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path ) -> None: _, sessions = database values = await _publication_fixture(sessions, tmp_path) projection = _MismatchingProjection(sessions, _Authorizer()) # type: ignore[arg-type] uow = SqlAlchemyFinalPublicationUnitOfWork(sessions, projection, FencedCommandGate(sessions)) try: with pytest.raises(PublicationReadbackMismatch): await uow.publish( script_build_id=values["build_id"], publication_id=values["publication"].publication_id, manifest_version=values["manifest"], structured_version=values["structured"], manifest=values["manifest_value"], structured_script=values["script"], direction_markdown="# Direction", owner_token=values["token"], ) async with sessions() as session: assert ( await session.scalar( select(script_build_record.c.status).where( script_build_record.c.id == values["build_id"] ) ) == BuildStatus.RUNNING.value ) assert ( await session.scalar( select(script_build_paragraph.c.name).where( script_build_paragraph.c.script_build_id == values["build_id"], script_build_paragraph.c.branch_id == 0, ) ) == "old" ) assert ( await session.scalar( select(mission_binding_table.c.accepted_root_artifact_version_id).where( mission_binding_table.c.script_build_id == values["build_id"] ) ) is None ) assert ( await session.scalar( select(publication_table.c.state).where( publication_table.c.id == values["publication"].publication_id ) ) == PublicationState.PENDING.value ) finally: await values["lease"].release() @pytest.mark.mysql @pytest.mark.asyncio async def test_mysql_final_publication_commits_one_atomic_projection(tmp_path: Path) -> None: runtime_dsn = os.environ.get("SCRIPT_BUILD_TEST_MYSQL_RUNTIME_DSN") or os.environ.get( "SCRIPT_BUILD_TEST_MYSQL_DSN" ) final_dsn = os.environ.get("SCRIPT_BUILD_TEST_MYSQL_FINAL_DSN") or runtime_dsn if not runtime_dsn or not final_dsn: pytest.skip("SCRIPT_BUILD_TEST_MYSQL_DSN is not configured") runtime_engine = create_async_engine( runtime_dsn, pool_pre_ping=True, isolation_level="READ COMMITTED" ) final_engine = create_async_engine( final_dsn, pool_pre_ping=True, isolation_level="READ COMMITTED" ) sessions = async_sessionmaker(runtime_engine, expire_on_commit=False) final_sessions = async_sessionmaker(final_engine, expire_on_commit=False) values = await _publication_fixture( sessions, tmp_path, publication_sessions=final_sessions, branch_sessions=final_sessions, prepare_publication=False, ) projection = LegacyDetailProjectionService( # type: ignore[arg-type] final_sessions, _Authorizer(), include_legacy_decision_traces=False, ) uow = SqlAlchemyFinalPublicationUnitOfWork( final_sessions, projection, FencedCommandGate(final_sessions) ) try: publication = await uow.reserve( script_build_id=values["build_id"], accept_decision_id=f"accept-{values['build_id']}", artifact_version_id=values["manifest"].artifact_version_id, expected_sha256=values["manifest"].canonical_sha256, owner_token=values["token"], ) result = await uow.publish( script_build_id=values["build_id"], publication_id=publication.publication_id, manifest_version=values["manifest"], structured_version=values["structured"], manifest=values["manifest_value"], structured_script=values["script"], direction_markdown="# Direction", owner_token=values["token"], ) assert result.committed is True async with final_sessions() as session: assert await session.scalar(text("SELECT @@transaction_isolation")) == ( "READ-COMMITTED" ) assert ( await session.scalar( select(script_build_record.c.status).where( script_build_record.c.id == values["build_id"] ) ) == BuildStatus.SUCCESS.value ) finally: await values["lease"].release() await runtime_engine.dispose() await final_engine.dispose()