test_phase_two_contracts.py 53 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341
  1. from __future__ import annotations
  2. import asyncio
  3. from dataclasses import replace
  4. from datetime import UTC, datetime, timedelta
  5. from hashlib import sha256
  6. from pathlib import Path
  7. from types import SimpleNamespace
  8. from typing import Any, cast
  9. import pytest
  10. from agent import AgentRunner, FileSystemArtifactStore, FileSystemTaskStore, FileSystemTraceStore
  11. from agent.orchestration import (
  12. AgentRole,
  13. ArtifactRef,
  14. AttemptSubmission,
  15. CriterionResult,
  16. DecisionAction,
  17. OperationStatus,
  18. OrchestrationConfig,
  19. TaskCoordinator,
  20. TaskStatus,
  21. ValidationVerdict,
  22. )
  23. from agent.orchestration.protocols import ValidatorRunResult, WorkerRunResult
  24. from script_build_host.application.mission_factory import ScriptMissionFactory
  25. from script_build_host.application.mission_service import ScriptMissionService
  26. from script_build_host.application.phase_two_planning import (
  27. PhasePolicyGuard,
  28. PhaseTwoPlanningService,
  29. _execution_seconds,
  30. )
  31. from script_build_host.domain.artifacts import DirectionArtifact, DirectionGoal
  32. from script_build_host.domain.errors import ProtocolViolation
  33. from script_build_host.domain.records import BuildStatus, MissionBinding, Principal
  34. from script_build_host.domain.task_contracts import (
  35. PhaseTwoLimits,
  36. ScriptTaskContractV1,
  37. TaskContractError,
  38. )
  39. from script_build_host.infrastructure.task_contract_store import FileScriptTaskContractStore
  40. def test_execution_budget_counts_stage_intervals_not_planner_idle_time() -> None:
  41. now = datetime.now(UTC)
  42. stages = (
  43. SimpleNamespace(
  44. started_at=(now - timedelta(hours=2)).isoformat(),
  45. completed_at=(now - timedelta(hours=2) + timedelta(seconds=12)).isoformat(),
  46. ),
  47. SimpleNamespace(
  48. started_at=(now - timedelta(minutes=1)).isoformat(),
  49. completed_at=(now - timedelta(minutes=1) + timedelta(seconds=8)).isoformat(),
  50. ),
  51. )
  52. assert _execution_seconds(stages) == pytest.approx(20.0)
  53. def _contract(
  54. kind: str,
  55. *,
  56. scope: str = "script-build://scopes/full",
  57. execution_ready: bool = True,
  58. ) -> dict[str, object]:
  59. output = {
  60. "direction": "script-direction/v1",
  61. "candidate-portfolio": "candidate-portfolio/v1",
  62. "compose": "structured-script/v1",
  63. "paragraph": "paragraph-artifact/v1",
  64. "element-set": "element-set-artifact/v1",
  65. "structure": "structure-artifact/v1",
  66. "compare": "comparison-artifact/v1",
  67. "decode-retrieval": "evidence-record/v1",
  68. }[kind]
  69. payload: dict[str, object] = {
  70. "schema_version": "script-task-contract/v1",
  71. "compiler_manifest_digest": "sha256:" + "0" * 64,
  72. "task_kind": kind,
  73. "scope_ref": scope,
  74. "intent_class": (
  75. "portfolio"
  76. if kind == "candidate-portfolio"
  77. else "compose"
  78. if kind == "compose"
  79. else "explore"
  80. ),
  81. "objective": f"produce one bounded {kind} increment",
  82. "input_decision_refs": [],
  83. "base_artifact_ref": None,
  84. "write_scope": [scope],
  85. "gap_ref": None,
  86. "output_schema": output,
  87. "criteria": [
  88. {
  89. "criterion_id": "closed",
  90. "description": "the bounded increment is independently verifiable",
  91. "hard": True,
  92. }
  93. ],
  94. "budget": {
  95. "max_attempts": 4,
  96. "max_tokens": 32000,
  97. "max_seconds": 900,
  98. "max_external_queries": 40,
  99. "max_no_improvement": 3,
  100. },
  101. "goal_ids": ([] if kind in {"direction", "decode-retrieval"} else ["goal-1"]),
  102. "supersedes_decision_ids": [],
  103. "candidate_closure_decision_refs": [],
  104. "adopted_decision_ids": [],
  105. "held_or_rejected_decision_ids": [],
  106. "compose_order": [],
  107. "comparison_decision_refs": [],
  108. }
  109. if execution_ready and kind in {"compose", "candidate-portfolio"}:
  110. ref = {
  111. "decision_id": "decision-1",
  112. "artifact_ref": {
  113. "uri": "script-build://artifact-versions/1",
  114. "kind": "structured_script" if kind == "candidate-portfolio" else "paragraph",
  115. "version": "1",
  116. "digest": "sha256:" + "a" * 64,
  117. },
  118. "scope_ref": scope,
  119. "expected_task_kind": "compose" if kind == "candidate-portfolio" else "paragraph",
  120. }
  121. payload["candidate_closure_decision_refs"] = [ref]
  122. payload["adopted_decision_ids"] = ["decision-1"]
  123. payload["compose_order"] = ["decision-1"]
  124. return payload
  125. def _planner_input(payload: dict[str, object]) -> dict[str, object]:
  126. input_refs = cast(list[dict[str, Any]], payload.get("input_decision_refs") or [])
  127. comparison_refs = cast(list[dict[str, Any]], payload.get("comparison_decision_refs") or [])
  128. base = cast(dict[str, Any] | None, payload.get("base_artifact_ref"))
  129. scope = str(payload["scope_ref"])
  130. prefix = "script-build://scopes/"
  131. path = scope.removeprefix(prefix).split("/") if scope.startswith(prefix) else ["full"]
  132. return {
  133. "task_kind": payload["task_kind"],
  134. "scope_selector": {"anchor": "mission", "path": path[:4]},
  135. "intent_class": payload["intent_class"],
  136. "objective": payload["objective"],
  137. "goal_ids": payload["goal_ids"],
  138. "criteria": [
  139. {
  140. "client_key": item["criterion_id"],
  141. "description": item["description"],
  142. "hard": item.get("hard", True),
  143. }
  144. for item in cast(list[dict[str, Any]], payload["criteria"])
  145. ],
  146. "input_decision_ids": [item["decision_id"] for item in input_refs],
  147. "base_decision_id": next(
  148. (
  149. item["decision_id"]
  150. for item in input_refs
  151. if base is not None and item["artifact_ref"]["uri"] == base["uri"]
  152. ),
  153. None,
  154. ),
  155. "comparison_decision_ids": [item["decision_id"] for item in comparison_refs],
  156. "supersedes_decision_ids": payload["supersedes_decision_ids"],
  157. "gap_key": (
  158. str(payload["gap_ref"]).rsplit("/", 1)[-1]
  159. if payload.get("gap_ref") is not None
  160. else None
  161. ),
  162. }
  163. def test_deep_control_uri_is_not_mistaken_for_base64_content() -> None:
  164. scope = (
  165. "script-build://mission/900032/direction/main/portfolio/main/compose/main/"
  166. "structure/main/paragraph/slot1"
  167. )
  168. contract = ScriptTaskContractV1.from_payload(_contract("paragraph", scope=scope))
  169. assert contract.scope_ref == scope
  170. encoded = _contract("paragraph", scope="script-build://scopes/" + "A" * 80)
  171. with pytest.raises(TaskContractError, match="encoded candidate content"):
  172. ScriptTaskContractV1.from_payload(encoded)
  173. def test_phase_policy_rejects_workspace_contract_without_capability_write_scope() -> None:
  174. guard = PhasePolicyGuard()
  175. ledger = SimpleNamespace(root_task_id="root")
  176. parent = SimpleNamespace(
  177. task_id="portfolio",
  178. current_spec=SimpleNamespace(
  179. context_refs=("script-build://task-kinds/candidate-portfolio",)
  180. ),
  181. )
  182. invalid = ScriptTaskContractV1.from_payload(_contract("structure"))
  183. with pytest.raises(TaskContractError, match="WRITE_SCOPE_VIOLATION"):
  184. guard.contracts(context={"phase": 2}, ledger=ledger, parent=parent, contracts=(invalid,))
  185. payload = _contract("structure")
  186. payload["write_scope"] = ["script-build://writes/paragraphs/main"]
  187. valid = ScriptTaskContractV1.from_payload(payload)
  188. guard.contracts(context={"phase": 2}, ledger=ledger, parent=parent, contracts=(valid,))
  189. def test_phase_policy_requires_broad_container_write_capability() -> None:
  190. guard = PhasePolicyGuard()
  191. ledger = SimpleNamespace(root_task_id="root")
  192. root = SimpleNamespace(task_id="root", current_spec=SimpleNamespace(context_refs=()))
  193. invalid = ScriptTaskContractV1.from_payload(
  194. _contract("candidate-portfolio", execution_ready=False)
  195. )
  196. with pytest.raises(TaskContractError, match="WRITE_SCOPE_VIOLATION"):
  197. guard.contracts(context={"phase": 2}, ledger=ledger, parent=root, contracts=(invalid,))
  198. payload = _contract("candidate-portfolio", execution_ready=False)
  199. payload["write_scope"] = ["script-build://writes"]
  200. valid = ScriptTaskContractV1.from_payload(payload)
  201. guard.contracts(context={"phase": 2}, ledger=ledger, parent=root, contracts=(valid,))
  202. def test_phase_policy_requires_structure_base_and_compose_coverage() -> None:
  203. guard = PhasePolicyGuard()
  204. structure_ref = {
  205. "decision_id": "structure-decision",
  206. "artifact_ref": {
  207. "uri": "script-build://artifact-versions/10",
  208. "kind": "structure",
  209. "version": "10",
  210. "digest": "sha256:" + "a" * 64,
  211. },
  212. "scope_ref": "script-build://scopes/full/structure",
  213. "expected_task_kind": "structure",
  214. }
  215. paragraph_ref = {
  216. "decision_id": "paragraph-decision",
  217. "artifact_ref": {
  218. "uri": "script-build://artifact-versions/11",
  219. "kind": "paragraph",
  220. "version": "11",
  221. "digest": "sha256:" + "b" * 64,
  222. },
  223. "scope_ref": "script-build://scopes/full/structure/paragraph",
  224. "expected_task_kind": "paragraph",
  225. }
  226. paragraph = _contract("paragraph")
  227. paragraph["write_scope"] = ["script-build://writes/paragraphs/main"]
  228. missing_structure_input = ScriptTaskContractV1.from_payload(paragraph)
  229. guard.revisions(context={"phase": 2}, contracts=(missing_structure_input,))
  230. paragraph["base_artifact_ref"] = structure_ref["artifact_ref"]
  231. invalid_paragraph = ScriptTaskContractV1.from_payload(paragraph)
  232. with pytest.raises(TaskContractError, match="accepted Structure decision"):
  233. guard.revisions(context={"phase": 2}, contracts=(invalid_paragraph,))
  234. paragraph["input_decision_refs"] = [structure_ref]
  235. paragraph["base_artifact_ref"] = {
  236. **structure_ref["artifact_ref"],
  237. "digest": "sha256:" + "f" * 64,
  238. }
  239. reconstructed_paragraph = ScriptTaskContractV1.from_payload(paragraph)
  240. with pytest.raises(TaskContractError, match="unchanged"):
  241. guard.revisions(context={"phase": 2}, contracts=(reconstructed_paragraph,))
  242. paragraph["base_artifact_ref"] = structure_ref["artifact_ref"]
  243. valid_paragraph = ScriptTaskContractV1.from_payload(paragraph)
  244. guard.revisions(context={"phase": 2}, contracts=(valid_paragraph,))
  245. compose = _contract("compose", execution_ready=False)
  246. compose["write_scope"] = ["script-build://writes"]
  247. guard.revisions(
  248. context={"phase": 2},
  249. contracts=(ScriptTaskContractV1.from_payload(compose),),
  250. )
  251. compose["candidate_closure_decision_refs"] = [paragraph_ref]
  252. compose["adopted_decision_ids"] = ["paragraph-decision"]
  253. compose["compose_order"] = ["paragraph-decision"]
  254. with pytest.raises(TaskContractError, match="accepted Direction"):
  255. guard.revisions(
  256. context={"phase": 2},
  257. contracts=(ScriptTaskContractV1.from_payload(compose),),
  258. )
  259. direction_ref = {
  260. "decision_id": "direction-decision",
  261. "artifact_ref": {
  262. "uri": "script-build://artifact-versions/9",
  263. "kind": "direction",
  264. "version": "9",
  265. "digest": "sha256:" + "d" * 64,
  266. },
  267. "scope_ref": "script-build://scopes/full",
  268. "expected_task_kind": "direction",
  269. }
  270. compose["input_decision_refs"] = [direction_ref]
  271. invalid_compose = ScriptTaskContractV1.from_payload(compose)
  272. with pytest.raises(TaskContractError, match="covering Structure"):
  273. guard.revisions(context={"phase": 2}, contracts=(invalid_compose,))
  274. compose["candidate_closure_decision_refs"] = [structure_ref, paragraph_ref]
  275. compose["adopted_decision_ids"] = ["structure-decision", "paragraph-decision"]
  276. compose["compose_order"] = ["structure-decision", "paragraph-decision"]
  277. valid_compose = ScriptTaskContractV1.from_payload(compose)
  278. guard.revisions(context={"phase": 2}, contracts=(valid_compose,))
  279. class _Bindings:
  280. def __init__(self, root: str) -> None:
  281. self.binding = MissionBinding(
  282. binding_id=1,
  283. script_build_id=7,
  284. root_trace_id=root,
  285. input_snapshot_id=11,
  286. active_direction_artifact_version_id=None,
  287. accepted_root_artifact_version_id=None,
  288. engine_version="test",
  289. schema_version="v1",
  290. created_at=datetime.now(UTC),
  291. updated_at=datetime.now(UTC),
  292. )
  293. async def get_by_root(self, root: str) -> MissionBinding:
  294. assert root == self.binding.root_trace_id
  295. return self.binding
  296. async def get_by_build(self, script_build_id: int) -> MissionBinding:
  297. assert script_build_id == self.binding.script_build_id
  298. return self.binding
  299. class _DirectionArtifacts:
  300. async def get_by_id(self, identifier: int, *, script_build_id: int) -> Any:
  301. assert script_build_id == 7
  302. return SimpleNamespace(
  303. artifact_version_id=identifier,
  304. artifact=DirectionArtifact(
  305. goals=(
  306. DirectionGoal(
  307. goal_id="goal-1",
  308. statement="complete the script",
  309. rationale="the mission needs one complete result",
  310. success_criteria=("the result is complete",),
  311. ),
  312. ),
  313. evidence_refs=("script-build://artifact-versions/evidence",),
  314. ),
  315. )
  316. class _PlaceholderThenReplacementExecutor:
  317. """Exercise Coordinator state transitions without bypassing durable operations."""
  318. def __init__(self, coordinator: TaskCoordinator) -> None:
  319. self.coordinator = coordinator
  320. self.artifact_version = 100
  321. async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
  322. refs = cast(dict[str, Any], context["task_spec"])["context_refs"]
  323. kind = next(item.rsplit("/", 1)[-1] for item in refs if "/task-kinds/" in item)
  324. artifact_kind = {
  325. "direction": "direction",
  326. "compose": "structured_script",
  327. "structure": "structure",
  328. "paragraph": "paragraph",
  329. "element-set": "element_set",
  330. }[kind]
  331. digest = "sha256:" + sha256(context["attempt_id"].encode()).hexdigest()
  332. version = str(self.artifact_version)
  333. self.artifact_version += 1
  334. ref = ArtifactRef(
  335. f"script-build://artifact-versions/{version}",
  336. artifact_kind,
  337. version,
  338. digest,
  339. )
  340. await self.coordinator.submit_attempt(
  341. {
  342. "role": AgentRole.WORKER.value,
  343. "root_trace_id": context["root_trace_id"],
  344. "task_id": context["task_id"],
  345. "attempt_id": context["attempt_id"],
  346. "spec_version": context["spec_version"],
  347. "trace_id": context["worker_trace_id"],
  348. "tool_call_id": f"submit:{context['attempt_id']}",
  349. "operation_id": context["operation_id"],
  350. "execution_epoch": context["execution_epoch"],
  351. },
  352. AttemptSubmission(
  353. summary=f"kind={artifact_kind};status=frozen",
  354. artifact_refs=[ref],
  355. ),
  356. )
  357. return WorkerRunResult(context["worker_trace_id"], "completed")
  358. async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
  359. refs = cast(dict[str, Any], context["task_spec"])["context_refs"]
  360. kind = next(item.rsplit("/", 1)[-1] for item in refs if "/task-kinds/" in item)
  361. placeholder = kind == "compose" and context["spec_version"] == 2
  362. verdict = ValidationVerdict.FAILED if placeholder else ValidationVerdict.PASSED
  363. criteria = cast(dict[str, Any], context["task_spec"])["acceptance_criteria"]
  364. await self.coordinator.submit_validation(
  365. {
  366. "role": AgentRole.VALIDATOR.value,
  367. "root_trace_id": context["root_trace_id"],
  368. "task_id": context["task_id"],
  369. "attempt_id": context["attempt_id"],
  370. "validation_id": context["validation_id"],
  371. "snapshot_id": context["snapshot_id"],
  372. "trace_id": context["validator_trace_id"],
  373. "tool_call_id": f"validate:{context['validation_id']}",
  374. "operation_id": context["operation_id"],
  375. "execution_epoch": context["execution_epoch"],
  376. },
  377. verdict,
  378. [
  379. CriterionResult(
  380. criterion_id=criterion["criterion_id"],
  381. verdict=verdict,
  382. reason=(
  383. "REALIZATION_PLACEHOLDER at artifact.paragraphs[0].description"
  384. if placeholder
  385. else "the replacement is concretely rendered"
  386. ),
  387. )
  388. for criterion in criteria
  389. ],
  390. (
  391. "hard_defects=1;defect_codes=REALIZATION_PLACEHOLDER"
  392. if placeholder
  393. else "hard_defects=0"
  394. ),
  395. (),
  396. ("artifact.paragraphs[0].description",) if placeholder else (),
  397. ("REALIZATION_PLACEHOLDER",) if placeholder else (),
  398. "split" if placeholder else "accept",
  399. )
  400. return ValidatorRunResult(context["validator_trace_id"], "completed")
  401. async def stop(self, trace_id: str) -> bool:
  402. del trace_id
  403. return True
  404. async def _service(
  405. tmp_path: Path, *, active_direction: bool = True
  406. ) -> tuple[PhaseTwoPlanningService, TaskCoordinator, str]:
  407. root = "root-contracts"
  408. coordinator = TaskCoordinator(
  409. FileSystemTaskStore(str(tmp_path / "ledger")),
  410. FileSystemArtifactStore(str(tmp_path / "snapshots")),
  411. FileSystemTraceStore(str(tmp_path / "traces")),
  412. )
  413. await coordinator.ensure_ledger(
  414. root,
  415. {
  416. "objective": "build a script",
  417. "acceptance_criteria": [
  418. {"criterion_id": "ready", "description": "ready", "hard": True}
  419. ],
  420. "context_refs": ["script-build://inputs/11"],
  421. },
  422. )
  423. bindings = _Bindings(root)
  424. service = PhaseTwoPlanningService(
  425. coordinator=coordinator,
  426. bindings=bindings,
  427. contracts=FileScriptTaskContractStore(tmp_path / "data"),
  428. artifacts=cast(Any, _DirectionArtifacts()),
  429. )
  430. if active_direction:
  431. await _activate_direction(service, coordinator, bindings, root)
  432. return service, coordinator, root
  433. async def _activate_direction(
  434. service: PhaseTwoPlanningService,
  435. coordinator: TaskCoordinator,
  436. bindings: _Bindings,
  437. root: str,
  438. ) -> None:
  439. coordinator.set_executor(_PlaceholderThenReplacementExecutor(coordinator))
  440. planned = await service.plan_script_tasks(
  441. contract_payloads=[_planner_input(_contract("direction"))],
  442. parent_task_id=None,
  443. context={"root_trace_id": root, "tool_call_id": "fixture:direction"},
  444. )
  445. direction_id = cast(str, planned["task_ids"][0])
  446. cycle = (
  447. await service.dispatch_script_tasks(
  448. task_ids=(direction_id,),
  449. context={"root_trace_id": root, "tool_call_id": "fixture:dispatch-direction"},
  450. )
  451. )[0]
  452. accepted = await service.decide_script_task(
  453. task_id=direction_id,
  454. action=DecisionAction.ACCEPT.value,
  455. reason="fixture Direction passed",
  456. validation_id=None,
  457. replacement_contract=None,
  458. child_contracts=(),
  459. context={"root_trace_id": root, "tool_call_id": "fixture:accept-direction"},
  460. )
  461. ledger = await coordinator.task_store.load(root)
  462. decision = ledger.decisions[cast(str, accepted["decision_id"])]
  463. assert decision.validation_id == cycle["validation_id"]
  464. attempt = ledger.attempts[cast(str, decision.attempt_id)]
  465. assert attempt.submission is not None
  466. direction_ref = attempt.submission.artifact_refs[0]
  467. bindings.binding = replace(
  468. bindings.binding,
  469. active_direction_artifact_version_id=int(direction_ref.version),
  470. )
  471. @pytest.mark.asyncio
  472. async def test_plan_inspection_returns_active_frontier_and_compact_accept_index(
  473. tmp_path: Path,
  474. ) -> None:
  475. service, _, root = await _service(tmp_path)
  476. view = await service.inspect_script_plan(context={"root_trace_id": root})
  477. assert "tasks" not in view
  478. assert view["root"]["status"] == TaskStatus.NEEDS_REPLAN.value
  479. assert [item["task_kind"] for item in view["active_tasks"]] == ["root"]
  480. assert len(view["accepted_decisions"]) == 1
  481. accepted = view["accepted_decisions"][0]
  482. assert accepted["task_kind"] == "direction"
  483. assert accepted["artifact_kind"] == "direction"
  484. assert "digest" not in accepted
  485. assert view["counts_by_status"] == {
  486. TaskStatus.NEEDS_REPLAN.value: 1,
  487. TaskStatus.COMPLETED.value: 1,
  488. }
  489. @pytest.mark.asyncio
  490. async def test_dispatch_explains_that_blocked_children_do_not_close_parent(
  491. tmp_path: Path,
  492. ) -> None:
  493. service, coordinator, root = await _service(tmp_path, active_direction=False)
  494. direction = await service.plan_script_tasks(
  495. contract_payloads=[_planner_input(_contract("direction"))],
  496. parent_task_id=None,
  497. context={"root_trace_id": root, "tool_call_id": "direction"},
  498. )
  499. direction_id = direction["task_ids"][0]
  500. retrieval = await service.plan_script_tasks(
  501. contract_payloads=[_planner_input(_contract("decode-retrieval"))],
  502. parent_task_id=direction_id,
  503. context={"root_trace_id": root, "tool_call_id": "retrieval"},
  504. )
  505. child_id = retrieval["task_ids"][0]
  506. await coordinator.decide_task(
  507. root,
  508. child_id,
  509. None,
  510. DecisionAction.BLOCK,
  511. {"reason": "unusable evidence"},
  512. "block-child",
  513. )
  514. with pytest.raises(TaskContractError, match="BLOCKED is resumable"):
  515. await service.dispatch_script_tasks(
  516. task_ids=[direction_id],
  517. context={"root_trace_id": root, "tool_call_id": "dispatch-parent"},
  518. )
  519. @pytest.mark.asyncio
  520. async def test_phase_one_direction_requires_an_accepted_retrieval_before_dispatch(
  521. tmp_path: Path,
  522. ) -> None:
  523. service, _coordinator, root = await _service(tmp_path, active_direction=False)
  524. context = {"root_trace_id": root, "phase": 1}
  525. direction = await service.plan_script_tasks(
  526. contract_payloads=[_planner_input(_contract("direction"))],
  527. parent_task_id=None,
  528. context={**context, "tool_call_id": "direction"},
  529. )
  530. await service.plan_script_tasks(
  531. contract_payloads=[_planner_input(_contract("decode-retrieval"))],
  532. parent_task_id=direction["task_ids"][0],
  533. context={**context, "tool_call_id": "retrieval"},
  534. )
  535. with pytest.raises(TaskContractError, match="Retrieval ACCEPT"):
  536. await service.dispatch_script_tasks(
  537. task_ids=[direction["task_ids"][0]],
  538. context={**context, "tool_call_id": "dispatch-direction"},
  539. )
  540. @pytest.mark.asyncio
  541. async def test_unique_portfolio_container_cannot_be_cancelled(tmp_path: Path) -> None:
  542. service, _coordinator, root = await _service(tmp_path)
  543. portfolio = await service.plan_script_tasks(
  544. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  545. parent_task_id=None,
  546. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  547. )
  548. with pytest.raises(TaskContractError, match="cannot be cancelled"):
  549. await service.decide_script_task(
  550. task_id=portfolio["task_ids"][0],
  551. action="cancel",
  552. reason="wrong structural recovery",
  553. validation_id=None,
  554. replacement_contract=None,
  555. child_contracts=(),
  556. context={"root_trace_id": root, "tool_call_id": "cancel-portfolio"},
  557. )
  558. @pytest.mark.asyncio
  559. async def test_accepted_input_scope_is_rejected_before_worker_dispatch(tmp_path: Path) -> None:
  560. service, _coordinator, root = await _service(tmp_path)
  561. structure_scope = "script-build://direction/main/compose/candidate-1/structure/main"
  562. structure = ScriptTaskContractV1.from_payload(_contract("structure", scope=structure_scope))
  563. paragraph_payload = _contract(
  564. "paragraph",
  565. scope="script-build://direction/main/compose/candidate-1/paragraph/p1",
  566. )
  567. accepted_ref = {
  568. "decision_id": "accepted-structure",
  569. "artifact_ref": {
  570. "uri": "script-build://artifact-versions/1",
  571. "kind": "structure",
  572. "version": "1",
  573. "digest": "sha256:" + "a" * 64,
  574. },
  575. "scope_ref": structure_scope,
  576. "expected_task_kind": "structure",
  577. }
  578. paragraph_payload["input_decision_refs"] = [accepted_ref]
  579. paragraph_payload["base_artifact_ref"] = accepted_ref["artifact_ref"]
  580. paragraph = ScriptTaskContractV1.from_payload(paragraph_payload)
  581. producer_task = SimpleNamespace(task_id="structure-task")
  582. ledger = SimpleNamespace(
  583. decisions={
  584. "accepted-structure": SimpleNamespace(
  585. action=DecisionAction.ACCEPT,
  586. attempt_id="structure-attempt",
  587. task_id=producer_task.task_id,
  588. )
  589. },
  590. tasks={producer_task.task_id: producer_task},
  591. )
  592. async def frozen_contract(_root_trace_id: str, task: Any) -> Any:
  593. assert task is producer_task
  594. return SimpleNamespace(contract=structure)
  595. service.contract_for_task = frozen_contract # type: ignore[method-assign]
  596. with pytest.raises(TaskContractError, match="INPUT_SCOPE_MISMATCH") as raised:
  597. await service._guard_accepted_input_scopes(root, ledger, (paragraph,))
  598. assert raised.value.details == {
  599. "decision_id": "accepted-structure",
  600. "consumer_kind": "paragraph",
  601. "consumer_scope": paragraph.scope_ref,
  602. "producer_kind": "structure",
  603. "producer_scope": structure.scope_ref,
  604. "source": "explicit",
  605. }
  606. class _BlockingExecutor:
  607. def __init__(self) -> None:
  608. self.started = asyncio.Event()
  609. self.stopped_trace_ids: list[str] = []
  610. async def run_worker(self, context: dict[str, Any]) -> WorkerRunResult:
  611. self.started.set()
  612. await asyncio.Event().wait()
  613. raise AssertionError(f"worker unexpectedly resumed: {context['attempt_id']}")
  614. async def run_validator(self, context: dict[str, Any]) -> ValidatorRunResult:
  615. raise AssertionError(f"validator must not start: {context['validation_id']}")
  616. async def stop(self, trace_id: str) -> bool:
  617. self.stopped_trace_ids.append(trace_id)
  618. return True
  619. class _StopState:
  620. def __init__(self) -> None:
  621. self.status = BuildStatus.RUNNING
  622. self.error_summary: str | None = None
  623. async def get_status(self, script_build_id: int) -> BuildStatus:
  624. assert script_build_id == 7
  625. return self.status
  626. async def set_status(
  627. self,
  628. script_build_id: int,
  629. status: BuildStatus,
  630. *,
  631. error_summary: str | None = None,
  632. ) -> None:
  633. assert script_build_id == 7
  634. self.status = status
  635. self.error_summary = error_summary
  636. @pytest.mark.asyncio
  637. async def test_real_operation_stop_converges_and_forbids_new_dispatch(tmp_path: Path) -> None:
  638. root = "root-stop-operation"
  639. trace_store = FileSystemTraceStore(str(tmp_path / "stop-traces"))
  640. coordinator = TaskCoordinator(
  641. FileSystemTaskStore(str(tmp_path / "stop-ledger")),
  642. FileSystemArtifactStore(str(tmp_path / "stop-artifacts")),
  643. trace_store,
  644. config=OrchestrationConfig(stop_grace_seconds=0.01),
  645. )
  646. await coordinator.ensure_ledger(
  647. root,
  648. {
  649. "objective": "build a script",
  650. "acceptance_criteria": [
  651. {"criterion_id": "ready", "description": "ready", "hard": True}
  652. ],
  653. "context_refs": ["script-build://inputs/11"],
  654. },
  655. )
  656. bindings = _Bindings(root)
  657. planning = PhaseTwoPlanningService(
  658. coordinator=coordinator,
  659. bindings=bindings,
  660. contracts=FileScriptTaskContractStore(tmp_path / "stop-contracts"),
  661. artifacts=cast(Any, _DirectionArtifacts()),
  662. )
  663. await _activate_direction(planning, coordinator, bindings, root)
  664. portfolio = await planning.plan_script_tasks(
  665. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  666. parent_task_id=None,
  667. context={"root_trace_id": root, "tool_call_id": "stop:portfolio"},
  668. )
  669. compose = await planning.plan_script_tasks(
  670. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  671. parent_task_id=portfolio["task_ids"][0],
  672. context={"root_trace_id": root, "tool_call_id": "stop:compose"},
  673. )
  674. paragraph = await planning.plan_script_tasks(
  675. contract_payloads=[_planner_input(_contract("paragraph"))],
  676. parent_task_id=compose["task_ids"][0],
  677. context={"root_trace_id": root, "tool_call_id": "stop:paragraph"},
  678. )
  679. executor = _BlockingExecutor()
  680. coordinator.set_executor(executor)
  681. async def unused_llm(**_kwargs: Any) -> Any:
  682. raise AssertionError("Planner LLM must not run during operation stop")
  683. runner = AgentRunner(trace_store=trace_store, llm_call=unused_llm)
  684. state = _StopState()
  685. class Authorizer:
  686. async def require_access(self, principal: Principal, script_build_id: int) -> None:
  687. assert principal.subject == "owner" and script_build_id == 7
  688. mission = ScriptMissionService(
  689. runner=runner,
  690. coordinator=coordinator,
  691. factory=ScriptMissionFactory(),
  692. input_snapshots=SimpleNamespace(),
  693. bindings=bindings,
  694. legacy_state=state,
  695. authorizer=Authorizer(),
  696. direction_reconciler=SimpleNamespace(),
  697. stop_timeout_seconds=1,
  698. )
  699. planning.dispatch_guard = mission
  700. dispatch = asyncio.create_task(
  701. planning.dispatch_script_tasks(
  702. task_ids=paragraph["task_ids"],
  703. context={"root_trace_id": root, "tool_call_id": "stop:dispatch"},
  704. )
  705. )
  706. await asyncio.wait_for(executor.started.wait(), timeout=1)
  707. stopped = await mission.stop(7, Principal("owner"))
  708. dispatch_result = await asyncio.gather(dispatch, return_exceptions=True)
  709. assert stopped.status is BuildStatus.STOPPED
  710. assert len(stopped.stopped_operation_ids) == 1
  711. assert executor.stopped_trace_ids
  712. assert not isinstance(dispatch_result[0], asyncio.CancelledError)
  713. ledger = await coordinator.task_store.load(root)
  714. operation = ledger.operations[stopped.stopped_operation_ids[0]]
  715. assert operation.status is OperationStatus.STOPPED
  716. assert ledger.tasks[paragraph["task_ids"][0]].status is TaskStatus.NEEDS_REPLAN
  717. assert (await mission.stop(7, Principal("owner"))).status is BuildStatus.STOPPED
  718. with pytest.raises(ProtocolViolation, match="forbidden after stop intent"):
  719. await planning.dispatch_script_tasks(
  720. task_ids=paragraph["task_ids"],
  721. context={"root_trace_id": root, "tool_call_id": "stop:late-dispatch"},
  722. )
  723. @pytest.mark.asyncio
  724. async def test_contract_store_is_content_addressed_and_detects_tampering(tmp_path: Path) -> None:
  725. store = FileScriptTaskContractStore(tmp_path)
  726. from script_build_host.domain.task_contracts import ScriptTaskContractV1
  727. contract = ScriptTaskContractV1.from_payload(_contract("paragraph"))
  728. first = await store.freeze("root-a", contract)
  729. second = await store.freeze("root-a", contract)
  730. assert first == second
  731. path = tmp_path / "script-task-contracts" / "root-a" / "sha256" / first.uri.rsplit("/", 1)[-1]
  732. path.write_bytes(b"{}")
  733. with pytest.raises(TaskContractError, match="TASK_CONTRACT_DIGEST_MISMATCH"):
  734. await store.read("root-a", first.uri)
  735. assert not await store.verify("root-a", first.uri, first.digest)
  736. @pytest.mark.asyncio
  737. async def test_orphan_contract_reclamation_uses_only_ledger_references(tmp_path: Path) -> None:
  738. service, _, root = await _service(tmp_path)
  739. referenced = await service.plan_script_tasks(
  740. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  741. parent_task_id=None,
  742. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  743. )
  744. orphan = await service.contracts.freeze(
  745. root, ScriptTaskContractV1.from_payload(_contract("paragraph"))
  746. )
  747. removed = await service.reclaim_orphan_contracts(root_trace_id=root)
  748. assert removed == (orphan.uri,)
  749. assert await service.contracts.verify(
  750. root,
  751. referenced["contracts"][0]["contract_ref"],
  752. referenced["contracts"][0]["contract_digest"],
  753. )
  754. assert not await service.contracts.verify(root, orphan.uri, orphan.digest)
  755. @pytest.mark.asyncio
  756. async def test_failed_operation_discards_workspace_but_recovery_required_preserves_it(
  757. tmp_path: Path,
  758. ) -> None:
  759. service, _, root = await _service(tmp_path)
  760. class Lifecycle:
  761. def __init__(self) -> None:
  762. self.calls: list[dict[str, str]] = []
  763. async def discard_attempt(self, **values: str) -> None:
  764. self.calls.append(values)
  765. lifecycle = Lifecycle()
  766. service.workspace_lifecycle = lifecycle
  767. attempts = {
  768. "failed-attempt": SimpleNamespace(attempt_id="failed-attempt", task_id="paragraph-task"),
  769. "recovery-attempt": SimpleNamespace(
  770. attempt_id="recovery-attempt", task_id="paragraph-task"
  771. ),
  772. }
  773. ledger = SimpleNamespace(
  774. attempts=attempts,
  775. operations={
  776. "failed": SimpleNamespace(
  777. status=OperationStatus.FAILED,
  778. error="worker failed before submit",
  779. attempt_ids=["failed-attempt"],
  780. ),
  781. "recovery": SimpleNamespace(
  782. status=OperationStatus.FAILED,
  783. error="MISSION_RECOVERY_REQUIRED",
  784. attempt_ids=["recovery-attempt"],
  785. ),
  786. },
  787. )
  788. await service._reconcile_terminal_workspaces(root, ledger)
  789. assert lifecycle.calls == [
  790. {
  791. "root_trace_id": root,
  792. "task_id": "paragraph-task",
  793. "attempt_id": "failed-attempt",
  794. "reason": "worker failed before submit",
  795. }
  796. ]
  797. def test_comparison_refs_are_kind_specific_and_limits_cannot_be_raised() -> None:
  798. candidate_refs = []
  799. for index, kind in ((1, "paragraph"), (2, "structure")):
  800. candidate_refs.append(
  801. {
  802. "decision_id": f"decision-{index}",
  803. "artifact_ref": {
  804. "uri": f"script-build://artifact-versions/{index}",
  805. "kind": kind,
  806. "version": str(index),
  807. "digest": "sha256:" + str(index) * 64,
  808. },
  809. "scope_ref": "script-build://scopes/full",
  810. "expected_task_kind": kind,
  811. }
  812. )
  813. compare = _contract("compare")
  814. compare["comparison_decision_refs"] = candidate_refs
  815. assert ScriptTaskContractV1.from_payload(compare).task_kind.value == "compare"
  816. comparison_ref = {
  817. "decision_id": "decision-compare",
  818. "artifact_ref": {
  819. "uri": "script-build://artifact-versions/3",
  820. "kind": "comparison",
  821. "version": "3",
  822. "digest": "sha256:" + "3" * 64,
  823. },
  824. "scope_ref": "script-build://scopes/full",
  825. "expected_task_kind": "compare",
  826. }
  827. compose = _contract("compose")
  828. compose["comparison_decision_refs"] = [comparison_ref]
  829. ScriptTaskContractV1.from_payload(compose)
  830. portfolio = _contract("candidate-portfolio")
  831. portfolio["comparison_decision_refs"] = [comparison_ref]
  832. with pytest.raises(TaskContractError, match="comparison_decision_refs"):
  833. ScriptTaskContractV1.from_payload(portfolio)
  834. with pytest.raises(TaskContractError, match="TASK_BUDGET_EXCEEDED"):
  835. PhaseTwoLimits(max_tasks=65)
  836. @pytest.mark.asyncio
  837. async def test_planning_freezes_contract_and_creates_only_one_portfolio(tmp_path: Path) -> None:
  838. service, coordinator, root = await _service(tmp_path)
  839. result = await service.plan_script_tasks(
  840. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  841. parent_task_id=None,
  842. context={"root_trace_id": root, "tool_call_id": "plan-portfolio"},
  843. )
  844. task_id = result["task_ids"][0]
  845. ledger = await coordinator.task_store.load(root)
  846. task = ledger.tasks[task_id]
  847. assert task.current_spec.context_refs[0].endswith("candidate-portfolio")
  848. assert task.current_spec.context_refs[1].startswith("script-build://task-contracts/sha256/")
  849. assert task.current_spec.context_refs[2] == "script-build://inputs/11"
  850. frozen = await service.contract_for_task(root, task)
  851. assert frozen.contract.scope_ref == "script-build://scopes/full"
  852. with pytest.raises(TaskContractError, match="one Portfolio"):
  853. await service.plan_script_tasks(
  854. contract_payloads=[
  855. _planner_input(_contract("candidate-portfolio", execution_ready=False))
  856. ],
  857. parent_task_id=None,
  858. context={"root_trace_id": root, "tool_call_id": "plan-portfolio-2"},
  859. )
  860. @pytest.mark.asyncio
  861. async def test_protected_phase_rejects_cross_phase_root_contracts(tmp_path: Path) -> None:
  862. service, coordinator, root = await _service(tmp_path)
  863. before = await coordinator.task_store.load(root)
  864. existing_task_ids = set(before.tasks)
  865. with pytest.raises(TaskContractError, match="PHASE_POLICY_VIOLATION"):
  866. await service.plan_script_tasks(
  867. contract_payloads=[
  868. _planner_input(_contract("candidate-portfolio", execution_ready=False))
  869. ],
  870. parent_task_id=None,
  871. context={"root_trace_id": root, "phase": 1},
  872. )
  873. ledger = await coordinator.task_store.load(root)
  874. assert set(ledger.tasks) == existing_task_ids
  875. @pytest.mark.asyncio
  876. async def test_phase_two_root_violation_describes_attempted_and_required_kind(
  877. tmp_path: Path,
  878. ) -> None:
  879. service, _, root = await _service(tmp_path)
  880. with pytest.raises(TaskContractError) as caught:
  881. await service.plan_script_tasks(
  882. contract_payloads=[_planner_input(_contract("structure"))],
  883. parent_task_id=None,
  884. context={"root_trace_id": root, "phase": 2},
  885. )
  886. error = caught.value
  887. assert "must first plan exactly one candidate-portfolio" in error.summary
  888. assert error.details["attempted_task_kinds"] == ["structure"]
  889. assert error.details["allowed_task_kinds"] == ["candidate-portfolio"]
  890. @pytest.mark.asyncio
  891. async def test_invalid_contract_batch_creates_zero_tasks(tmp_path: Path) -> None:
  892. service, coordinator, root = await _service(tmp_path)
  893. before = await coordinator.task_store.load(root)
  894. existing_task_ids = set(before.tasks)
  895. invalid = _planner_input(_contract("paragraph"))
  896. invalid["output_schema"] = "structured-script/v1"
  897. with pytest.raises(TaskContractError, match="Host-owned fields"):
  898. await service.plan_script_tasks(
  899. contract_payloads=[
  900. _planner_input(_contract("candidate-portfolio")),
  901. invalid,
  902. ],
  903. parent_task_id=None,
  904. context={"root_trace_id": root},
  905. )
  906. ledger = await coordinator.task_store.load(root)
  907. assert set(ledger.tasks) == existing_task_ids
  908. @pytest.mark.asyncio
  909. async def test_planning_rejects_unresolved_supersession_before_freeze(tmp_path: Path) -> None:
  910. service, coordinator, root = await _service(tmp_path)
  911. portfolio = await service.plan_script_tasks(
  912. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  913. parent_task_id=None,
  914. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  915. )
  916. compose = await service.plan_script_tasks(
  917. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  918. parent_task_id=portfolio["task_ids"][0],
  919. context={"root_trace_id": root, "tool_call_id": "compose"},
  920. )
  921. replacement = _contract("paragraph")
  922. replacement["supersedes_decision_ids"] = ["missing-accept"]
  923. before = await coordinator.task_store.load(root)
  924. existing_task_ids = set(before.tasks)
  925. with pytest.raises(TaskContractError, match="not an ACCEPT decision"):
  926. await service.plan_script_tasks(
  927. contract_payloads=[_planner_input(replacement)],
  928. parent_task_id=compose["task_ids"][0],
  929. context={"root_trace_id": root, "tool_call_id": "invalid-replacement"},
  930. )
  931. ledger = await coordinator.task_store.load(root)
  932. assert set(ledger.tasks) == existing_task_ids
  933. @pytest.mark.asyncio
  934. async def test_open_compose_container_cannot_dispatch(tmp_path: Path) -> None:
  935. service, coordinator, root = await _service(tmp_path)
  936. portfolio = await service.plan_script_tasks(
  937. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  938. parent_task_id=None,
  939. context={"root_trace_id": root, "tool_call_id": "p"},
  940. )
  941. compose = await service.plan_script_tasks(
  942. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  943. parent_task_id=portfolio["task_ids"][0],
  944. context={"root_trace_id": root, "tool_call_id": "c"},
  945. )
  946. before = await coordinator.task_store.load(root)
  947. operation_ids = set(before.operations)
  948. with pytest.raises(TaskContractError, match="first accept child candidates"):
  949. await service.dispatch_script_tasks(
  950. task_ids=compose["task_ids"], context={"root_trace_id": root}
  951. )
  952. ledger = await coordinator.task_store.load(root)
  953. assert set(ledger.operations) == operation_ids
  954. @pytest.mark.asyncio
  955. async def test_compose_placeholder_split_replacement_then_revised_contract_passes(
  956. tmp_path: Path,
  957. ) -> None:
  958. service, coordinator, root = await _service(tmp_path)
  959. coordinator.set_executor(_PlaceholderThenReplacementExecutor(coordinator))
  960. portfolio = await service.plan_script_tasks(
  961. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  962. parent_task_id=None,
  963. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  964. )
  965. compose = await service.plan_script_tasks(
  966. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  967. parent_task_id=portfolio["task_ids"][0],
  968. context={"root_trace_id": root, "tool_call_id": "compose-v1"},
  969. )
  970. compose_id = compose["task_ids"][0]
  971. accepted_children: dict[str, str] = {}
  972. for kind in ("structure", "paragraph", "element-set"):
  973. planned = await service.plan_script_tasks(
  974. contract_payloads=[_planner_input(_contract(kind))],
  975. parent_task_id=compose_id,
  976. context={"root_trace_id": root, "tool_call_id": f"plan:{kind}"},
  977. )
  978. child_id = cast(str, planned["task_ids"][0])
  979. child_cycle = (
  980. await service.dispatch_script_tasks(
  981. task_ids=(child_id,),
  982. context={"root_trace_id": root, "tool_call_id": f"dispatch:{kind}"},
  983. )
  984. )[0]
  985. accepted = await service.decide_script_task(
  986. task_id=child_id,
  987. action=DecisionAction.ACCEPT.value,
  988. reason=f"{kind} passed",
  989. validation_id=cast(str, child_cycle["validation_id"]),
  990. replacement_contract=None,
  991. child_contracts=(),
  992. context={"root_trace_id": root, "tool_call_id": f"accept:{kind}"},
  993. )
  994. accepted_children[kind] = cast(str, accepted["decision_id"])
  995. await service.decide_script_task(
  996. task_id=compose_id,
  997. action=DecisionAction.REVISE.value,
  998. reason="freeze the initial complete frontier",
  999. validation_id=None,
  1000. replacement_contract=None,
  1001. child_contracts=(),
  1002. selected_decision_ids=tuple(accepted_children.values()),
  1003. context={"root_trace_id": root, "tool_call_id": "revise-compose-v1"},
  1004. )
  1005. first_cycle = (
  1006. await service.dispatch_script_tasks(
  1007. task_ids=(compose_id,),
  1008. context={"root_trace_id": root, "tool_call_id": "dispatch-compose-v1"},
  1009. )
  1010. )[0]
  1011. assert first_cycle["error"] is None
  1012. ledger = await coordinator.task_store.load(root)
  1013. first_validation = ledger.validations[first_cycle["validation_id"]]
  1014. assert first_validation.verdict is ValidationVerdict.FAILED
  1015. assert first_validation.unverified_claims == ["artifact.paragraphs[0].description"]
  1016. assert first_validation.risks == ["REALIZATION_PLACEHOLDER"]
  1017. split = await service.decide_script_task(
  1018. task_id=compose_id,
  1019. action=DecisionAction.SPLIT.value,
  1020. reason="replace only the unrealized opening scope",
  1021. validation_id=first_validation.validation_id,
  1022. replacement_contract=None,
  1023. child_contracts=(
  1024. _planner_input(
  1025. _contract(
  1026. "paragraph",
  1027. scope="script-build://scopes/full/replacement",
  1028. )
  1029. ),
  1030. ),
  1031. context={"root_trace_id": root, "tool_call_id": "split-placeholder"},
  1032. )
  1033. replacement_id = split["payload"]["child_task_ids"][0]
  1034. replacement_cycle = (
  1035. await service.dispatch_script_tasks(
  1036. task_ids=(replacement_id,),
  1037. context={"root_trace_id": root, "tool_call_id": "dispatch-replacement"},
  1038. )
  1039. )[0]
  1040. replacement_accept = await service.decide_script_task(
  1041. task_id=replacement_id,
  1042. action=DecisionAction.ACCEPT.value,
  1043. reason="the local replacement passed independent validation",
  1044. validation_id=replacement_cycle["validation_id"],
  1045. replacement_contract=None,
  1046. child_contracts=(),
  1047. context={"root_trace_id": root, "tool_call_id": "accept-replacement"},
  1048. )
  1049. ledger = await coordinator.task_store.load(root)
  1050. assert ledger.tasks[compose_id].status is TaskStatus.NEEDS_REPLAN
  1051. replacement_decision_id = replacement_accept["decision_id"]
  1052. await service.decide_script_task(
  1053. task_id=compose_id,
  1054. action=DecisionAction.REVISE.value,
  1055. reason="pin the accepted replacement and retry the complete render",
  1056. validation_id=None,
  1057. replacement_contract=None,
  1058. child_contracts=(),
  1059. selected_decision_ids=(
  1060. accepted_children["structure"],
  1061. replacement_decision_id,
  1062. accepted_children["element-set"],
  1063. ),
  1064. context={"root_trace_id": root, "tool_call_id": "revise-compose-v2"},
  1065. )
  1066. second_cycle = (
  1067. await service.dispatch_script_tasks(
  1068. task_ids=(compose_id,),
  1069. context={"root_trace_id": root, "tool_call_id": "dispatch-compose-v2"},
  1070. )
  1071. )[0]
  1072. compose_accept = await service.decide_script_task(
  1073. task_id=compose_id,
  1074. action=DecisionAction.ACCEPT.value,
  1075. reason="the revised full render passed independent validation",
  1076. validation_id=second_cycle["validation_id"],
  1077. replacement_contract=None,
  1078. child_contracts=(),
  1079. context={"root_trace_id": root, "tool_call_id": "accept-compose-v2"},
  1080. )
  1081. reloaded = await FileSystemTaskStore(str(tmp_path / "ledger")).load(root)
  1082. compose_task = reloaded.tasks[compose_id]
  1083. assert compose_task.status is TaskStatus.COMPLETED
  1084. assert compose_task.current_spec_version == 3
  1085. assert len(compose_task.attempt_ids) == 2
  1086. assert [
  1087. reloaded.validations[validation_id].verdict for validation_id in compose_task.validation_ids
  1088. ] == [ValidationVerdict.FAILED, ValidationVerdict.PASSED]
  1089. assert [reloaded.decisions[item].action for item in compose_task.decision_ids] == [
  1090. DecisionAction.REVISE,
  1091. DecisionAction.SPLIT,
  1092. DecisionAction.REVISE,
  1093. DecisionAction.ACCEPT,
  1094. ]
  1095. assert compose_accept["status"] == TaskStatus.COMPLETED.value
  1096. assert reloaded.tasks[replacement_id].status is TaskStatus.COMPLETED
  1097. @pytest.mark.asyncio
  1098. @pytest.mark.parametrize(
  1099. "creation_order",
  1100. [
  1101. ("element-set", "paragraph", "structure"),
  1102. ("paragraph", "element-set", "structure"),
  1103. ("structure", "paragraph", "element-set"),
  1104. ],
  1105. )
  1106. async def test_exploration_order_and_completion_order_do_not_choose_compose_order(
  1107. tmp_path: Path,
  1108. creation_order: tuple[str, str, str],
  1109. ) -> None:
  1110. service, coordinator, root = await _service(tmp_path)
  1111. coordinator.set_executor(_PlaceholderThenReplacementExecutor(coordinator))
  1112. portfolio = await service.plan_script_tasks(
  1113. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  1114. parent_task_id=None,
  1115. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  1116. )
  1117. compose = await service.plan_script_tasks(
  1118. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  1119. parent_task_id=portfolio["task_ids"][0],
  1120. context={"root_trace_id": root, "tool_call_id": "compose"},
  1121. )
  1122. compose_id = compose["task_ids"][0]
  1123. task_by_kind: dict[str, str] = {}
  1124. for kind in creation_order:
  1125. planned = await service.plan_script_tasks(
  1126. contract_payloads=[_planner_input(_contract(kind))],
  1127. parent_task_id=compose_id,
  1128. context={"root_trace_id": root, "tool_call_id": f"plan:{kind}"},
  1129. )
  1130. task_by_kind[kind] = planned["task_ids"][0]
  1131. decision_by_kind: dict[str, str] = {}
  1132. completion_order = tuple(reversed(creation_order))
  1133. for kind in completion_order:
  1134. task_id = task_by_kind[kind]
  1135. cycle = (
  1136. await service.dispatch_script_tasks(
  1137. task_ids=(task_id,),
  1138. context={"root_trace_id": root, "tool_call_id": f"dispatch:{kind}"},
  1139. )
  1140. )[0]
  1141. accepted = await service.decide_script_task(
  1142. task_id=task_id,
  1143. action=DecisionAction.ACCEPT.value,
  1144. reason=f"the bounded {kind} increment passed validation",
  1145. validation_id=cycle["validation_id"],
  1146. replacement_contract=None,
  1147. child_contracts=(),
  1148. context={"root_trace_id": root, "tool_call_id": f"accept:{kind}"},
  1149. )
  1150. decision_id = accepted["decision_id"]
  1151. decision_by_kind[kind] = decision_id
  1152. ledger = await coordinator.task_store.load(root)
  1153. assert ledger.tasks[compose_id].status is TaskStatus.NEEDS_REPLAN
  1154. adoption_kinds = ("paragraph", "structure", "element-set")
  1155. selected = tuple(decision_by_kind[kind] for kind in adoption_kinds)
  1156. await service.decide_script_task(
  1157. task_id=compose_id,
  1158. action=DecisionAction.REVISE.value,
  1159. reason="formal adoption is based on scope and intent, not timing",
  1160. validation_id=None,
  1161. replacement_contract=None,
  1162. child_contracts=(),
  1163. selected_decision_ids=selected,
  1164. context={"root_trace_id": root, "tool_call_id": "close-compose"},
  1165. )
  1166. reloaded = await FileSystemTaskStore(str(tmp_path / "ledger")).load(root)
  1167. frozen = await service.contract_for_task(root, reloaded.tasks[compose_id])
  1168. assert frozen.contract.compose_order == selected
  1169. assert [
  1170. next(kind for kind, task_id in task_by_kind.items() if task_id == child_id)
  1171. for child_id in reloaded.tasks[compose_id].child_task_ids
  1172. ] == list(creation_order)
  1173. assert completion_order == tuple(reversed(creation_order))
  1174. @pytest.mark.asyncio
  1175. async def test_phase_two_retrieval_counts_toward_task_and_depth_limits(tmp_path: Path) -> None:
  1176. service, _, root = await _service(tmp_path)
  1177. service.limits = PhaseTwoLimits(max_tasks=2)
  1178. portfolio = await service.plan_script_tasks(
  1179. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  1180. parent_task_id=None,
  1181. context={"root_trace_id": root, "tool_call_id": "portfolio"},
  1182. )
  1183. compose = await service.plan_script_tasks(
  1184. contract_payloads=[_planner_input(_contract("compose", execution_ready=False))],
  1185. parent_task_id=portfolio["task_ids"][0],
  1186. context={"root_trace_id": root, "tool_call_id": "compose"},
  1187. )
  1188. with pytest.raises(TaskContractError, match="Task budget"):
  1189. await service.plan_script_tasks(
  1190. contract_payloads=[_planner_input(_contract("decode-retrieval"))],
  1191. parent_task_id=compose["task_ids"][0],
  1192. context={"root_trace_id": root, "tool_call_id": "retrieval"},
  1193. )
  1194. service.limits = PhaseTwoLimits(max_depth=1)
  1195. with pytest.raises(TaskContractError, match="depth"):
  1196. await service.plan_script_tasks(
  1197. contract_payloads=[
  1198. _planner_input(_contract("structure", scope="script-build://scopes/full/local"))
  1199. ],
  1200. parent_task_id=compose["task_ids"][0],
  1201. context={"root_trace_id": root, "tool_call_id": "too-deep"},
  1202. )
  1203. @pytest.mark.asyncio
  1204. async def test_task_contract_cannot_be_silently_changed_in_task_spec(tmp_path: Path) -> None:
  1205. service, coordinator, root = await _service(tmp_path)
  1206. result = await service.plan_script_tasks(
  1207. contract_payloads=[_planner_input(_contract("candidate-portfolio", execution_ready=False))],
  1208. parent_task_id=None,
  1209. context={"root_trace_id": root},
  1210. )
  1211. ledger = await coordinator.task_store.load(root)
  1212. task = ledger.tasks[result["task_ids"][0]]
  1213. task.specs[-1] = replace(task.specs[-1], objective="tampered")
  1214. await coordinator.task_store.commit(ledger, expected_revision=ledger.revision)
  1215. with pytest.raises(TaskContractError, match="DIGEST_MISMATCH"):
  1216. await service.contract_for_task(root, task)