test_mission_root.py 17 KB

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