test_mission_root.py 16 KB

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