|
|
@@ -0,0 +1,393 @@
|
|
|
+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,
|
|
|
+ DirectionGoal,
|
|
|
+ ScriptDirectionArtifactV1,
|
|
|
+)
|
|
|
+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 (
|
|
|
+ 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=(),
|
|
|
+ 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=ScriptDirectionArtifactV1(
|
|
|
+ goals=(DirectionGoal("goal", "deliver", "source"),),
|
|
|
+ evidence_refs=("script-build://artifact-versions/777",),
|
|
|
+ legacy_markdown="# Direction",
|
|
|
+ ),
|
|
|
+ )
|
|
|
+ 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(final_sessions, _Authorizer()) # type: ignore[arg-type]
|
|
|
+ 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()
|