| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407 |
- 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()
|