test_mission_root.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427
  1. import json
  2. import pytest
  3. from agent.orchestration.coordinator import TaskConflict, TaskCoordinator
  4. from agent.orchestration.models import DecisionAction, TaskLedger, TaskStatus, ValidationVerdict
  5. from agent.orchestration.store import (
  6. FileSystemArtifactStore,
  7. FileSystemTaskStore,
  8. TaskStoreError,
  9. )
  10. from agent.trace.store import FileSystemTraceStore
  11. from test_coordinator_integration import (
  12. FakeExecutor,
  13. ROOT_TASK_SPEC,
  14. create_task,
  15. make_coordinator,
  16. )
  17. @pytest.mark.asyncio
  18. async def test_new_ledger_has_one_stable_root_and_inspect_exposes_it(tmp_path):
  19. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  20. ledger = await store.load("root")
  21. root_id = ledger.root_task_id
  22. assert root_id is not None
  23. assert ledger.focused_task_id == root_id
  24. assert ledger.tasks[root_id].parent_task_id is None
  25. assert ledger.tasks[root_id].display_path == "0"
  26. assert ledger.tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  27. assert [task.task_id for task in ledger.tasks.values() if task.parent_task_id is None] == [root_id]
  28. same = await coordinator.ensure_ledger("root", None)
  29. assert same.root_task_id == root_id
  30. inspected = await coordinator.inspect_tasks("root")
  31. assert inspected["root_task_id"] == root_id
  32. assert inspected["tasks"][0]["current_spec"]["acceptance_criteria"][0][
  33. "criterion_id"
  34. ] == "mission-done"
  35. @pytest.mark.asyncio
  36. async def test_new_ledger_requires_root_spec_and_rootless_policy_is_explicit(tmp_path):
  37. store = FileSystemTaskStore(str(tmp_path))
  38. coordinator = TaskCoordinator(
  39. store,
  40. FileSystemArtifactStore(str(tmp_path)),
  41. FileSystemTraceStore(str(tmp_path)),
  42. )
  43. with pytest.raises(ValueError, match="root_task_spec is required"):
  44. await coordinator.ensure_ledger("missing")
  45. with pytest.raises(TaskConflict, match="has no Root Task ledger"):
  46. await coordinator.ensure_ledger(
  47. "missing-resume", ROOT_TASK_SPEC, allow_create=False
  48. )
  49. legacy = TaskLedger(root_trace_id="legacy", mission="legacy")
  50. committed = await store.commit(legacy, expected_revision=-1)
  51. with pytest.raises(TaskConflict, match="Rootless"):
  52. await coordinator.ensure_ledger("legacy", ROOT_TASK_SPEC)
  53. unchanged = await store.load("legacy")
  54. assert unchanged.root_task_id is None
  55. assert unchanged.revision == committed.ledger.revision
  56. invalid_path = store.ledger_path("legacy-invalid")
  57. invalid_path.parent.mkdir(parents=True, exist_ok=True)
  58. invalid_path.write_text(json.dumps({
  59. "root_trace_id": "legacy-invalid",
  60. "mission": "legacy invalid",
  61. "tasks": {
  62. "old-task": {
  63. "task_id": "old-task",
  64. "parent_task_id": None,
  65. "display_path": "1",
  66. "specs": [{
  67. "version": 1,
  68. "objective": "old task",
  69. "acceptance_criteria": [],
  70. }],
  71. }
  72. },
  73. }), encoding="utf-8")
  74. with pytest.raises(TaskStoreError, match="at least one criterion"):
  75. await store.load("legacy-invalid")
  76. @pytest.mark.asyncio
  77. async def test_parent_waits_until_every_child_is_terminal(tmp_path):
  78. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  79. first = await create_task(coordinator, "first")
  80. second = await create_task(coordinator, "second")
  81. third = await create_task(coordinator, "third")
  82. root_id = (await store.load("root")).root_task_id
  83. await coordinator.decide_task(
  84. "root", first, None, DecisionAction.CANCEL, {"reason": "drop first"}, "cancel-first"
  85. )
  86. assert (await store.load("root")).tasks[root_id].status == TaskStatus.WAITING_CHILDREN
  87. superseded = await coordinator.decide_task(
  88. "root",
  89. second,
  90. None,
  91. DecisionAction.SUPERSEDE,
  92. {
  93. "reason": "replace second",
  94. "replacement": {
  95. "objective": "second replacement",
  96. "acceptance_criteria": [
  97. {"criterion_id": "replacement", "description": "replacement passes"}
  98. ],
  99. },
  100. },
  101. "supersede-second",
  102. )
  103. replacement = superseded["payload"]["replacement_task_id"]
  104. assert (await store.load("root")).tasks[root_id].status == TaskStatus.WAITING_CHILDREN
  105. await coordinator.decide_task(
  106. "root", replacement, None, DecisionAction.CANCEL, {"reason": "drop replacement"}, "cancel-replacement"
  107. )
  108. assert (await store.load("root")).tasks[root_id].status == TaskStatus.WAITING_CHILDREN
  109. await coordinator.decide_task(
  110. "root", third, None, DecisionAction.CANCEL, {"reason": "drop third"}, "cancel-third"
  111. )
  112. assert (await store.load("root")).tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  113. @pytest.mark.asyncio
  114. async def test_root_executes_after_children_and_receives_only_accepted_results(tmp_path):
  115. class CapturingExecutor(FakeExecutor):
  116. def __init__(self):
  117. super().__init__([ValidationVerdict.PASSED, ValidationVerdict.PASSED])
  118. self.worker_contexts = []
  119. async def run_worker(self, context):
  120. self.worker_contexts.append(context)
  121. return await super().run_worker(context)
  122. executor = CapturingExecutor()
  123. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  124. child_id = await create_task(coordinator, "child")
  125. root_id = (await store.load("root")).root_task_id
  126. blocked_dispatch = (await coordinator.dispatch_tasks("root", [root_id]))[0]
  127. assert "cannot dispatch" in blocked_dispatch.error
  128. child_cycle = (await coordinator.dispatch_tasks("root", [child_id]))[0]
  129. await coordinator.decide_task(
  130. "root",
  131. child_id,
  132. child_cycle.validation_id,
  133. DecisionAction.ACCEPT,
  134. {"reason": "child accepted"},
  135. "accept-child",
  136. )
  137. accepted_ledger = await store.load("root")
  138. assert accepted_ledger.tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  139. assert [
  140. item["task_id"]
  141. for item in coordinator._accepted_child_results(
  142. accepted_ledger, accepted_ledger.tasks[root_id]
  143. )
  144. ] == [child_id]
  145. root_cycle = (await coordinator.dispatch_tasks("root", [root_id]))[0]
  146. assert root_cycle.error is None
  147. assert len(executor.worker_contexts) == 2
  148. child_results = executor.worker_contexts[-1]["accepted_child_results"]
  149. assert [item["task_id"] for item in child_results] == [child_id]
  150. assert child_results[0]["attempt_id"] == child_cycle.attempt_id
  151. assert child_results[0]["validation"]["validation_id"] == child_cycle.validation_id
  152. await coordinator.decide_task(
  153. "root",
  154. root_id,
  155. root_cycle.validation_id,
  156. DecisionAction.ACCEPT,
  157. {"reason": "root accepted"},
  158. "accept-root",
  159. )
  160. completion = await coordinator.mission_completion("root")
  161. assert completion["status"] == TaskStatus.COMPLETED.value
  162. assert completion["result_summary"] == "done"
  163. @pytest.mark.asyncio
  164. async def test_root_identity_is_not_cancellable_or_supersedable(tmp_path):
  165. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  166. root_id = (await store.load("root")).root_task_id
  167. with pytest.raises(TaskConflict, match="root task cannot be cancelled"):
  168. await coordinator.decide_task(
  169. "root", root_id, None, DecisionAction.CANCEL, {"reason": "invalid"}, "cancel-root"
  170. )
  171. with pytest.raises(TaskConflict, match="root task cannot be superseded"):
  172. await coordinator.decide_task(
  173. "root",
  174. root_id,
  175. None,
  176. DecisionAction.SUPERSEDE,
  177. {
  178. "reason": "invalid",
  179. "replacement": {
  180. "objective": "replacement",
  181. "acceptance_criteria": [
  182. {"criterion_id": "replacement", "description": "replacement passes"}
  183. ],
  184. },
  185. },
  186. "supersede-root",
  187. )
  188. @pytest.mark.asyncio
  189. async def test_parent_with_active_descendant_cannot_be_terminated(tmp_path):
  190. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  191. parent_id = await create_task(coordinator, "parent")
  192. descendant_id = await create_task(
  193. coordinator, "descendant", parent_task_id=parent_id
  194. )
  195. root_id = (await store.load("root")).root_task_id
  196. with pytest.raises(TaskConflict, match="non-terminal descendants"):
  197. await coordinator.decide_task(
  198. "root",
  199. parent_id,
  200. None,
  201. DecisionAction.CANCEL,
  202. {"reason": "invalid while descendant is active"},
  203. "cancel-active-parent",
  204. )
  205. with pytest.raises(TaskConflict, match="non-terminal descendants"):
  206. await coordinator.decide_task(
  207. "root",
  208. parent_id,
  209. None,
  210. DecisionAction.SUPERSEDE,
  211. {
  212. "reason": "invalid while descendant is active",
  213. "replacement": {
  214. "objective": "replacement",
  215. "acceptance_criteria": [
  216. {"criterion_id": "replacement", "description": "replacement passes"}
  217. ],
  218. },
  219. },
  220. "supersede-active-parent",
  221. )
  222. root_dispatch = (await coordinator.dispatch_tasks("root", [root_id]))[0]
  223. assert "cannot dispatch" in root_dispatch.error
  224. await coordinator.decide_task(
  225. "root",
  226. descendant_id,
  227. None,
  228. DecisionAction.CANCEL,
  229. {"reason": "finish descendant"},
  230. "cancel-descendant",
  231. )
  232. await coordinator.decide_task(
  233. "root",
  234. parent_id,
  235. None,
  236. DecisionAction.CANCEL,
  237. {"reason": "parent can now terminate"},
  238. "cancel-parent",
  239. )
  240. assert (await store.load("root")).tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  241. @pytest.mark.asyncio
  242. async def test_blocked_descendants_allow_mission_to_report_blocked(tmp_path):
  243. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  244. parent_id = await create_task(coordinator, "parent")
  245. child_id = await create_task(coordinator, "child", parent_task_id=parent_id)
  246. root_id = (await store.load("root")).root_task_id
  247. await coordinator.decide_task(
  248. "root",
  249. child_id,
  250. None,
  251. DecisionAction.BLOCK,
  252. {"reason": "child external dependency"},
  253. "block-child",
  254. )
  255. await coordinator.decide_task(
  256. "root",
  257. parent_id,
  258. None,
  259. DecisionAction.BLOCK,
  260. {"reason": "parent waits on blocked child"},
  261. "block-parent",
  262. )
  263. await coordinator.decide_task(
  264. "root",
  265. root_id,
  266. None,
  267. DecisionAction.BLOCK,
  268. {"reason": "mission external dependency"},
  269. "block-root",
  270. )
  271. completion = await coordinator.mission_completion("root")
  272. assert completion["status"] == TaskStatus.BLOCKED.value
  273. assert completion["blocked_reason"] == "mission external dependency"
  274. await coordinator.decide_task(
  275. "root",
  276. root_id,
  277. None,
  278. DecisionAction.UNBLOCK,
  279. {},
  280. "unblock-root",
  281. )
  282. assert (await store.load("root")).tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  283. @pytest.mark.asyncio
  284. async def test_parent_cannot_block_while_descendant_attempt_is_active(tmp_path):
  285. coordinator, _, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  286. parent_id = await create_task(coordinator, "parent")
  287. child_id = await create_task(coordinator, "child", parent_task_id=parent_id)
  288. await coordinator._create_attempt("root", child_id, "worker", None)
  289. with pytest.raises(TaskConflict, match="active descendants"):
  290. await coordinator.decide_task(
  291. "root",
  292. parent_id,
  293. None,
  294. DecisionAction.BLOCK,
  295. {"reason": "must not hide active work"},
  296. "block-active-parent",
  297. )
  298. @pytest.mark.asyncio
  299. async def test_inserting_a_sibling_rebases_its_existing_subtree(tmp_path):
  300. coordinator, store, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  301. first_id = await create_task(coordinator, "first")
  302. moved_id = await create_task(coordinator, "moved")
  303. descendant_id = await create_task(
  304. coordinator, "moved descendant", parent_task_id=moved_id
  305. )
  306. await coordinator.create_tasks(
  307. "root",
  308. [{
  309. "objective": "inserted",
  310. "acceptance_criteria": [
  311. {"criterion_id": "inserted", "description": "inserted passes"}
  312. ],
  313. }],
  314. placement={"after_task_id": first_id},
  315. )
  316. ledger = await store.load("root")
  317. assert ledger.tasks[moved_id].display_path == "0.3"
  318. assert ledger.tasks[descendant_id].display_path == "0.3.1"
  319. @pytest.mark.asyncio
  320. async def test_invalid_batch_and_revision_are_atomic(tmp_path):
  321. executor = FakeExecutor([ValidationVerdict.FAILED])
  322. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  323. root_id = (await store.load("root")).root_task_id
  324. with pytest.raises(ValueError, match="at least one criterion"):
  325. await coordinator.create_tasks(
  326. "root",
  327. [
  328. {
  329. "objective": "valid first draft",
  330. "acceptance_criteria": [
  331. {"criterion_id": "valid", "description": "valid passes"}
  332. ],
  333. },
  334. {"objective": "invalid second draft"},
  335. ],
  336. )
  337. ledger = await store.load("root")
  338. assert list(ledger.tasks) == [root_id]
  339. assert ledger.tasks[root_id].status == TaskStatus.NEEDS_REPLAN
  340. with pytest.raises(ValueError, match="context_refs"):
  341. await coordinator.create_tasks(
  342. "root",
  343. [{
  344. "objective": "invalid context refs",
  345. "acceptance_criteria": [
  346. {"criterion_id": "context", "description": "context passes"}
  347. ],
  348. "context_refs": "memory://not-a-list",
  349. }],
  350. )
  351. task_id = await create_task(coordinator, "revise atomically")
  352. cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  353. for index, criteria in enumerate(([], None, "criterion")):
  354. with pytest.raises(ValueError, match="acceptance_criteria|at least one criterion"):
  355. await coordinator.decide_task(
  356. "root",
  357. task_id,
  358. cycle.validation_id,
  359. DecisionAction.REVISE,
  360. {"reason": "invalid revision", "acceptance_criteria": criteria},
  361. f"invalid-revision-{index}",
  362. )
  363. task = (await store.load("root")).tasks[task_id]
  364. assert task.current_spec_version == 1
  365. assert len(task.specs) == 1
  366. assert task.status == TaskStatus.AWAITING_DECISION
  367. for index, objective in enumerate((123, None, "", " ")):
  368. with pytest.raises(ValueError, match="objective"):
  369. await coordinator.decide_task(
  370. "root",
  371. task_id,
  372. cycle.validation_id,
  373. DecisionAction.REVISE,
  374. {"reason": "invalid objective", "objective": objective},
  375. f"invalid-objective-{index}",
  376. )
  377. unchanged = (await store.load("root")).tasks[task_id]
  378. assert unchanged.current_spec_version == 1
  379. assert unchanged.status == TaskStatus.AWAITING_DECISION