test_phase_three_publication.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397
  1. from __future__ import annotations
  2. import os
  3. from dataclasses import replace
  4. from datetime import UTC, datetime
  5. from pathlib import Path
  6. from typing import Any
  7. from uuid import uuid4
  8. import pytest
  9. from sqlalchemy import insert, select, text
  10. from sqlalchemy.ext.asyncio import (
  11. AsyncEngine,
  12. AsyncSession,
  13. async_sessionmaker,
  14. create_async_engine,
  15. )
  16. from script_build_host.application.legacy_projection import LegacyDetailProjectionService
  17. from script_build_host.domain.artifacts import (
  18. ArtifactState,
  19. DirectionGoal,
  20. ScriptDirectionArtifactV1,
  21. )
  22. from script_build_host.domain.errors import PublicationReadbackMismatch
  23. from script_build_host.domain.phase_three_artifacts import (
  24. LegacyProjectionCanonicalV1,
  25. RootDeliveryManifestV1,
  26. canonicalize_legacy_projection,
  27. )
  28. from script_build_host.domain.phase_two_artifacts import (
  29. ScriptElementV1,
  30. ScriptParagraphElementLinkV1,
  31. ScriptParagraphV1,
  32. StructuredScriptArtifactV1,
  33. )
  34. from script_build_host.domain.records import BuildStatus, PublicationState, PublicationType
  35. from script_build_host.infrastructure.final_publication import (
  36. SqlAlchemyFinalPublicationUnitOfWork,
  37. )
  38. from script_build_host.infrastructure.legacy_tables import (
  39. script_build_element,
  40. script_build_paragraph,
  41. script_build_record,
  42. )
  43. from script_build_host.infrastructure.ownership import FencedCommandGate, OwnerLease
  44. from script_build_host.infrastructure.tables import (
  45. artifact_version_table,
  46. mission_binding_table,
  47. publication_table,
  48. )
  49. from script_build_host.repositories.legacy_state import SqlAlchemyLegacyBuildStateRepository
  50. from script_build_host.repositories.sqlalchemy import (
  51. SqlAlchemyMissionBindingRepository,
  52. SqlAlchemyPublicationRepository,
  53. SqlAlchemyScriptBusinessArtifactRepository,
  54. )
  55. _CLOSURE = "sha256:" + "b" * 64
  56. class _Authorizer:
  57. async def require_access(self, principal: Any, script_build_id: int) -> None:
  58. return None
  59. class _MismatchingProjection(LegacyDetailProjectionService):
  60. def to_canonical(self, detail: Any) -> LegacyProjectionCanonicalV1:
  61. return replace(super().to_canonical(detail), summary="tampered")
  62. def _structured(direction_ref: str) -> StructuredScriptArtifactV1:
  63. return StructuredScriptArtifactV1(
  64. direction_ref=direction_ref,
  65. input_closure_digest=_CLOSURE,
  66. paragraphs=(
  67. ScriptParagraphV1(10, 1, 1, None, "opening", {"start": 0, "end": 4}),
  68. ScriptParagraphV1(20, 2, 2, 10, "detail", {"start": 4, "end": 8}),
  69. ),
  70. elements=(ScriptElementV1(30, "hook", "形式", "opening"),),
  71. paragraph_element_links=(ScriptParagraphElementLinkV1(20, 30),),
  72. source_artifact_refs=(),
  73. evidence_refs=(),
  74. acceptance_notes=(),
  75. )
  76. async def _publication_fixture(
  77. sessions: async_sessionmaker[AsyncSession],
  78. tmp_path: Path,
  79. *,
  80. publication_sessions: async_sessionmaker[AsyncSession] | None = None,
  81. branch_sessions: async_sessionmaker[AsyncSession] | None = None,
  82. prepare_publication: bool = True,
  83. ) -> dict[str, Any]:
  84. state = SqlAlchemyLegacyBuildStateRepository(sessions)
  85. root_trace_id = f"root-publication-{uuid4()}"
  86. build_id = await state.create(
  87. execution_id=1,
  88. topic_build_id=2,
  89. topic_id=3,
  90. agent_type="AigcAgent",
  91. agent_config={},
  92. data_source_url=None,
  93. strategies_config={},
  94. root_trace_id=root_trace_id,
  95. )
  96. await SqlAlchemyMissionBindingRepository(sessions).create(
  97. script_build_id=build_id,
  98. root_trace_id=root_trace_id,
  99. input_snapshot_id=1,
  100. engine_version="test",
  101. schema_version="test/v1",
  102. )
  103. now = datetime.now(UTC)
  104. async with (branch_sessions or sessions)() as session, session.begin():
  105. await session.execute(
  106. insert(script_build_paragraph).values(
  107. script_build_id=build_id,
  108. branch_id=0,
  109. paragraph_index=99,
  110. level=1,
  111. name="old",
  112. content_range={},
  113. is_active=True,
  114. created_at=now,
  115. updated_at=now,
  116. )
  117. )
  118. await session.execute(
  119. insert(script_build_element).values(
  120. script_build_id=build_id,
  121. branch_id=0,
  122. name="old",
  123. dimension_primary="形式",
  124. dimension_secondary="old",
  125. is_active=True,
  126. created_at=now,
  127. updated_at=now,
  128. )
  129. )
  130. artifacts = SqlAlchemyScriptBusinessArtifactRepository(sessions)
  131. direction, _ = await artifacts.freeze(
  132. script_build_id=build_id,
  133. task_id="direction",
  134. attempt_id=f"direction-{build_id}",
  135. spec_version=1,
  136. artifact=ScriptDirectionArtifactV1(
  137. goals=(DirectionGoal("goal", "deliver", "source"),),
  138. evidence_refs=("script-build://artifact-versions/777",),
  139. legacy_markdown="# Direction",
  140. ),
  141. )
  142. script = _structured(f"script-build://artifact-versions/{direction.artifact_version_id}")
  143. structured, _ = await artifacts.freeze(
  144. script_build_id=build_id,
  145. task_id="structured",
  146. attempt_id=f"structured-{build_id}",
  147. spec_version=1,
  148. artifact=script,
  149. )
  150. canonical = canonicalize_legacy_projection(
  151. script, direction="# Direction", summary="final summary"
  152. )
  153. manifest_value = RootDeliveryManifestV1(
  154. direction_ref=f"script-build://artifact-versions/{direction.artifact_version_id}",
  155. candidate_portfolio_ref="script-build://artifact-versions/999",
  156. structured_script_ref=(
  157. f"script-build://artifact-versions/{structured.artifact_version_id}"
  158. ),
  159. input_closure_digest=_CLOSURE,
  160. legacy_projection_digest=canonical.canonical_sha256,
  161. build_summary="final summary",
  162. )
  163. manifest, _ = await artifacts.freeze(
  164. script_build_id=build_id,
  165. task_id="root",
  166. attempt_id=f"root-{build_id}",
  167. spec_version=1,
  168. artifact=manifest_value,
  169. )
  170. publication = None
  171. if prepare_publication:
  172. publication = await SqlAlchemyPublicationRepository(
  173. publication_sessions or sessions
  174. ).prepare(
  175. script_build_id=build_id,
  176. publication_type=PublicationType.FINAL,
  177. accept_decision_id=f"accept-{build_id}",
  178. artifact_version_id=manifest.artifact_version_id,
  179. expected_sha256=manifest.canonical_sha256,
  180. )
  181. lease = OwnerLease(sessions, tmp_path / "leases", owner_instance_id=f"owner-{build_id}")
  182. token = await lease.acquire(build_id)
  183. return {
  184. "build_id": build_id,
  185. "script": script,
  186. "structured": structured,
  187. "manifest": manifest,
  188. "manifest_value": manifest.artifact,
  189. "publication": publication,
  190. "lease": lease,
  191. "token": token,
  192. }
  193. @pytest.mark.asyncio
  194. async def test_final_publication_commits_projection_pointer_artifacts_and_success(
  195. database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path
  196. ) -> None:
  197. _, sessions = database
  198. values = await _publication_fixture(sessions, tmp_path)
  199. projection = LegacyDetailProjectionService(sessions, _Authorizer()) # type: ignore[arg-type]
  200. uow = SqlAlchemyFinalPublicationUnitOfWork(sessions, projection, FencedCommandGate(sessions))
  201. try:
  202. result = await uow.publish(
  203. script_build_id=values["build_id"],
  204. publication_id=values["publication"].publication_id,
  205. manifest_version=values["manifest"],
  206. structured_version=values["structured"],
  207. manifest=values["manifest_value"],
  208. structured_script=values["script"],
  209. direction_markdown="# Direction",
  210. owner_token=values["token"],
  211. )
  212. assert result.committed is True
  213. replay = await uow.publish(
  214. script_build_id=values["build_id"],
  215. publication_id=values["publication"].publication_id,
  216. manifest_version=values["manifest"],
  217. structured_version=values["structured"],
  218. manifest=values["manifest_value"],
  219. structured_script=values["script"],
  220. direction_markdown="# Direction",
  221. owner_token=values["token"],
  222. )
  223. assert replay.committed is True
  224. async with sessions() as session:
  225. assert (
  226. await session.scalar(
  227. select(script_build_record.c.status).where(
  228. script_build_record.c.id == values["build_id"]
  229. )
  230. )
  231. == BuildStatus.SUCCESS.value
  232. )
  233. assert (
  234. await session.scalar(
  235. select(mission_binding_table.c.accepted_root_artifact_version_id).where(
  236. mission_binding_table.c.script_build_id == values["build_id"]
  237. )
  238. )
  239. == values["manifest"].artifact_version_id
  240. )
  241. assert set(
  242. (
  243. await session.execute(
  244. select(artifact_version_table.c.state).where(
  245. artifact_version_table.c.id.in_(
  246. [
  247. values["manifest"].artifact_version_id,
  248. values["structured"].artifact_version_id,
  249. ]
  250. )
  251. )
  252. )
  253. ).scalars()
  254. ) == {ArtifactState.PUBLISHED.value}
  255. finally:
  256. await values["lease"].release()
  257. @pytest.mark.asyncio
  258. async def test_final_publication_readback_failure_rolls_back_every_final_write(
  259. database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path
  260. ) -> None:
  261. _, sessions = database
  262. values = await _publication_fixture(sessions, tmp_path)
  263. projection = _MismatchingProjection(sessions, _Authorizer()) # type: ignore[arg-type]
  264. uow = SqlAlchemyFinalPublicationUnitOfWork(sessions, projection, FencedCommandGate(sessions))
  265. try:
  266. with pytest.raises(PublicationReadbackMismatch):
  267. await uow.publish(
  268. script_build_id=values["build_id"],
  269. publication_id=values["publication"].publication_id,
  270. manifest_version=values["manifest"],
  271. structured_version=values["structured"],
  272. manifest=values["manifest_value"],
  273. structured_script=values["script"],
  274. direction_markdown="# Direction",
  275. owner_token=values["token"],
  276. )
  277. async with sessions() as session:
  278. assert (
  279. await session.scalar(
  280. select(script_build_record.c.status).where(
  281. script_build_record.c.id == values["build_id"]
  282. )
  283. )
  284. == BuildStatus.RUNNING.value
  285. )
  286. assert (
  287. await session.scalar(
  288. select(script_build_paragraph.c.name).where(
  289. script_build_paragraph.c.script_build_id == values["build_id"],
  290. script_build_paragraph.c.branch_id == 0,
  291. )
  292. )
  293. == "old"
  294. )
  295. assert (
  296. await session.scalar(
  297. select(mission_binding_table.c.accepted_root_artifact_version_id).where(
  298. mission_binding_table.c.script_build_id == values["build_id"]
  299. )
  300. )
  301. is None
  302. )
  303. assert (
  304. await session.scalar(
  305. select(publication_table.c.state).where(
  306. publication_table.c.id == values["publication"].publication_id
  307. )
  308. )
  309. == PublicationState.PENDING.value
  310. )
  311. finally:
  312. await values["lease"].release()
  313. @pytest.mark.mysql
  314. @pytest.mark.asyncio
  315. async def test_mysql_final_publication_commits_one_atomic_projection(tmp_path: Path) -> None:
  316. runtime_dsn = os.environ.get("SCRIPT_BUILD_TEST_MYSQL_RUNTIME_DSN") or os.environ.get(
  317. "SCRIPT_BUILD_TEST_MYSQL_DSN"
  318. )
  319. final_dsn = os.environ.get("SCRIPT_BUILD_TEST_MYSQL_FINAL_DSN") or runtime_dsn
  320. if not runtime_dsn or not final_dsn:
  321. pytest.skip("SCRIPT_BUILD_TEST_MYSQL_DSN is not configured")
  322. runtime_engine = create_async_engine(
  323. runtime_dsn, pool_pre_ping=True, isolation_level="READ COMMITTED"
  324. )
  325. final_engine = create_async_engine(
  326. final_dsn, pool_pre_ping=True, isolation_level="READ COMMITTED"
  327. )
  328. sessions = async_sessionmaker(runtime_engine, expire_on_commit=False)
  329. final_sessions = async_sessionmaker(final_engine, expire_on_commit=False)
  330. values = await _publication_fixture(
  331. sessions,
  332. tmp_path,
  333. publication_sessions=final_sessions,
  334. branch_sessions=final_sessions,
  335. prepare_publication=False,
  336. )
  337. projection = LegacyDetailProjectionService( # type: ignore[arg-type]
  338. final_sessions,
  339. _Authorizer(),
  340. include_legacy_decision_traces=False,
  341. )
  342. uow = SqlAlchemyFinalPublicationUnitOfWork(
  343. final_sessions, projection, FencedCommandGate(final_sessions)
  344. )
  345. try:
  346. publication = await uow.reserve(
  347. script_build_id=values["build_id"],
  348. accept_decision_id=f"accept-{values['build_id']}",
  349. artifact_version_id=values["manifest"].artifact_version_id,
  350. expected_sha256=values["manifest"].canonical_sha256,
  351. owner_token=values["token"],
  352. )
  353. result = await uow.publish(
  354. script_build_id=values["build_id"],
  355. publication_id=publication.publication_id,
  356. manifest_version=values["manifest"],
  357. structured_version=values["structured"],
  358. manifest=values["manifest_value"],
  359. structured_script=values["script"],
  360. direction_markdown="# Direction",
  361. owner_token=values["token"],
  362. )
  363. assert result.committed is True
  364. async with final_sessions() as session:
  365. assert await session.scalar(text("SELECT @@transaction_isolation")) == (
  366. "READ-COMMITTED"
  367. )
  368. assert (
  369. await session.scalar(
  370. select(script_build_record.c.status).where(
  371. script_build_record.c.id == values["build_id"]
  372. )
  373. )
  374. == BuildStatus.SUCCESS.value
  375. )
  376. finally:
  377. await values["lease"].release()
  378. await runtime_engine.dispose()
  379. await final_engine.dispose()