test_coordinator_integration.py 33 KB

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