test_coordinator_integration.py 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813
  1. import asyncio
  2. from collections import deque
  3. import pytest
  4. from agent.orchestration.config import OrchestrationConfig
  5. from agent.orchestration.coordinator import TaskConflict, TaskCoordinator
  6. from agent.orchestration.models import (
  7. AgentRole,
  8. ArtifactRef,
  9. AttemptSubmission,
  10. CriterionResult,
  11. DecisionAction,
  12. OperationKind,
  13. OperationStatus,
  14. TaskStatus,
  15. ValidationVerdict,
  16. )
  17. from agent.orchestration.protocols import ValidatorRunResult, WorkerRunResult
  18. from agent.orchestration.store import FileSystemArtifactStore, FileSystemTaskStore, TraceEventSink
  19. from agent.trace.models import Trace
  20. from agent.trace.store import FileSystemTraceStore
  21. from agent.core.runner import AgentRunner
  22. from agent.orchestration.wiring import wire_orchestration
  23. ROOT_TASK_SPEC = {
  24. "objective": "mission",
  25. "acceptance_criteria": [
  26. {"criterion_id": "mission-done", "description": "mission is complete", "hard": True}
  27. ],
  28. }
  29. class FakeExecutor:
  30. def __init__(self, verdicts, submit_worker=True, submit_validator=True, delay=0):
  31. self.verdicts = deque(verdicts)
  32. self.submit_worker = submit_worker
  33. self.submit_validator = submit_validator
  34. self.delay = delay
  35. self.coordinator = None
  36. self.active = 0
  37. self.max_active = 0
  38. self.worker_calls = 0
  39. self.validator_calls = 0
  40. async def run_worker(self, context):
  41. self.worker_calls += 1
  42. self.active += 1
  43. self.max_active = max(self.max_active, self.active)
  44. if self.delay:
  45. await asyncio.sleep(self.delay)
  46. try:
  47. if self.submit_worker:
  48. await self.coordinator.submit_attempt(
  49. {
  50. **context,
  51. "role": AgentRole.WORKER.value,
  52. "trace_id": context["worker_trace_id"],
  53. "tool_call_id": f"submit-{context['attempt_id']}",
  54. },
  55. AttemptSubmission(
  56. summary="done",
  57. artifact_refs=[ArtifactRef(uri=f"memory://{context['attempt_id']}", version="1")],
  58. ),
  59. )
  60. return WorkerRunResult(context["worker_trace_id"], "completed")
  61. finally:
  62. self.active -= 1
  63. async def run_validator(self, context):
  64. self.validator_calls += 1
  65. verdict = self.verdicts.popleft() if self.verdicts else ValidationVerdict.PASSED
  66. if self.submit_validator:
  67. criteria = [
  68. CriterionResult(
  69. criterion_id=item["criterion_id"], verdict=verdict,
  70. reason="checked",
  71. )
  72. for item in context["task_spec"]["acceptance_criteria"]
  73. ]
  74. await self.coordinator.submit_validation(
  75. {
  76. **context,
  77. "role": AgentRole.VALIDATOR.value,
  78. "trace_id": context["validator_trace_id"],
  79. "tool_call_id": f"validate-{context['validation_id']}",
  80. },
  81. verdict,
  82. criteria,
  83. "independent report",
  84. [], [], [], "accept" if verdict == ValidationVerdict.PASSED else "replan",
  85. )
  86. return ValidatorRunResult(context["validator_trace_id"], "completed")
  87. async def make_coordinator(
  88. tmp_path,
  89. executor,
  90. artifact_store=None,
  91. trace_store=None,
  92. task_store=None,
  93. ):
  94. trace_store = trace_store or FileSystemTraceStore(str(tmp_path))
  95. await trace_store.create_trace(Trace(trace_id="root", mode="agent", task="mission", agent_role="planner"))
  96. task_store = task_store or FileSystemTaskStore(str(tmp_path))
  97. coordinator = TaskCoordinator(
  98. task_store,
  99. artifact_store or FileSystemArtifactStore(str(tmp_path)),
  100. trace_store,
  101. OrchestrationConfig(max_parallel_tasks=4),
  102. TraceEventSink(str(tmp_path)),
  103. executor,
  104. )
  105. executor.coordinator = coordinator
  106. await coordinator.ensure_ledger("root", ROOT_TASK_SPEC)
  107. return coordinator, task_store, trace_store
  108. async def create_task(coordinator, objective="task", parent_task_id=None):
  109. result = await coordinator.create_tasks(
  110. "root",
  111. [{
  112. "objective": objective,
  113. "acceptance_criteria": [{"criterion_id": "c1", "description": "must pass", "hard": True}],
  114. }],
  115. parent_task_id=parent_task_id,
  116. )
  117. return result["tasks"][0]["task_id"]
  118. @pytest.mark.asyncio
  119. async def test_passed_validation_requires_planner_accept(tmp_path):
  120. executor = FakeExecutor([ValidationVerdict.PASSED])
  121. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  122. task_id = await create_task(coordinator)
  123. cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  124. assert cycle.task_status == TaskStatus.AWAITING_DECISION
  125. ledger = await store.load("root")
  126. assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION
  127. await coordinator.decide_task(
  128. "root", task_id, cycle.validation_id, DecisionAction.ACCEPT,
  129. {"reason": "all hard criteria passed"}, "decision-1",
  130. )
  131. assert (await store.load("root")).tasks[task_id].status == TaskStatus.COMPLETED
  132. @pytest.mark.asyncio
  133. async def test_failed_or_inconclusive_cannot_be_accepted(tmp_path):
  134. executor = FakeExecutor([ValidationVerdict.FAILED])
  135. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  136. task_id = await create_task(coordinator)
  137. cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  138. with pytest.raises(TaskConflict, match="passed"):
  139. await coordinator.decide_task(
  140. "root", task_id, cycle.validation_id, DecisionAction.ACCEPT,
  141. {"reason": "override"}, "bad-accept",
  142. )
  143. assert (await store.load("root")).tasks[task_id].status == TaskStatus.AWAITING_DECISION
  144. @pytest.mark.asyncio
  145. async def test_repair_is_limited_to_once_per_spec_version(tmp_path):
  146. executor = FakeExecutor([ValidationVerdict.FAILED, ValidationVerdict.FAILED])
  147. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  148. task_id = await create_task(coordinator)
  149. first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  150. await coordinator.decide_task(
  151. "root", task_id, first.validation_id, DecisionAction.REPAIR,
  152. {"reason": "small local correction"}, "repair-1",
  153. )
  154. second = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  155. ledger = await store.load("root")
  156. attempts = [ledger.attempts[x] for x in ledger.tasks[task_id].attempt_ids]
  157. assert attempts[0].worker_trace_id == attempts[1].worker_trace_id
  158. with pytest.raises(TaskConflict, match="limit"):
  159. await coordinator.decide_task(
  160. "root", task_id, second.validation_id, DecisionAction.REPAIR,
  161. {"reason": "try again"}, "repair-2",
  162. )
  163. @pytest.mark.asyncio
  164. async def test_worker_without_submit_attempt_needs_replan(tmp_path):
  165. executor = FakeExecutor([], submit_worker=False)
  166. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  167. task_id = await create_task(coordinator)
  168. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  169. ledger = await store.load("root")
  170. assert result.task_status == TaskStatus.NEEDS_REPLAN
  171. assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
  172. assert ledger.attempts[result.attempt_id].status.value == "failed"
  173. @pytest.mark.asyncio
  174. async def test_validator_without_submit_is_error_then_revalidates_with_new_trace(tmp_path):
  175. executor = FakeExecutor([ValidationVerdict.PASSED], submit_validator=False)
  176. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  177. task_id = await create_task(coordinator)
  178. first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  179. ledger = await store.load("root")
  180. first_validation = ledger.validations[first.validation_id]
  181. assert first.task_status == TaskStatus.NEEDS_REPLAN
  182. assert first_validation.status.value == "error"
  183. assert first_validation.verdict is None
  184. executor.submit_validator = True
  185. executor.verdicts.append(ValidationVerdict.PASSED)
  186. second = await coordinator.revalidate_attempt(
  187. "root", task_id, first.attempt_id, "revalidate-1"
  188. )
  189. assert second.task_status == TaskStatus.AWAITING_DECISION
  190. assert second.validation_id != first.validation_id
  191. assert second.validation.validator_trace_id != first_validation.validator_trace_id
  192. @pytest.mark.asyncio
  193. async def test_retry_uses_new_trace_and_revise_invalidates_old_validation(tmp_path):
  194. executor = FakeExecutor([
  195. ValidationVerdict.FAILED,
  196. ValidationVerdict.INCONCLUSIVE,
  197. ValidationVerdict.PASSED,
  198. ])
  199. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  200. task_id = await create_task(coordinator)
  201. first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  202. await coordinator.decide_task(
  203. "root", task_id, first.validation_id, DecisionAction.RETRY,
  204. {"reason": "use a fresh worker"}, "retry-1",
  205. )
  206. second = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  207. ledger = await store.load("root")
  208. first_attempt = ledger.attempts[first.attempt_id]
  209. second_attempt = ledger.attempts[second.attempt_id]
  210. assert first_attempt.worker_trace_id != second_attempt.worker_trace_id
  211. await coordinator.decide_task(
  212. "root", task_id, second.validation_id, DecisionAction.REVISE,
  213. {
  214. "reason": "clarify criterion",
  215. "objective": "revised task",
  216. "acceptance_criteria": [{"criterion_id": "c1", "description": "revised", "hard": True}],
  217. },
  218. "revise-1",
  219. )
  220. ledger = await store.load("root")
  221. assert ledger.tasks[task_id].current_spec_version == 2
  222. third = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  223. with pytest.raises(TaskConflict, match="obsolete"):
  224. await coordinator.decide_task(
  225. "root", task_id, first.validation_id, DecisionAction.ACCEPT,
  226. {"reason": "stale"}, "stale-accept",
  227. )
  228. await coordinator.decide_task(
  229. "root", task_id, third.validation_id, DecisionAction.ACCEPT,
  230. {"reason": "current validation passed"}, "current-accept",
  231. )
  232. @pytest.mark.asyncio
  233. async def test_block_unblock_cancel_and_supersede(tmp_path):
  234. executor = FakeExecutor([])
  235. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  236. task_id = await create_task(coordinator, "blocked")
  237. await coordinator.decide_task(
  238. "root", task_id, None, DecisionAction.BLOCK,
  239. {"reason": "external dependency"}, "block-1",
  240. )
  241. assert (await store.load("root")).tasks[task_id].blocked_reason == "external dependency"
  242. await coordinator.decide_task(
  243. "root", task_id, None, DecisionAction.UNBLOCK, {}, "unblock-1"
  244. )
  245. assert (await store.load("root")).tasks[task_id].status == TaskStatus.NEEDS_REPLAN
  246. await coordinator.decide_task(
  247. "root", task_id, None, DecisionAction.CANCEL,
  248. {"reason": "no longer needed"}, "cancel-1",
  249. )
  250. assert (await store.load("root")).tasks[task_id].status == TaskStatus.CANCELLED
  251. old_id = await create_task(coordinator, "old")
  252. result = await coordinator.decide_task(
  253. "root", old_id, None, DecisionAction.SUPERSEDE,
  254. {
  255. "reason": "replace spec",
  256. "replacement": {
  257. "objective": "replacement",
  258. "acceptance_criteria": [
  259. {"criterion_id": "replacement-c1", "description": "replacement passes"}
  260. ],
  261. "context_refs": ["replacement-context"],
  262. },
  263. },
  264. "supersede-1",
  265. )
  266. replacement_id = result["payload"]["replacement_task_id"]
  267. ledger = await store.load("root")
  268. assert ledger.tasks[old_id].status == TaskStatus.SUPERSEDED
  269. assert ledger.tasks[replacement_id].status == TaskStatus.PENDING
  270. replacement_spec = ledger.tasks[replacement_id].current_spec
  271. assert replacement_spec.objective == "replacement"
  272. assert replacement_spec.acceptance_criteria[0].criterion_id == "replacement-c1"
  273. assert replacement_spec.context_refs == ("replacement-context",)
  274. @pytest.mark.asyncio
  275. async def test_insert_after_keeps_stable_ids_and_updates_display_order(tmp_path):
  276. executor = FakeExecutor([])
  277. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  278. first = await create_task(coordinator, "first")
  279. third = await create_task(coordinator, "third")
  280. inserted = await coordinator.create_tasks(
  281. "root",
  282. [{
  283. "objective": "second",
  284. "acceptance_criteria": [
  285. {"criterion_id": "second-c1", "description": "second passes"}
  286. ],
  287. "context_refs": ["second-context"],
  288. }],
  289. placement={"after_task_id": first, "focus": True},
  290. idempotency_key="insert",
  291. )
  292. second = inserted["tasks"][0]["task_id"]
  293. ledger = await store.load("root")
  294. assert ledger.tasks[first].display_path == "0.1"
  295. assert ledger.tasks[second].display_path == "0.2"
  296. assert ledger.tasks[third].display_path == "0.3"
  297. assert ledger.tasks[second].current_spec.context_refs == ("second-context",)
  298. assert ledger.focused_task_id == second
  299. assert len({first, second, third}) == 3
  300. repeated = await coordinator.create_tasks(
  301. "root",
  302. [{
  303. "objective": "second",
  304. "acceptance_criteria": [
  305. {"criterion_id": "second-c1", "description": "second passes"}
  306. ],
  307. "context_refs": ["second-context"],
  308. }],
  309. placement={"after_task_id": first, "focus": True},
  310. idempotency_key="insert",
  311. )
  312. assert repeated["tasks"][0]["task_id"] == second
  313. @pytest.mark.asyncio
  314. async def test_split_children_do_not_auto_complete_parent(tmp_path):
  315. executor = FakeExecutor([ValidationVerdict.FAILED, ValidationVerdict.PASSED, ValidationVerdict.PASSED])
  316. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  317. parent_id = await create_task(coordinator, "parent")
  318. first = (await coordinator.dispatch_tasks("root", [parent_id]))[0]
  319. decision = await coordinator.decide_task(
  320. "root", parent_id, first.validation_id, DecisionAction.SPLIT,
  321. {
  322. "reason": "split work",
  323. "tasks": [
  324. {
  325. "objective": "child one",
  326. "acceptance_criteria": [{"criterion_id": "c1", "description": "pass"}],
  327. "context_refs": ["child-one-context"],
  328. },
  329. {
  330. "objective": "child two",
  331. "acceptance_criteria": [{"criterion_id": "c2", "description": "pass"}],
  332. "context_refs": ["child-two-context"],
  333. },
  334. ],
  335. },
  336. "split-1",
  337. )
  338. child_ids = decision["payload"]["child_task_ids"]
  339. ledger = await store.load("root")
  340. assert ledger.tasks[parent_id].child_task_ids == child_ids
  341. assert [ledger.tasks[x].display_path for x in child_ids] == ["0.1.1", "0.1.2"]
  342. assert [ledger.tasks[x].current_spec.context_refs for x in child_ids] == [
  343. ("child-one-context",),
  344. ("child-two-context",),
  345. ]
  346. assert [
  347. ledger.tasks[x].current_spec.acceptance_criteria[0].criterion_id
  348. for x in child_ids
  349. ] == ["c1", "c2"]
  350. cycles = await coordinator.dispatch_tasks("root", child_ids)
  351. for child_id, cycle in zip(child_ids, cycles):
  352. await coordinator.decide_task(
  353. "root", child_id, cycle.validation_id, DecisionAction.ACCEPT,
  354. {"reason": "passed"}, f"accept-{child_id}",
  355. )
  356. ledger = await store.load("root")
  357. assert all(ledger.tasks[x].status == TaskStatus.COMPLETED for x in child_ids)
  358. assert ledger.tasks[parent_id].status == TaskStatus.NEEDS_REPLAN
  359. @pytest.mark.asyncio
  360. async def test_parallel_tasks_are_isolated_and_bounded(tmp_path):
  361. executor = FakeExecutor([ValidationVerdict.PASSED] * 4, delay=0.02)
  362. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  363. task_ids = [await create_task(coordinator, f"task-{i}") for i in range(4)]
  364. results = await coordinator.dispatch_tasks("root", task_ids, idempotency_key="batch")
  365. assert [x.task_id for x in results] == task_ids
  366. assert executor.max_active <= 4
  367. assert all(x.task_status == TaskStatus.AWAITING_DECISION for x in results)
  368. ledger = await store.load("root")
  369. assert len({ledger.attempts[ledger.tasks[x].attempt_ids[-1]].worker_trace_id for x in task_ids}) == 4
  370. repeated = await coordinator.dispatch_tasks("root", task_ids, idempotency_key="batch")
  371. assert [x.attempt_id for x in repeated] == [x.attempt_id for x in results]
  372. @pytest.mark.asyncio
  373. async def test_real_runner_local_executor_creates_independent_terminal_traces(tmp_path):
  374. import json
  375. async def fake_llm(messages, tools, **kwargs):
  376. names = {item["function"]["name"] for item in tools or []}
  377. if "submit_attempt" in names:
  378. arguments = {
  379. "summary": "implemented",
  380. "artifact_refs": [{"uri": "memory://result", "version": "1"}],
  381. "evidence_refs": [],
  382. }
  383. tool_name = "submit_attempt"
  384. elif "submit_validation" in names:
  385. arguments = {
  386. "verdict": "passed",
  387. "criterion_results": [{"criterion_id": "c1", "verdict": "passed", "reason": "verified"}],
  388. "summary": "independent pass",
  389. "evidence_refs": [],
  390. "unverified_claims": [],
  391. "risks": [],
  392. "recommendation": "accept",
  393. }
  394. tool_name = "submit_validation"
  395. else:
  396. raise AssertionError(f"unexpected tools: {names}")
  397. return {
  398. "content": "",
  399. "tool_calls": [{
  400. "id": f"call-{tool_name}",
  401. "type": "function",
  402. "function": {"name": tool_name, "arguments": json.dumps(arguments)},
  403. }],
  404. "finish_reason": "tool_calls",
  405. }
  406. trace_store = FileSystemTraceStore(str(tmp_path))
  407. await trace_store.create_trace(Trace(trace_id="root", mode="agent", task="mission", agent_role="planner"))
  408. runner = AgentRunner(trace_store=trace_store, llm_call=fake_llm)
  409. coordinator = wire_orchestration(
  410. runner,
  411. FileSystemTaskStore(str(tmp_path)),
  412. FileSystemArtifactStore(str(tmp_path)),
  413. )
  414. await coordinator.ensure_ledger("root", ROOT_TASK_SPEC)
  415. task_id = await create_task(coordinator)
  416. cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  417. ledger = await coordinator.task_store.load("root")
  418. worker = await trace_store.get_trace(ledger.attempts[cycle.attempt_id].worker_trace_id)
  419. validator = await trace_store.get_trace(cycle.validation.validator_trace_id)
  420. assert cycle.task_status == TaskStatus.AWAITING_DECISION
  421. assert worker.agent_role == "worker" and worker.result_summary
  422. assert validator.agent_role == "validator" and validator.result_summary
  423. assert worker.trace_id != validator.trace_id
  424. class OneShotGetFailureArtifactStore(FileSystemArtifactStore):
  425. def __init__(self, base_path):
  426. super().__init__(base_path)
  427. self.get_calls = 0
  428. def for_root(self, root_trace_id):
  429. self.root_trace_id = root_trace_id
  430. return self
  431. async def get(self, root_trace_id, snapshot_id=None):
  432. self.get_calls += 1
  433. if self.get_calls == 2:
  434. raise RuntimeError("injected artifact read failure")
  435. return await super().get(root_trace_id, snapshot_id)
  436. class FailingDispatchCompletionTaskStore(FileSystemTaskStore):
  437. def __init__(self, base_path):
  438. super().__init__(base_path)
  439. self.fail_next_dispatch_completion = True
  440. async def commit(self, ledger, expected_revision, idempotency_key=None, event=None):
  441. if (
  442. self.fail_next_dispatch_completion
  443. and idempotency_key == "root:batch"
  444. ):
  445. self.fail_next_dispatch_completion = False
  446. raise RuntimeError("injected dispatch completion persistence failure")
  447. return await super().commit(
  448. ledger,
  449. expected_revision,
  450. idempotency_key=idempotency_key,
  451. event=event,
  452. )
  453. class OneShotLoadFailureTaskStore(FileSystemTaskStore):
  454. def __init__(self, base_path):
  455. super().__init__(base_path)
  456. self.fail_next_load = False
  457. self._active_operation_loads_before_failure = 1
  458. async def load(self, root_trace_id):
  459. ledger = await super().load(root_trace_id)
  460. if self.fail_next_load and any(
  461. operation.status.value == "running"
  462. for operation in ledger.operations.values()
  463. ):
  464. if self._active_operation_loads_before_failure:
  465. self._active_operation_loads_before_failure -= 1
  466. else:
  467. self.fail_next_load = False
  468. raise RuntimeError("injected reservation load failure")
  469. return ledger
  470. class FailingEventSink:
  471. async def emit(self, root_trace_id, event_type, payload):
  472. raise RuntimeError("injected event sink failure")
  473. @pytest.mark.asyncio
  474. async def test_parallel_unexpected_branch_error_does_not_cancel_siblings(tmp_path):
  475. executor = FakeExecutor([ValidationVerdict.PASSED] * 4, delay=0.01)
  476. artifact_store = OneShotGetFailureArtifactStore(str(tmp_path))
  477. coordinator, store, _ = await make_coordinator(
  478. tmp_path, executor, artifact_store=artifact_store
  479. )
  480. task_ids = [await create_task(coordinator, f"isolated-{index}") for index in range(4)]
  481. results = await coordinator.dispatch_tasks("root", task_ids)
  482. assert [result.task_id for result in results] == task_ids
  483. failures = [result for result in results if result.error]
  484. assert len(failures) == 1
  485. assert "artifact read failure" in failures[0].error
  486. assert failures[0].task_status == TaskStatus.NEEDS_REPLAN
  487. successes = [result for result in results if not result.error]
  488. assert len(successes) == 3
  489. assert all(result.task_status == TaskStatus.AWAITING_DECISION for result in successes)
  490. ledger = await store.load("root")
  491. assert all(
  492. ledger.tasks[task_id].status not in {TaskStatus.RUNNING, TaskStatus.VALIDATING}
  493. for task_id in task_ids
  494. )
  495. @pytest.mark.asyncio
  496. async def test_batch_reservation_conflict_does_not_strand_other_tasks(tmp_path):
  497. executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED])
  498. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  499. task_ids = [await create_task(coordinator, f"reserve-{index}") for index in range(3)]
  500. existing = await coordinator._create_attempt(
  501. "root", task_ids[1], "worker", "existing-reservation"
  502. )
  503. results = await coordinator.dispatch_tasks("root", task_ids)
  504. assert [result.task_id for result in results] == task_ids
  505. assert results[0].task_status == TaskStatus.AWAITING_DECISION
  506. assert results[1].attempt_id is None
  507. assert "Dispatch conflict" in results[1].error
  508. assert results[2].task_status == TaskStatus.AWAITING_DECISION
  509. ledger = await store.load("root")
  510. assert ledger.tasks[task_ids[0]].status == TaskStatus.AWAITING_DECISION
  511. assert ledger.tasks[task_ids[2]].status == TaskStatus.AWAITING_DECISION
  512. assert ledger.attempts[existing["attempt_id"]].status.value == "running"
  513. @pytest.mark.asyncio
  514. async def test_explicit_lifecycle_never_creates_goal_tree(tmp_path):
  515. executor = FakeExecutor([ValidationVerdict.PASSED])
  516. coordinator, store, trace_store = await make_coordinator(tmp_path, executor)
  517. task_id = await create_task(coordinator, "task-ledger-only")
  518. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  519. assert result.task_status == TaskStatus.AWAITING_DECISION
  520. assert (await store.load("root")).tasks[task_id].status == TaskStatus.AWAITING_DECISION
  521. assert await trace_store.get_goal_tree("root") is None
  522. @pytest.mark.asyncio
  523. async def test_task_context_is_rendered_from_ledger(tmp_path):
  524. coordinator, _, _ = await make_coordinator(tmp_path, FakeExecutor([]))
  525. await create_task(coordinator, "draft the answer")
  526. context = await coordinator.task_context("root")
  527. assert "## Current Task Plan" in context
  528. assert "- 0 [waiting_children] mission" in context
  529. assert "- 0.1 [pending] draft the answer" in context
  530. @pytest.mark.asyncio
  531. async def test_dispatch_returns_results_when_batch_result_persistence_fails(tmp_path):
  532. executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED])
  533. task_store = FailingDispatchCompletionTaskStore(str(tmp_path))
  534. coordinator, store, _ = await make_coordinator(
  535. tmp_path,
  536. executor,
  537. task_store=task_store,
  538. )
  539. task_ids = [await create_task(coordinator, f"batch-{index}") for index in range(2)]
  540. first = await coordinator.dispatch_tasks(
  541. "root",
  542. task_ids,
  543. idempotency_key="batch",
  544. )
  545. assert [result.task_status for result in first] == [
  546. TaskStatus.AWAITING_DECISION,
  547. TaskStatus.AWAITING_DECISION,
  548. ]
  549. assert executor.worker_calls == 2
  550. assert executor.validator_calls == 2
  551. repeated = await coordinator.dispatch_tasks(
  552. "root",
  553. task_ids,
  554. idempotency_key="batch",
  555. )
  556. assert [result.attempt_id for result in repeated] == [
  557. result.attempt_id for result in first
  558. ]
  559. assert executor.worker_calls == 2
  560. assert executor.validator_calls == 2
  561. ledger = await store.load("root")
  562. assert all(len(ledger.tasks[task_id].attempt_ids) == 1 for task_id in task_ids)
  563. assert all(len(ledger.tasks[task_id].validation_ids) == 1 for task_id in task_ids)
  564. @pytest.mark.asyncio
  565. async def test_dispatch_idempotency_key_is_bound_to_tasks_and_presets(tmp_path):
  566. executor = FakeExecutor([ValidationVerdict.PASSED])
  567. coordinator, _, _ = await make_coordinator(tmp_path, executor)
  568. first_task = await create_task(coordinator, "first-binding")
  569. second_task = await create_task(coordinator, "second-binding")
  570. await coordinator.dispatch_tasks(
  571. "root",
  572. [first_task],
  573. worker_presets=["worker"],
  574. idempotency_key="bound-batch",
  575. )
  576. with pytest.raises(TaskConflict, match="different task_ids"):
  577. await coordinator.dispatch_tasks(
  578. "root",
  579. [second_task],
  580. worker_presets=["worker"],
  581. idempotency_key="bound-batch",
  582. )
  583. with pytest.raises(TaskConflict, match="different worker_presets"):
  584. await coordinator.dispatch_tasks(
  585. "root",
  586. [first_task],
  587. worker_presets=["alternate-worker"],
  588. idempotency_key="bound-batch",
  589. )
  590. assert executor.worker_calls == 1
  591. assert executor.validator_calls == 1
  592. @pytest.mark.asyncio
  593. async def test_concurrent_dispatch_replay_does_not_duplicate_agent_execution(tmp_path):
  594. executor = FakeExecutor([ValidationVerdict.PASSED], delay=0.02)
  595. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  596. task_id = await create_task(coordinator, "concurrent-idempotency")
  597. first, replay = await asyncio.gather(
  598. coordinator.dispatch_tasks("root", [task_id], idempotency_key="concurrent"),
  599. coordinator.dispatch_tasks("root", [task_id], idempotency_key="concurrent"),
  600. )
  601. assert executor.worker_calls == 1
  602. assert executor.validator_calls == 1
  603. assert any(
  604. result[0].task_status == TaskStatus.AWAITING_DECISION
  605. for result in (first, replay)
  606. )
  607. assert first[0].attempt_id == replay[0].attempt_id
  608. assert all(result[0].task_status == TaskStatus.AWAITING_DECISION for result in (first, replay))
  609. final = await coordinator.dispatch_tasks(
  610. "root",
  611. [task_id],
  612. idempotency_key="concurrent",
  613. )
  614. ledger = await store.load("root")
  615. assert final[0].task_status == TaskStatus.AWAITING_DECISION
  616. assert len(ledger.tasks[task_id].attempt_ids) == 1
  617. assert len(ledger.tasks[task_id].validation_ids) == 1
  618. @pytest.mark.asyncio
  619. async def test_reservation_load_failure_does_not_corrupt_existing_attempt(tmp_path):
  620. executor = FakeExecutor([ValidationVerdict.PASSED])
  621. task_store = OneShotLoadFailureTaskStore(str(tmp_path))
  622. coordinator, store, _ = await make_coordinator(
  623. tmp_path,
  624. executor,
  625. task_store=task_store,
  626. )
  627. running_task = await create_task(coordinator, "already-running")
  628. existing = await coordinator._create_attempt(
  629. "root", running_task, "worker", "existing-running-attempt"
  630. )
  631. sibling_task = await create_task(coordinator, "unrelated-sibling")
  632. task_store.fail_next_load = True
  633. results = await coordinator.dispatch_tasks(
  634. "root",
  635. [running_task, sibling_task],
  636. )
  637. assert [result.task_id for result in results] == [running_task, sibling_task]
  638. assert results[0].attempt_id is None
  639. assert results[0].task_status == TaskStatus.RUNNING
  640. assert "reservation load failure" in results[0].error
  641. assert results[1].task_status == TaskStatus.AWAITING_DECISION
  642. ledger = await store.load("root")
  643. assert ledger.tasks[running_task].status == TaskStatus.RUNNING
  644. assert ledger.attempts[existing["attempt_id"]].status.value == "running"
  645. @pytest.mark.asyncio
  646. async def test_event_sink_failure_does_not_rollback_committed_ledger(tmp_path):
  647. executor = FakeExecutor([])
  648. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  649. coordinator.event_sink = FailingEventSink()
  650. task_id = await create_task(coordinator, "event-outage")
  651. ledger = await store.load("root")
  652. assert ledger.tasks[task_id].status == TaskStatus.PENDING
  653. @pytest.mark.asyncio
  654. async def test_background_operation_is_durable_and_idempotently_bound(tmp_path):
  655. executor = FakeExecutor([ValidationVerdict.PASSED])
  656. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  657. task_id = await create_task(coordinator, "durable-operation")
  658. operation = await coordinator.start_operation(
  659. "root",
  660. OperationKind.DISPATCH,
  661. task_ids=[task_id],
  662. idempotency_key="operation-key",
  663. )
  664. completed = await coordinator.await_operation("root", operation.operation_id)
  665. replay = await coordinator.start_operation(
  666. "root",
  667. OperationKind.DISPATCH,
  668. task_ids=[task_id],
  669. idempotency_key="operation-key",
  670. )
  671. ledger = await store.load("root")
  672. assert completed.status == OperationStatus.COMPLETED
  673. assert replay.operation_id == completed.operation_id
  674. assert ledger.operations[completed.operation_id].attempt_ids
  675. assert ledger.operations[completed.operation_id].validation_ids
  676. assert executor.worker_calls == executor.validator_calls == 1
  677. other = await create_task(coordinator, "different-operation")
  678. with pytest.raises(TaskConflict, match="different task_ids"):
  679. await coordinator.start_operation(
  680. "root",
  681. OperationKind.DISPATCH,
  682. task_ids=[other],
  683. idempotency_key="operation-key",
  684. )
  685. @pytest.mark.asyncio
  686. async def test_submitted_attempt_advances_to_validation_without_worker_replay(tmp_path):
  687. executor = FakeExecutor([ValidationVerdict.PASSED])
  688. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  689. task_id = await create_task(coordinator, "resume-after-submit")
  690. reserved = await coordinator._create_attempt(
  691. "root", task_id, "worker", "resume-reservation"
  692. )
  693. await coordinator.submit_attempt(
  694. {
  695. "role": AgentRole.WORKER.value,
  696. "root_trace_id": "root",
  697. "task_id": task_id,
  698. "attempt_id": reserved["attempt_id"],
  699. "trace_id": reserved["worker_trace_id"],
  700. "spec_version": 1,
  701. "tool_call_id": "resume-submission",
  702. },
  703. AttemptSubmission(
  704. summary="already executed",
  705. artifact_refs=[ArtifactRef(uri="memory://submitted", version="1")],
  706. ),
  707. )
  708. result = await coordinator.advance_cycle("root", task_id, reserved["attempt_id"])
  709. assert result.task_status == TaskStatus.AWAITING_DECISION
  710. assert executor.worker_calls == 0
  711. assert executor.validator_calls == 1
  712. assert len((await store.load("root")).tasks[task_id].attempt_ids) == 1