test_coordinator_integration.py 30 KB

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