test_phase_two_sql_flow.py 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992
  1. from __future__ import annotations
  2. import asyncio
  3. import re
  4. from collections.abc import Mapping, Sequence
  5. from datetime import UTC, datetime
  6. from pathlib import Path
  7. from types import SimpleNamespace
  8. from typing import Any, cast
  9. import pytest
  10. from agent import FileSystemArtifactStore, FileSystemTaskStore, FileSystemTraceStore
  11. from agent.orchestration import (
  12. AgentRole,
  13. ArtifactRef,
  14. AttemptSubmission,
  15. CriterionResult,
  16. DecisionAction,
  17. TaskCoordinator,
  18. TaskStatus,
  19. ValidationVerdict,
  20. )
  21. from agent.orchestration.protocols import ValidatorRunResult, WorkerRunResult
  22. from sqlalchemy import event, func, select
  23. from script_build_host.application.phase_two_boundary import ScriptPhaseTwoBoundaryVerifier
  24. from script_build_host.application.phase_two_candidates import PhaseTwoCandidateService
  25. from script_build_host.application.phase_two_inputs import (
  26. AcceptedInputResolver,
  27. ActiveFrontierResolver,
  28. StoredTaskContractReader,
  29. )
  30. from script_build_host.application.phase_two_planning import (
  31. PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
  32. PhaseTwoPlanningService,
  33. )
  34. from script_build_host.domain.artifacts import (
  35. ArtifactKind,
  36. DirectionArtifact,
  37. DirectionGoal,
  38. EvidenceRecordV1,
  39. )
  40. from script_build_host.domain.errors import PhaseTwoBoundaryNotReady
  41. from script_build_host.domain.phase_two_artifacts import (
  42. CandidatePortfolioArtifactV1,
  43. StructuredScriptArtifactV1,
  44. )
  45. from script_build_host.domain.records import BuildStatus, MissionBinding
  46. from script_build_host.domain.task_contracts import ScriptTaskKind
  47. from script_build_host.infrastructure.legacy_tables import (
  48. script_build_element,
  49. script_build_paragraph,
  50. script_build_paragraph_element,
  51. script_build_record,
  52. )
  53. from script_build_host.infrastructure.tables import (
  54. artifact_version_table,
  55. publication_table,
  56. )
  57. from script_build_host.infrastructure.task_contract_store import FileScriptTaskContractStore
  58. from script_build_host.repositories.legacy_state import SqlAlchemyLegacyBuildStateRepository
  59. from script_build_host.repositories.sqlalchemy import (
  60. SqlAlchemyMissionBindingRepository,
  61. SqlAlchemyScriptBusinessArtifactRepository,
  62. )
  63. from script_build_host.repositories.workspace import SqlAlchemyCandidateWorkspaceRepository
  64. ROOT = "phase-two-sql-root"
  65. SNAPSHOT_ID = 11
  66. FULL_SCOPE = "script-build://scopes/full"
  67. FULL_WRITE = "script-build://writes/full"
  68. def _contract(
  69. kind: ScriptTaskKind,
  70. *,
  71. scope: str = FULL_SCOPE,
  72. write_scope: Sequence[str] = (FULL_WRITE,),
  73. input_refs: Sequence[dict[str, Any]] = (),
  74. base_ref: ArtifactRef | None = None,
  75. supersedes: Sequence[str] = (),
  76. closure: Sequence[dict[str, Any]] = (),
  77. adopted: Sequence[str] = (),
  78. held: Sequence[str] = (),
  79. order: Sequence[str] = (),
  80. comparison_refs: Sequence[dict[str, Any]] = (),
  81. execution_ready: bool = True,
  82. ) -> dict[str, Any]:
  83. del write_scope, closure, adopted, held, order, execution_ready
  84. if kind in {ScriptTaskKind.COMPOSE, ScriptTaskKind.CANDIDATE_PORTFOLIO}:
  85. intent = "compose" if kind is ScriptTaskKind.COMPOSE else "portfolio"
  86. elif supersedes:
  87. intent = "replace"
  88. else:
  89. intent = "explore"
  90. payload: dict[str, Any] = {
  91. "task_kind": kind.value,
  92. "scope_ref": scope,
  93. "intent_class": intent,
  94. "objective": f"produce independently verifiable {kind.value} content",
  95. "input_decision_ids": [str(item["decision_id"]) for item in input_refs],
  96. "base_decision_id": next(
  97. (
  98. str(item["decision_id"])
  99. for item in input_refs
  100. if base_ref is not None
  101. and cast(dict[str, Any], item["artifact_ref"])["uri"] == base_ref.uri
  102. ),
  103. None,
  104. ),
  105. "gap_ref": None,
  106. "criteria": [
  107. {
  108. "criterion_id": "closed",
  109. "description": "the frozen increment is concrete and independently verifiable",
  110. "hard": True,
  111. }
  112. ],
  113. "goal_ids": (
  114. []
  115. if kind in {ScriptTaskKind.DIRECTION, ScriptTaskKind.DECODE_RETRIEVAL}
  116. else ["goal-1"]
  117. ),
  118. "supersedes_decision_ids": list(supersedes),
  119. "comparison_decision_ids": [str(item["decision_id"]) for item in comparison_refs],
  120. }
  121. return payload
  122. def _decision_ref(
  123. decision_id: str,
  124. ref: ArtifactRef,
  125. *,
  126. scope: str,
  127. kind: ScriptTaskKind,
  128. ) -> dict[str, Any]:
  129. return {
  130. "decision_id": decision_id,
  131. "artifact_ref": _ref_payload(ref),
  132. "scope_ref": scope,
  133. "expected_task_kind": kind.value,
  134. }
  135. def _ref_payload(ref: ArtifactRef) -> dict[str, Any]:
  136. return {
  137. "uri": ref.uri,
  138. "kind": ref.kind,
  139. "version": ref.version,
  140. "digest": ref.digest,
  141. "summary": ref.summary,
  142. "metadata": ref.metadata,
  143. }
  144. def _task_kind(context: Mapping[str, Any]) -> ScriptTaskKind:
  145. refs = cast(dict[str, Any], context["task_spec"])["context_refs"]
  146. values = [value.rsplit("/", 1)[-1] for value in refs if "/task-kinds/" in value]
  147. assert len(values) == 1
  148. return ScriptTaskKind(values[0])
  149. class _ScriptedExecutor:
  150. """Drive real Coordinator stages through the same protected Host services."""
  151. def __init__(
  152. self,
  153. *,
  154. coordinator: TaskCoordinator,
  155. candidates: PhaseTwoCandidateService,
  156. artifacts: SqlAlchemyScriptBusinessArtifactRepository,
  157. script_build_id: int,
  158. ) -> None:
  159. self.coordinator = coordinator
  160. self.candidates = candidates
  161. self.artifacts = artifacts
  162. self.script_build_id = script_build_id
  163. async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
  164. kind = _task_kind(context)
  165. candidate_context = {
  166. "root_trace_id": context["root_trace_id"],
  167. "task_id": context["task_id"],
  168. "attempt_id": context["attempt_id"],
  169. "spec_version": context["spec_version"],
  170. }
  171. if kind is ScriptTaskKind.DECODE_RETRIEVAL:
  172. await self.artifacts.freeze(
  173. script_build_id=self.script_build_id,
  174. task_id=context["task_id"],
  175. attempt_id=context["attempt_id"],
  176. spec_version=context["spec_version"],
  177. artifact=EvidenceRecordV1(
  178. evidence_id=f"evidence-{context['attempt_id']}",
  179. source_type="decode",
  180. tool_name="retrieve_decode",
  181. query={"query": "opening tension", "top_k": 3},
  182. source_refs=("decode-index://fixture/1",),
  183. raw_artifact_ref=None,
  184. summary="A concrete opening needs a visible tension and an observable detail.",
  185. supports=("direction",),
  186. confidence="high",
  187. limitations=(),
  188. content_sha256="",
  189. created_at=datetime.now(UTC),
  190. ),
  191. )
  192. elif kind is ScriptTaskKind.DIRECTION:
  193. child_results = cast(list[dict[str, Any]], context["accepted_child_results"])
  194. assert len(child_results) == 1
  195. evidence_refs = child_results[0]["submission"]["evidence_refs"]
  196. await self.artifacts.freeze(
  197. script_build_id=self.script_build_id,
  198. task_id=context["task_id"],
  199. attempt_id=context["attempt_id"],
  200. spec_version=context["spec_version"],
  201. artifact=DirectionArtifact(
  202. goals=(
  203. DirectionGoal(
  204. "goal-1",
  205. "Use a concrete reversal to reveal the topic",
  206. "A visible reversal makes the direction observable",
  207. success_criteria=("the opening reveals one concrete reversal",),
  208. ),
  209. ),
  210. evidence_refs=(str(evidence_refs[0]["uri"]),),
  211. ),
  212. )
  213. elif kind in {ScriptTaskKind.STRUCTURE, ScriptTaskKind.PARAGRAPH}:
  214. workspace = await self.candidates.read_attempt_workspace(context=candidate_context)
  215. if workspace["paragraphs"]:
  216. paragraph_id = workspace["paragraphs"][0]["paragraph_id"]
  217. else:
  218. atom = (
  219. {"原子点": "可观察的反转", "维度": "开场", "维度类型": "主维度"},
  220. )
  221. created = await self.candidates.create_script_paragraphs(
  222. paragraphs=(
  223. {
  224. "client_key": "opening",
  225. "paragraph_index": 1,
  226. "name": f"{kind.value} opening",
  227. "content_range": {"scope": "opening"},
  228. "level": 1,
  229. "theme_elements": atom,
  230. "form_elements": atom,
  231. "function_elements": atom,
  232. "feeling_elements": atom,
  233. },
  234. ),
  235. context=candidate_context,
  236. )
  237. paragraph_id = created["paragraph_ids_by_client_key"]["opening"]
  238. if kind is ScriptTaskKind.PARAGRAPH:
  239. await self.candidates.batch_update_script_paragraphs(
  240. updates=(
  241. {
  242. "paragraph_id": paragraph_id,
  243. "theme": "a specific reversal",
  244. "form": "contrast",
  245. "function": "hook",
  246. "feeling": "curiosity",
  247. "description": (
  248. "The observable detail changes the initial interpretation."
  249. ),
  250. "full_description": (
  251. "The opening states a familiar assumption, then overturns it "
  252. "with one visible and source-grounded detail."
  253. ),
  254. },
  255. ),
  256. context=candidate_context,
  257. )
  258. elif kind is ScriptTaskKind.ELEMENT_SET:
  259. created = await self.candidates.create_script_element(
  260. payload={
  261. "name": "observable reversal detail",
  262. "dimension_primary": "实质",
  263. "dimension_secondary": "narrative evidence",
  264. "topic_support": {"scope": "opening"},
  265. "weight_score": {"score": 0.9},
  266. },
  267. context=candidate_context,
  268. )
  269. workspace = await self.candidates.read_attempt_workspace(context=candidate_context)
  270. if workspace["paragraphs"]:
  271. await self.candidates.batch_link_paragraph_elements(
  272. links=(
  273. {
  274. "paragraph_id": workspace["paragraphs"][0]["paragraph_id"],
  275. "element_ids": [created["element_id"]],
  276. },
  277. ),
  278. context=candidate_context,
  279. )
  280. elif kind is ScriptTaskKind.COMPARE:
  281. pinned = await self.candidates.read_pinned_candidates(context=candidate_context)
  282. candidate_refs = [
  283. str(item["artifact_ref"]["uri"])
  284. for item in pinned
  285. if item["artifact_ref"]["kind"] == ArtifactKind.PARAGRAPH.value
  286. ]
  287. await self.candidates.save_comparison_candidate(
  288. payload={
  289. "criterion_results": [
  290. {
  291. "criterion_id": "closed",
  292. "candidate_results": [
  293. {"artifact_ref": ref, "reason": "immutable candidate reviewed"}
  294. for ref in candidate_refs
  295. ],
  296. }
  297. ],
  298. "conflicts": ["the replacement has the more concrete opening"],
  299. "recommendation": candidate_refs[-1],
  300. },
  301. context=candidate_context,
  302. )
  303. elif kind is ScriptTaskKind.COMPOSE:
  304. await self.candidates.save_structured_script_candidate(
  305. acceptance_notes=("all adopted increments are pinned by ACCEPT decisions",),
  306. context=candidate_context,
  307. )
  308. elif kind is ScriptTaskKind.CANDIDATE_PORTFOLIO:
  309. await self.candidates.save_candidate_portfolio(
  310. payload={"unresolved_defects": []}, context=candidate_context
  311. )
  312. else: # pragma: no cover - this executor is intentionally phase-two bounded
  313. raise AssertionError(f"unexpected worker kind: {kind.value}")
  314. manifest = await self.candidates.resolve_attempt_manifest(context=candidate_context)
  315. safe_ref = ArtifactRef(
  316. manifest.artifact_ref.uri,
  317. manifest.artifact_ref.kind,
  318. manifest.artifact_ref.version,
  319. manifest.artifact_ref.digest,
  320. )
  321. await self.coordinator.submit_attempt(
  322. {
  323. "role": AgentRole.WORKER.value,
  324. "root_trace_id": context["root_trace_id"],
  325. "task_id": context["task_id"],
  326. "attempt_id": context["attempt_id"],
  327. "spec_version": context["spec_version"],
  328. "trace_id": context["worker_trace_id"],
  329. "tool_call_id": f"submit:{context['attempt_id']}",
  330. "operation_id": context["operation_id"],
  331. "execution_epoch": context["execution_epoch"],
  332. },
  333. AttemptSubmission(
  334. summary=f"kind={safe_ref.kind};status=frozen",
  335. artifact_refs=[] if safe_ref.kind == ArtifactKind.EVIDENCE.value else [safe_ref],
  336. evidence_refs=[safe_ref] if safe_ref.kind == ArtifactKind.EVIDENCE.value else [],
  337. ),
  338. )
  339. return WorkerRunResult(context["worker_trace_id"], "completed")
  340. async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
  341. criteria = cast(dict[str, Any], context["task_spec"])["acceptance_criteria"]
  342. results = [
  343. CriterionResult(
  344. criterion_id=item["criterion_id"],
  345. verdict=ValidationVerdict.PASSED,
  346. reason="the immutable SQL artifact is concrete and owner/digest closed",
  347. )
  348. for item in criteria
  349. ]
  350. await self.coordinator.submit_validation(
  351. {
  352. "role": AgentRole.VALIDATOR.value,
  353. "root_trace_id": context["root_trace_id"],
  354. "task_id": context["task_id"],
  355. "attempt_id": context["attempt_id"],
  356. "validation_id": context["validation_id"],
  357. "snapshot_id": context["snapshot_id"],
  358. "trace_id": context["validator_trace_id"],
  359. "tool_call_id": f"validate:{context['validation_id']}",
  360. "operation_id": context["operation_id"],
  361. "execution_epoch": context["execution_epoch"],
  362. },
  363. ValidationVerdict.PASSED,
  364. results,
  365. "all hard criteria passed against the immutable snapshot",
  366. (),
  367. (),
  368. (),
  369. "accept",
  370. )
  371. return ValidatorRunResult(context["validator_trace_id"], "completed")
  372. async def stop(self, trace_id: str) -> bool:
  373. del trace_id
  374. return True
  375. async def _plan_one(
  376. planning: PhaseTwoPlanningService,
  377. payload: Mapping[str, Any],
  378. *,
  379. parent_task_id: str | None,
  380. call_id: str,
  381. ) -> str:
  382. result = await planning.plan_script_tasks(
  383. contract_payloads=(payload,),
  384. parent_task_id=parent_task_id,
  385. context={"root_trace_id": ROOT, "tool_call_id": call_id},
  386. )
  387. return cast(str, result["task_ids"][0])
  388. async def _dispatch_accept(
  389. planning: PhaseTwoPlanningService, coordinator: TaskCoordinator, task_id: str
  390. ) -> tuple[str, ArtifactRef]:
  391. cycles = await planning.dispatch_script_tasks(
  392. task_ids=(task_id,),
  393. context={"root_trace_id": ROOT, "tool_call_id": f"dispatch:{task_id}"},
  394. )
  395. assert cycles[0]["error"] is None
  396. validation_id = cast(str, cycles[0]["validation_id"])
  397. accepted = await planning.decide_script_task(
  398. task_id=task_id,
  399. action=DecisionAction.ACCEPT.value,
  400. reason="independent validation passed",
  401. validation_id=validation_id,
  402. replacement_contract=None,
  403. child_contracts=(),
  404. context={"root_trace_id": ROOT, "tool_call_id": f"accept:{task_id}"},
  405. )
  406. decision_id = cast(str, accepted["decision_id"])
  407. ledger = await coordinator.task_store.load(ROOT)
  408. decision = ledger.decisions[decision_id]
  409. attempt = ledger.attempts[cast(str, decision.attempt_id)]
  410. assert attempt.submission is not None
  411. refs = [*attempt.submission.artifact_refs, *attempt.submission.evidence_refs]
  412. assert len(refs) == 1
  413. return decision_id, refs[0]
  414. @pytest.mark.asyncio
  415. async def test_sql_phase_two_dynamic_replacement_compose_portfolio_and_boundary(
  416. database: Any, tmp_path: Path
  417. ) -> None:
  418. engine, sessions = database
  419. mutation_statements: list[str] = []
  420. def observe_sql(
  421. _connection: Any,
  422. _cursor: Any,
  423. statement: str,
  424. _parameters: Any,
  425. _context: Any,
  426. _executemany: bool,
  427. ) -> None:
  428. if statement.lstrip().upper().startswith(("INSERT ", "UPDATE ", "DELETE ")):
  429. mutation_statements.append(statement)
  430. event.listen(engine.sync_engine, "before_cursor_execute", observe_sql)
  431. build_states = SqlAlchemyLegacyBuildStateRepository(sessions)
  432. script_build_id = await build_states.create(
  433. execution_id=101,
  434. topic_build_id=202,
  435. topic_id=303,
  436. agent_type="script-planner",
  437. agent_config={},
  438. data_source_url=None,
  439. strategies_config={},
  440. )
  441. bindings = SqlAlchemyMissionBindingRepository(sessions)
  442. await bindings.create(
  443. script_build_id=script_build_id,
  444. root_trace_id=ROOT,
  445. input_snapshot_id=SNAPSHOT_ID,
  446. engine_version="phase-two-test",
  447. schema_version="v1",
  448. )
  449. task_path = tmp_path / "task-ledger"
  450. artifact_path = tmp_path / "framework-artifacts"
  451. contract_path = tmp_path / "contract-data"
  452. task_store = FileSystemTaskStore(str(task_path))
  453. framework_artifacts = FileSystemArtifactStore(str(artifact_path))
  454. coordinator = TaskCoordinator(
  455. task_store,
  456. framework_artifacts,
  457. FileSystemTraceStore(str(tmp_path / "traces")),
  458. )
  459. await coordinator.ensure_ledger(
  460. ROOT,
  461. {
  462. "objective": "produce a governed phase-two candidate portfolio",
  463. "acceptance_criteria": [
  464. {"criterion_id": "portfolio", "description": "portfolio closed", "hard": True}
  465. ],
  466. "context_refs": [f"script-build://inputs/{SNAPSHOT_ID}"],
  467. },
  468. )
  469. contract_store = FileScriptTaskContractStore(contract_path)
  470. contract_reader = StoredTaskContractReader(contract_store)
  471. artifacts = SqlAlchemyScriptBusinessArtifactRepository(sessions)
  472. accepted_inputs = AcceptedInputResolver(
  473. task_store=task_store,
  474. artifact_store=framework_artifacts,
  475. artifacts=artifacts,
  476. contracts=contract_reader,
  477. )
  478. workspaces = SqlAlchemyCandidateWorkspaceRepository(sessions, artifacts)
  479. candidates = PhaseTwoCandidateService(
  480. bindings=bindings,
  481. task_store=task_store,
  482. framework_artifact_store=framework_artifacts,
  483. artifacts=artifacts,
  484. accepted_inputs=accepted_inputs,
  485. active_frontier=ActiveFrontierResolver(),
  486. workspaces=workspaces,
  487. )
  488. boundary = ScriptPhaseTwoBoundaryVerifier(
  489. task_store=task_store,
  490. artifact_store=framework_artifacts,
  491. artifacts=artifacts,
  492. bindings=bindings,
  493. accepted_inputs=accepted_inputs,
  494. )
  495. planning = PhaseTwoPlanningService(
  496. coordinator=coordinator,
  497. bindings=bindings,
  498. contracts=contract_store,
  499. artifacts=artifacts,
  500. closure_gate=boundary,
  501. )
  502. coordinator.set_executor(
  503. _ScriptedExecutor(
  504. coordinator=coordinator,
  505. candidates=candidates,
  506. artifacts=artifacts,
  507. script_build_id=script_build_id,
  508. )
  509. )
  510. direction = await _plan_one(
  511. planning, _contract(ScriptTaskKind.DIRECTION), parent_task_id=None, call_id="direction"
  512. )
  513. retrieval = await _plan_one(
  514. planning,
  515. _contract(ScriptTaskKind.DECODE_RETRIEVAL, write_scope=()),
  516. parent_task_id=direction,
  517. call_id="retrieval",
  518. )
  519. await _dispatch_accept(planning, coordinator, retrieval)
  520. _, direction_ref = await _dispatch_accept(planning, coordinator, direction)
  521. await bindings.set_active_direction(
  522. script_build_id=script_build_id,
  523. artifact_version_id=int(direction_ref.version),
  524. )
  525. portfolio = await _plan_one(
  526. planning,
  527. _contract(
  528. ScriptTaskKind.CANDIDATE_PORTFOLIO,
  529. write_scope=("script-build://writes",),
  530. execution_ready=False,
  531. ),
  532. parent_task_id=None,
  533. call_id="portfolio",
  534. )
  535. compose = await _plan_one(
  536. planning,
  537. _contract(
  538. ScriptTaskKind.COMPOSE,
  539. write_scope=("script-build://writes",),
  540. execution_ready=False,
  541. ),
  542. parent_task_id=portfolio,
  543. call_id="compose",
  544. )
  545. # Exploration order is intentionally Element -> Paragraph -> Structure.
  546. old_element = await _plan_one(
  547. planning,
  548. _contract(
  549. ScriptTaskKind.ELEMENT_SET,
  550. scope="script-build://scopes/full/elements",
  551. write_scope=("script-build://writes/elements",),
  552. ),
  553. parent_task_id=compose,
  554. call_id="element-first",
  555. )
  556. old_element_decision, old_element_ref = await _dispatch_accept(
  557. planning, coordinator, old_element
  558. )
  559. old_paragraph = await _plan_one(
  560. planning,
  561. _contract(
  562. ScriptTaskKind.PARAGRAPH,
  563. scope="script-build://scopes/full/opening",
  564. write_scope=("script-build://writes/paragraphs/opening",),
  565. ),
  566. parent_task_id=compose,
  567. call_id="paragraph-second",
  568. )
  569. old_paragraph_decision, old_paragraph_ref = await _dispatch_accept(
  570. planning, coordinator, old_paragraph
  571. )
  572. structure = await _plan_one(
  573. planning,
  574. _contract(
  575. ScriptTaskKind.STRUCTURE,
  576. scope="script-build://scopes/full/opening",
  577. write_scope=("script-build://writes/paragraphs/opening",),
  578. ),
  579. parent_task_id=compose,
  580. call_id="structure-third",
  581. )
  582. structure_decision, structure_ref = await _dispatch_accept(planning, coordinator, structure)
  583. paragraph_replacement = await _plan_one(
  584. planning,
  585. _contract(
  586. ScriptTaskKind.PARAGRAPH,
  587. scope="script-build://scopes/full/opening",
  588. write_scope=("script-build://writes/paragraphs/opening",),
  589. input_refs=(
  590. _decision_ref(
  591. structure_decision,
  592. structure_ref,
  593. scope="script-build://scopes/full/opening",
  594. kind=ScriptTaskKind.STRUCTURE,
  595. ),
  596. _decision_ref(
  597. old_paragraph_decision,
  598. old_paragraph_ref,
  599. scope="script-build://scopes/full/opening",
  600. kind=ScriptTaskKind.PARAGRAPH,
  601. ),
  602. ),
  603. base_ref=structure_ref,
  604. supersedes=(old_paragraph_decision,),
  605. ),
  606. parent_task_id=compose,
  607. call_id="replace-paragraph",
  608. )
  609. paragraph_decision, paragraph_ref = await _dispatch_accept(
  610. planning, coordinator, paragraph_replacement
  611. )
  612. element_replacement = await _plan_one(
  613. planning,
  614. _contract(
  615. ScriptTaskKind.ELEMENT_SET,
  616. scope="script-build://scopes/full/elements",
  617. write_scope=("script-build://writes/elements",),
  618. input_refs=(
  619. _decision_ref(
  620. old_element_decision,
  621. old_element_ref,
  622. scope="script-build://scopes/full/elements",
  623. kind=ScriptTaskKind.ELEMENT_SET,
  624. ),
  625. _decision_ref(
  626. paragraph_decision,
  627. paragraph_ref,
  628. scope="script-build://scopes/full/opening",
  629. kind=ScriptTaskKind.PARAGRAPH,
  630. ),
  631. ),
  632. base_ref=paragraph_ref,
  633. supersedes=(old_element_decision,),
  634. ),
  635. parent_task_id=compose,
  636. call_id="replace-element",
  637. )
  638. element_decision, _element_ref = await _dispatch_accept(
  639. planning, coordinator, element_replacement
  640. )
  641. comparison = await _plan_one(
  642. planning,
  643. _contract(
  644. ScriptTaskKind.COMPARE,
  645. scope="script-build://scopes/full/opening",
  646. write_scope=(),
  647. comparison_refs=(
  648. _decision_ref(
  649. old_paragraph_decision,
  650. old_paragraph_ref,
  651. scope="script-build://scopes/full/opening",
  652. kind=ScriptTaskKind.PARAGRAPH,
  653. ),
  654. _decision_ref(
  655. paragraph_decision,
  656. paragraph_ref,
  657. scope="script-build://scopes/full/opening",
  658. kind=ScriptTaskKind.PARAGRAPH,
  659. ),
  660. ),
  661. ),
  662. parent_task_id=compose,
  663. call_id="compare-opening-candidates",
  664. )
  665. await _dispatch_accept(planning, coordinator, comparison)
  666. with pytest.raises(PhaseTwoBoundaryNotReady, match="CandidatePortfolio"):
  667. await planning.decide_script_task(
  668. task_id=(await task_store.load(ROOT)).root_task_id,
  669. action=DecisionAction.BLOCK.value,
  670. reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
  671. validation_id=None,
  672. replacement_contract=None,
  673. child_contracts=(),
  674. context={"root_trace_id": ROOT, "tool_call_id": "early-boundary"},
  675. )
  676. adopted = (structure_decision, paragraph_decision, element_decision)
  677. await planning.decide_script_task(
  678. task_id=compose,
  679. action=DecisionAction.REVISE.value,
  680. reason="freeze the selected active frontier and formal compose order",
  681. validation_id=None,
  682. replacement_contract=None,
  683. child_contracts=(),
  684. selected_decision_ids=adopted,
  685. context={"root_trace_id": ROOT, "tool_call_id": "close-compose"},
  686. )
  687. compose_decision, compose_ref = await _dispatch_accept(planning, coordinator, compose)
  688. await planning.decide_script_task(
  689. task_id=portfolio,
  690. action=DecisionAction.REVISE.value,
  691. reason="freeze the uniquely adopted StructuredScript",
  692. validation_id=None,
  693. replacement_contract=None,
  694. child_contracts=(),
  695. selected_decision_ids=(compose_decision,),
  696. context={"root_trace_id": ROOT, "tool_call_id": "close-portfolio"},
  697. )
  698. portfolio_decision, _ = await _dispatch_accept(planning, coordinator, portfolio)
  699. ledger = await task_store.load(ROOT)
  700. compose_attempt = ledger.attempts[cast(str, ledger.decisions[compose_decision].attempt_id)]
  701. frontier = await candidates.read_active_frontier(
  702. context={
  703. "root_trace_id": ROOT,
  704. "task_id": compose,
  705. "attempt_id": compose_attempt.attempt_id,
  706. "spec_version": compose_attempt.spec_version,
  707. }
  708. )
  709. assert [item["decision_id"] for item in frontier["inputs"]] == list(adopted)
  710. assert old_element_decision not in {item["decision_id"] for item in frontier["inputs"]}
  711. assert old_paragraph_decision not in {item["decision_id"] for item in frontier["inputs"]}
  712. await planning.decide_script_task(
  713. task_id=ledger.root_task_id,
  714. action=DecisionAction.BLOCK.value,
  715. reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
  716. validation_id=None,
  717. replacement_contract=None,
  718. child_contracts=(),
  719. context={"root_trace_id": ROOT, "tool_call_id": "phase-two-boundary"},
  720. )
  721. await build_states.set_checkpoint(
  722. script_build_id,
  723. checkpoint_code=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
  724. summary="one validated CandidatePortfolio is ready for phase three",
  725. root_trace_id=ROOT,
  726. )
  727. reloaded = await FileSystemTaskStore(str(task_path)).load(ROOT)
  728. assert reloaded.tasks[reloaded.root_task_id].status is TaskStatus.BLOCKED
  729. assert reloaded.tasks[reloaded.root_task_id].blocked_reason == (
  730. PHASE_TWO_CANDIDATE_PORTFOLIO_READY
  731. )
  732. compose_children = [reloaded.tasks[item] for item in reloaded.tasks[compose].child_task_ids]
  733. stored_reader = StoredTaskContractReader(FileScriptTaskContractStore(contract_path))
  734. assert [
  735. (
  736. await stored_reader.read_for_task(
  737. root_trace_id=ROOT,
  738. task=item,
  739. spec_version=item.current_spec_version,
  740. )
  741. ).task_kind
  742. for item in compose_children[:3]
  743. ] == [
  744. ScriptTaskKind.ELEMENT_SET,
  745. ScriptTaskKind.PARAGRAPH,
  746. ScriptTaskKind.STRUCTURE,
  747. ]
  748. compose_version = await artifacts.get_by_attempt(
  749. script_build_id=script_build_id,
  750. task_id=compose,
  751. attempt_id=compose_attempt.attempt_id,
  752. )
  753. assert isinstance(compose_version.artifact, StructuredScriptArtifactV1)
  754. assert len(compose_version.artifact.paragraphs) == 1
  755. assert compose_version.artifact.elements
  756. assert compose_version.artifact.paragraph_element_links
  757. portfolio_attempt = reloaded.attempts[
  758. cast(str, reloaded.decisions[portfolio_decision].attempt_id)
  759. ]
  760. portfolio_version = await artifacts.get_by_attempt(
  761. script_build_id=script_build_id,
  762. task_id=portfolio,
  763. attempt_id=portfolio_attempt.attempt_id,
  764. )
  765. assert isinstance(portfolio_version.artifact, CandidatePortfolioArtifactV1)
  766. assert portfolio_version.artifact.accepted_decision_ids == (compose_decision,)
  767. assert portfolio_version.artifact.adopted_structured_script_ref == compose_ref.uri
  768. async with sessions() as session:
  769. build = (
  770. (
  771. await session.execute(
  772. select(script_build_record).where(script_build_record.c.id == script_build_id)
  773. )
  774. )
  775. .mappings()
  776. .one()
  777. )
  778. branch_values = [
  779. *(
  780. await session.scalars(
  781. select(script_build_paragraph.c.branch_id).where(
  782. script_build_paragraph.c.script_build_id == script_build_id
  783. )
  784. )
  785. ),
  786. *(
  787. await session.scalars(
  788. select(script_build_element.c.branch_id).where(
  789. script_build_element.c.script_build_id == script_build_id
  790. )
  791. )
  792. ),
  793. *(
  794. await session.scalars(
  795. select(script_build_paragraph_element.c.branch_id).where(
  796. script_build_paragraph_element.c.script_build_id == script_build_id
  797. )
  798. )
  799. ),
  800. ]
  801. final_publications = await session.scalar(
  802. select(func.count())
  803. .select_from(publication_table)
  804. .where(
  805. publication_table.c.script_build_id == script_build_id,
  806. publication_table.c.publication_type == "final",
  807. )
  808. )
  809. phase_two_rows = (
  810. (
  811. await session.execute(
  812. select(artifact_version_table).where(
  813. artifact_version_table.c.script_build_id == script_build_id,
  814. artifact_version_table.c.artifact_type.in_(
  815. [
  816. ArtifactKind.STRUCTURE.value,
  817. ArtifactKind.PARAGRAPH.value,
  818. ArtifactKind.ELEMENT_SET.value,
  819. ArtifactKind.COMPARISON.value,
  820. ArtifactKind.STRUCTURED_SCRIPT.value,
  821. ArtifactKind.CANDIDATE_PORTFOLIO.value,
  822. ]
  823. ),
  824. )
  825. )
  826. )
  827. .mappings()
  828. .all()
  829. )
  830. mutation_targets = {
  831. match.group(1).strip('`"').lower()
  832. for statement in mutation_statements
  833. if (
  834. match := re.match(
  835. r"\s*(?:INSERT\s+INTO|UPDATE|DELETE\s+FROM)\s+([`\"\w]+)",
  836. statement,
  837. re.IGNORECASE,
  838. )
  839. )
  840. }
  841. assert mutation_targets <= {
  842. "script_build_record",
  843. "script_build_mission_binding",
  844. "script_build_artifact_version",
  845. "script_build_publication",
  846. "script_build_paragraph",
  847. "script_build_element",
  848. "script_build_paragraph_element",
  849. "script_build_task_plan_step",
  850. }
  851. assert not mutation_targets & {
  852. "script_build_branch",
  853. "script_build_round",
  854. "script_build_multipath",
  855. }
  856. assert BuildStatus(str(build["status"])) is BuildStatus.PARTIAL
  857. assert build["error_message"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
  858. assert build["reson_trace_id"] == ROOT and build["end_time"] is not None
  859. assert branch_values and min(branch_values) > 0 and 0 not in branch_values
  860. assert final_publications == 0
  861. assert (
  862. sum(
  863. row["artifact_type"] == ArtifactKind.CANDIDATE_PORTFOLIO.value for row in phase_two_rows
  864. )
  865. == 1
  866. )
  867. for row in phase_two_rows:
  868. if row["artifact_type"] in {
  869. ArtifactKind.STRUCTURE.value,
  870. ArtifactKind.PARAGRAPH.value,
  871. ArtifactKind.ELEMENT_SET.value,
  872. ArtifactKind.STRUCTURED_SCRIPT.value,
  873. }:
  874. assert row["legacy_branch_id"] == row["id"] > 0
  875. else:
  876. assert row["legacy_branch_id"] is None
  877. class _SingleBinding:
  878. def __init__(self, root: str) -> None:
  879. now = datetime.now(UTC)
  880. self.value = MissionBinding(1, 1, root, 11, None, None, "test", "v1", now, now)
  881. async def get_by_root(self, root_trace_id: str) -> MissionBinding:
  882. assert root_trace_id == self.value.root_trace_id
  883. return self.value
  884. class _BlockingExecutor:
  885. def __init__(self) -> None:
  886. self.started = asyncio.Event()
  887. self.release = asyncio.Event()
  888. async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
  889. self.started.set()
  890. await self.release.wait()
  891. return WorkerRunResult(context["worker_trace_id"], "failed", error="test release")
  892. async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
  893. raise AssertionError(f"unexpected validation: {context}")
  894. async def stop(self, trace_id: str) -> bool:
  895. del trace_id
  896. self.release.set()
  897. return True
  898. @pytest.mark.asyncio
  899. async def test_same_task_cannot_reserve_parallel_attempts(tmp_path: Path) -> None:
  900. root = "parallel-attempt-root"
  901. task_store = FileSystemTaskStore(str(tmp_path / "ledger"))
  902. executor = _BlockingExecutor()
  903. coordinator = TaskCoordinator(
  904. task_store,
  905. FileSystemArtifactStore(str(tmp_path / "artifacts")),
  906. FileSystemTraceStore(str(tmp_path / "traces")),
  907. executor=executor,
  908. )
  909. await coordinator.ensure_ledger(
  910. root,
  911. {
  912. "objective": "parallel attempt guard",
  913. "acceptance_criteria": [
  914. {"criterion_id": "guard", "description": "guard", "hard": True}
  915. ],
  916. "context_refs": ["script-build://inputs/11"],
  917. },
  918. )
  919. planning = PhaseTwoPlanningService(
  920. coordinator=coordinator,
  921. bindings=cast(Any, _SingleBinding(root)),
  922. contracts=FileScriptTaskContractStore(tmp_path / "contracts"),
  923. artifacts=cast(Any, SimpleNamespace()),
  924. )
  925. task = await planning.plan_script_tasks(
  926. contract_payloads=(_contract(ScriptTaskKind.DIRECTION),),
  927. parent_task_id=None,
  928. context={"root_trace_id": root},
  929. )
  930. task_id = task["task_ids"][0]
  931. first = asyncio.create_task(
  932. planning.dispatch_script_tasks(task_ids=(task_id,), context={"root_trace_id": root})
  933. )
  934. await asyncio.wait_for(executor.started.wait(), timeout=2)
  935. second = await planning.dispatch_script_tasks(
  936. task_ids=(task_id,), context={"root_trace_id": root}
  937. )
  938. assert "Dispatch conflict" in cast(str, second[0]["error"])
  939. executor.release.set()
  940. await first
  941. ledger = await task_store.load(root)
  942. assert len(ledger.tasks[task_id].attempt_ids) == 1