test_phase_three_publication.py 15 KB

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