test_phase_two_sql_flow.py 37 KB

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