test_phase_two_sql_flow.py 38 KB

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