test_orchestration_v2_control.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675
  1. import asyncio
  2. from dataclasses import replace
  3. from types import SimpleNamespace
  4. import pytest
  5. from agent import AgentExecutor, FailureDetail, FailureDisposition
  6. from agent.orchestration.config import OrchestrationConfig
  7. from agent.orchestration.coordinator import TaskCoordinator
  8. from agent.orchestration.executor import LocalAgentExecutor
  9. from agent.orchestration.models import (
  10. AttemptStatus,
  11. AttemptSubmission,
  12. DecisionAction,
  13. ExecutionStats,
  14. FailureCode,
  15. TaskAttempt,
  16. TaskStatus,
  17. ValidationVerdict,
  18. utc_now,
  19. )
  20. from agent.orchestration.protocols import AgentExecutor as ProtocolAgentExecutor
  21. from agent.orchestration.protocols import WorkerRunResult
  22. from agent.orchestration.store import (
  23. FileSystemArtifactStore,
  24. FileSystemTaskStore,
  25. TraceEventSink,
  26. )
  27. from test_coordinator_integration import FakeExecutor, create_task, make_coordinator
  28. def test_v2_control_config_defaults_preserve_unbounded_v1_execution():
  29. config = OrchestrationConfig()
  30. assert config.worker_timeout_seconds is None
  31. assert config.validator_timeout_seconds is None
  32. assert config.stop_grace_seconds == 5.0
  33. def test_agent_executor_port_is_available_from_public_package():
  34. assert AgentExecutor is ProtocolAgentExecutor
  35. @pytest.mark.parametrize(
  36. ("field", "value"),
  37. [
  38. ("worker_timeout_seconds", 0),
  39. ("worker_timeout_seconds", -0.1),
  40. ("validator_timeout_seconds", 0),
  41. ("validator_timeout_seconds", -1),
  42. ("validator_timeout_seconds", float("nan")),
  43. ("worker_timeout_seconds", float("inf")),
  44. ("worker_timeout_seconds", True),
  45. ("stop_grace_seconds", -0.1),
  46. ("stop_grace_seconds", float("nan")),
  47. ("stop_grace_seconds", float("inf")),
  48. ("stop_grace_seconds", False),
  49. ],
  50. )
  51. def test_v2_control_config_rejects_invalid_durations(field, value):
  52. with pytest.raises(ValueError, match=field):
  53. OrchestrationConfig(**{field: value})
  54. @pytest.mark.parametrize("field", ["max_parallel_tasks", "max_repair_continuations"])
  55. def test_v2_control_config_rejects_boolean_integer_limits(field):
  56. with pytest.raises(ValueError, match=field):
  57. OrchestrationConfig(**{field: True})
  58. def test_v2_control_config_allows_immediate_stop_and_explicit_timeouts():
  59. config = OrchestrationConfig(
  60. worker_timeout_seconds=1.5,
  61. validator_timeout_seconds=2,
  62. stop_grace_seconds=0,
  63. )
  64. assert config.worker_timeout_seconds == 1.5
  65. assert config.validator_timeout_seconds == 2
  66. assert config.stop_grace_seconds == 0
  67. @pytest.mark.asyncio
  68. async def test_local_executor_stop_delegates_every_request():
  69. class Runner:
  70. def __init__(self):
  71. self.calls = []
  72. async def stop(self, trace_id):
  73. self.calls.append(trace_id)
  74. return len(self.calls) == 1
  75. runner = Runner()
  76. executor = LocalAgentExecutor(runner)
  77. assert await executor.stop("worker-trace") is True
  78. assert await executor.stop("worker-trace") is False
  79. assert runner.calls == ["worker-trace", "worker-trace"]
  80. @pytest.mark.asyncio
  81. async def test_local_executor_stop_does_not_hide_runner_failure():
  82. class Runner:
  83. async def stop(self, trace_id):
  84. raise RuntimeError(f"cannot stop {trace_id}")
  85. with pytest.raises(RuntimeError, match="cannot stop validator-trace"):
  86. await LocalAgentExecutor(Runner()).stop("validator-trace")
  87. @pytest.mark.asyncio
  88. async def test_hard_cancellation_is_not_normalized_as_agent_failure():
  89. class Runner:
  90. async def run_result(self, **kwargs):
  91. raise asyncio.CancelledError
  92. with pytest.raises(asyncio.CancelledError):
  93. await LocalAgentExecutor(Runner()).run_worker({
  94. "worker_preset": "worker",
  95. "worker_trace_id": "worker-trace",
  96. "root_trace_id": "root",
  97. "task_id": "task",
  98. "spec_version": 1,
  99. "attempt_id": "attempt",
  100. "task_spec": {},
  101. "continue_trace_id": None,
  102. })
  103. @pytest.mark.asyncio
  104. async def test_failed_stats_lookup_does_not_mask_original_runner_error():
  105. class Store:
  106. async def get_trace(self, trace_id):
  107. raise RuntimeError("stats store unavailable")
  108. class Runner:
  109. trace_store = Store()
  110. async def run_result(self, **kwargs):
  111. raise RuntimeError("original runner failure")
  112. result = await LocalAgentExecutor(Runner()).run_worker({
  113. "worker_preset": "worker",
  114. "worker_trace_id": "worker-trace",
  115. "root_trace_id": "root",
  116. "task_id": "task",
  117. "spec_version": 1,
  118. "attempt_id": "attempt",
  119. "task_spec": {},
  120. "continue_trace_id": None,
  121. })
  122. assert result.error == "original runner failure"
  123. assert result.execution_stats.failure_code == FailureCode.EXECUTOR_ERROR
  124. @pytest.mark.asyncio
  125. async def test_worker_protects_operation_epoch_and_deadline():
  126. class Runner:
  127. def __init__(self):
  128. self.config = None
  129. async def run_result(self, *, messages, config):
  130. self.config = config
  131. return {"status": "completed", "summary": "done"}
  132. runner = Runner()
  133. executor = LocalAgentExecutor(runner)
  134. result = await executor.run_worker({
  135. "worker_preset": "worker",
  136. "worker_trace_id": "worker-trace",
  137. "root_trace_id": "root",
  138. "task_id": "task",
  139. "spec_version": 2,
  140. "attempt_id": "attempt",
  141. "task_spec": {"objective": "test"},
  142. "continue_trace_id": None,
  143. "operation_id": "operation",
  144. "execution_epoch": 7,
  145. "deadline": "2030-01-01T00:00:00+00:00",
  146. "untrusted": "must-not-leak",
  147. })
  148. assert result.trace_id == "worker-trace"
  149. assert result.status == "completed"
  150. assert runner.config.trace_id is None
  151. assert runner.config.new_trace_id == "worker-trace"
  152. protected = dict(runner.config.context)
  153. packing = protected.pop("prompt_packing")
  154. assert protected == {
  155. "root_trace_id": "root",
  156. "task_id": "task",
  157. "spec_version": 2,
  158. "attempt_id": "attempt",
  159. "completion_policy": "explicit_validation",
  160. "operation_id": "operation",
  161. "execution_epoch": 7,
  162. "deadline": "2030-01-01T00:00:00+00:00",
  163. }
  164. assert packing == {
  165. "mode": "full",
  166. "estimated_tokens": 110,
  167. "token_limit": 32_000,
  168. "output_reserve": 8_192,
  169. "deduplicated_segments": 0,
  170. }
  171. @pytest.mark.asyncio
  172. async def test_validator_omits_absent_v2_context_and_normalizes_empty_result():
  173. class Runner:
  174. async def run_result(self, *, messages, config):
  175. self.config = config
  176. return {}
  177. runner = Runner()
  178. result = await LocalAgentExecutor(runner).run_validator({
  179. "validator_preset": "validator",
  180. "validator_trace_id": "validator-trace",
  181. "root_trace_id": "root",
  182. "task_id": "task",
  183. "spec_version": 1,
  184. "attempt_id": "attempt",
  185. "snapshot_id": "snapshot",
  186. "validation_id": "validation",
  187. "task_spec": {},
  188. "artifact_snapshot": {},
  189. })
  190. assert result.trace_id == "validator-trace"
  191. assert result.status == "failed"
  192. assert result.summary == ""
  193. assert result.error is None
  194. assert runner.config.trace_id is None
  195. assert runner.config.new_trace_id == "validator-trace"
  196. assert "operation_id" not in runner.config.context
  197. assert "execution_epoch" not in runner.config.context
  198. assert "deadline" not in runner.config.context
  199. @pytest.mark.parametrize(
  200. "kwargs",
  201. [
  202. {"primary_model": " "},
  203. {"primary_model": 123},
  204. {"total_tokens": -1},
  205. {"total_tokens": True},
  206. {"total_cost": -0.1},
  207. {"total_cost": float("nan")},
  208. {"total_cost": float("inf")},
  209. {"total_cost": True},
  210. {"failure_code": "unknown"},
  211. {"failure_code": 123},
  212. ],
  213. )
  214. def test_execution_stats_reject_invalid_values(kwargs):
  215. with pytest.raises(ValueError, match="ExecutionStats"):
  216. ExecutionStats(**kwargs)
  217. def test_execution_stats_normalizes_known_failure_code_string():
  218. assert ExecutionStats(failure_code="timeout").failure_code == FailureCode.TIMEOUT
  219. def test_old_stage_without_stats_roundtrips_and_duration_is_derived():
  220. attempt = TaskAttempt.from_dict({
  221. "attempt_id": "attempt",
  222. "task_id": "task",
  223. "spec_version": 1,
  224. "worker_trace_id": "trace",
  225. "worker_preset": "worker",
  226. "execution_mode": "new",
  227. "started_at": "2026-07-18T00:00:00Z",
  228. "completed_at": "2026-07-18T00:00:01.250+00:00",
  229. })
  230. assert attempt.execution_stats is None
  231. assert attempt.duration_ms == 1250
  232. assert replace(attempt, completed_at="2026-07-17T23:59:59Z").duration_ms == 0
  233. assert replace(attempt, completed_at=None).duration_ms is None
  234. @pytest.mark.asyncio
  235. async def test_repair_continuation_records_only_trace_usage_delta():
  236. trace = SimpleNamespace(
  237. context={}, total_tokens=100, total_cost=1.0, model="repair-model"
  238. )
  239. class Coordinator:
  240. async def validate_continue_from(self, *args):
  241. return "trace"
  242. class Store:
  243. async def get_trace(self, trace_id):
  244. return trace
  245. async def update_trace(self, trace_id, **updates):
  246. trace.context = updates["context"]
  247. class Runner:
  248. task_coordinator = Coordinator()
  249. trace_store = Store()
  250. async def run_result(self, **kwargs):
  251. trace.total_tokens = 130
  252. trace.total_cost = 1.4
  253. return {
  254. "trace_id": "trace",
  255. "status": "completed",
  256. "stats": {"total_tokens": 130, "total_cost": 1.4},
  257. }
  258. result = await LocalAgentExecutor(Runner()).run_worker({
  259. "worker_preset": "worker",
  260. "worker_trace_id": "trace",
  261. "root_trace_id": "root",
  262. "task_id": "task",
  263. "spec_version": 1,
  264. "attempt_id": "new-attempt",
  265. "prior_attempt_id": "old-attempt",
  266. "task_spec": {},
  267. "continue_trace_id": "trace",
  268. })
  269. assert result.execution_stats.primary_model == "repair-model"
  270. assert result.execution_stats.total_tokens == 30
  271. assert result.execution_stats.total_cost == pytest.approx(0.4)
  272. class StatsExecutor(FakeExecutor):
  273. async def run_worker(self, context):
  274. result = await super().run_worker(context)
  275. return replace(
  276. result,
  277. execution_stats=ExecutionStats(
  278. primary_model="worker-model", total_tokens=12, total_cost=0.12
  279. ),
  280. )
  281. async def run_validator(self, context):
  282. result = await super().run_validator(context)
  283. return replace(
  284. result,
  285. execution_stats=ExecutionStats(
  286. primary_model="validator-model", total_tokens=7, total_cost=0.07
  287. ),
  288. )
  289. @pytest.mark.asyncio
  290. async def test_successful_stage_stats_persist_once_and_failed_verdict_is_not_run_failure(
  291. tmp_path,
  292. ):
  293. executor = StatsExecutor([ValidationVerdict.FAILED])
  294. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  295. task_id = await create_task(coordinator, "persist execution stats")
  296. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  297. ledger = await store.load("root")
  298. attempt = ledger.attempts[result.attempt_id]
  299. validation = ledger.validations[result.validation_id]
  300. assert attempt.execution_stats.primary_model == "worker-model"
  301. assert attempt.execution_stats.total_tokens == 12
  302. assert validation.verdict == ValidationVerdict.FAILED
  303. assert validation.execution_stats.primary_model == "validator-model"
  304. assert validation.execution_stats.failure_code is None
  305. assert attempt.duration_ms is not None
  306. assert validation.duration_ms is not None
  307. @pytest.mark.asyncio
  308. async def test_completed_worker_without_submit_is_protocol_failure(tmp_path):
  309. executor = StatsExecutor([], submit_worker=False)
  310. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  311. task_id = await create_task(coordinator, "protocol failure")
  312. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  313. attempt = (await store.load("root")).attempts[result.attempt_id]
  314. assert attempt.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
  315. assert attempt.execution_stats.total_tokens == 12
  316. @pytest.mark.asyncio
  317. async def test_protocol_failure_overrides_executor_supplied_failure_code(tmp_path):
  318. class MisclassifiedExecutor(StatsExecutor):
  319. async def run_worker(self, context):
  320. result = await super().run_worker(context)
  321. return replace(
  322. result,
  323. execution_stats=replace(
  324. result.execution_stats,
  325. failure_code=FailureCode.AGENT_FAILED,
  326. ),
  327. )
  328. executor = MisclassifiedExecutor([], submit_worker=False)
  329. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  330. task_id = await create_task(coordinator, "protocol code wins")
  331. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  332. attempt = (await store.load("root")).attempts[result.attempt_id]
  333. assert attempt.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
  334. @pytest.mark.asyncio
  335. async def test_structured_worker_failure_persists_and_feeds_next_attempt(tmp_path):
  336. failure = FailureDetail(
  337. code="INPUT_SCOPE_MISMATCH",
  338. message="Paragraph has no covering Structure",
  339. disposition=FailureDisposition.REPLAN_TASK,
  340. source_tool="save_structured_script_candidate",
  341. details={"scope": "paragraph/3"},
  342. )
  343. class StructuredFailureExecutor:
  344. def __init__(self):
  345. self.contexts = []
  346. self.coordinator = None
  347. async def run_worker(self, context):
  348. self.contexts.append(context)
  349. return WorkerRunResult(
  350. trace_id=context["worker_trace_id"],
  351. status="failed",
  352. summary="save failed before submit",
  353. error=f"{failure.code}: {failure.message}",
  354. failure=failure,
  355. execution_stats=ExecutionStats(
  356. primary_model="worker-model",
  357. total_tokens=9,
  358. failure_code=FailureCode.TOOL_FAILURE,
  359. ),
  360. )
  361. async def run_validator(self, _context):
  362. raise AssertionError("validator must not run")
  363. async def stop(self, _trace_id):
  364. return True
  365. executor = StructuredFailureExecutor()
  366. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  367. task_id = await create_task(coordinator, "structured failure")
  368. first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  369. ledger = await store.load("root")
  370. first_attempt = ledger.attempts[first.attempt_id]
  371. assert first.failure == failure
  372. assert first_attempt.failure == failure
  373. assert first_attempt.worker_summary == "save failed before submit"
  374. assert first_attempt.execution_stats.failure_code == FailureCode.TOOL_FAILURE
  375. await coordinator.decide_task(
  376. "root",
  377. task_id,
  378. None,
  379. DecisionAction.RETRY,
  380. {"reason": "retry after inspecting failure"},
  381. "retry-structured-failure",
  382. )
  383. await coordinator.dispatch_tasks("root", [task_id])
  384. feedback = executor.contexts[1]["prior_feedback"]
  385. assert feedback["attempt"]["failure"] == failure.to_dict()
  386. assert feedback["attempt"]["worker_summary"] == "save failed before submit"
  387. assert feedback["attempt"]["artifact_submitted"] is False
  388. assert feedback.get("validation") is None
  389. @pytest.mark.asyncio
  390. async def test_local_executor_uses_structured_failure_code_instead_of_agent_failed():
  391. failure = FailureDetail(
  392. code="INPUT_SCOPE_MISMATCH",
  393. message="revise contract",
  394. disposition=FailureDisposition.REPLAN_TASK,
  395. )
  396. class Runner:
  397. trace_store = None
  398. async def run_result(self, **_kwargs):
  399. return {
  400. "status": "failed",
  401. "summary": "",
  402. "error": "protocol_violation",
  403. "failure": failure.to_dict(),
  404. }
  405. result = await LocalAgentExecutor(Runner()).run_worker(
  406. {
  407. "worker_preset": "worker",
  408. "worker_trace_id": "worker-trace",
  409. "root_trace_id": "root",
  410. "task_id": "task",
  411. "spec_version": 1,
  412. "attempt_id": "attempt",
  413. "task_spec": {},
  414. "continue_trace_id": None,
  415. }
  416. )
  417. assert result.failure == failure
  418. assert result.execution_stats.failure_code == FailureCode.TOOL_FAILURE
  419. @pytest.mark.asyncio
  420. async def test_late_failure_cannot_overwrite_concurrent_submission(tmp_path):
  421. executor = FakeExecutor([])
  422. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  423. task_id = await create_task(coordinator, "concurrent submission")
  424. operation = await coordinator.start_operation(
  425. "root", "dispatch", task_ids=[task_id]
  426. )
  427. completed = await coordinator.await_operation("root", operation.operation_id)
  428. ledger = await store.load("root")
  429. attempt_id = completed.attempt_ids[0]
  430. original = ledger.attempts[attempt_id]
  431. await coordinator._mark_worker_failure(
  432. "root",
  433. task_id,
  434. attempt_id,
  435. "failed",
  436. "late failure",
  437. ExecutionStats(failure_code=FailureCode.AGENT_FAILED),
  438. operation.operation_id,
  439. completed.execution_epoch,
  440. )
  441. after = (await store.load("root")).attempts[attempt_id]
  442. assert after.status == original.status
  443. assert after.error == original.error
  444. assert after.execution_stats == original.execution_stats
  445. @pytest.mark.asyncio
  446. async def test_timeout_persists_machine_failure_code(tmp_path):
  447. class SlowExecutor(FakeExecutor):
  448. async def run_worker(self, context):
  449. await asyncio.sleep(1)
  450. return WorkerRunResult(context["worker_trace_id"], "completed")
  451. executor = SlowExecutor([])
  452. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  453. coordinator.config = OrchestrationConfig(worker_timeout_seconds=0.001)
  454. task_id = await create_task(coordinator, "timeout stats")
  455. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  456. attempt = (await store.load("root")).attempts[result.attempt_id]
  457. assert attempt.status == AttemptStatus.EXPIRED
  458. assert attempt.execution_stats.failure_code == FailureCode.TIMEOUT
  459. @pytest.mark.asyncio
  460. async def test_remote_stop_wins_over_late_worker_failure(tmp_path):
  461. class LateFailureExecutor:
  462. def __init__(self):
  463. self.started = asyncio.Event()
  464. self.release = asyncio.Event()
  465. async def run_worker(self, context):
  466. self.started.set()
  467. await self.release.wait()
  468. return WorkerRunResult(
  469. context["worker_trace_id"],
  470. "failed",
  471. error="late failure",
  472. execution_stats=ExecutionStats(
  473. total_tokens=3, failure_code=FailureCode.AGENT_FAILED
  474. ),
  475. )
  476. async def run_validator(self, context):
  477. raise AssertionError("validator must not run")
  478. async def stop(self, trace_id):
  479. return True
  480. executor = LateFailureExecutor()
  481. coordinator_a, store, trace_store = await make_coordinator(tmp_path, executor)
  482. task_id = await create_task(coordinator_a, "remote stop stats")
  483. operation = await coordinator_a.start_operation("root", "dispatch", task_ids=[task_id])
  484. await asyncio.wait_for(executor.started.wait(), timeout=1)
  485. observer_executor = FakeExecutor([])
  486. coordinator_b = TaskCoordinator(
  487. store,
  488. FileSystemArtifactStore(str(tmp_path)),
  489. trace_store,
  490. OrchestrationConfig(stop_grace_seconds=0),
  491. TraceEventSink(str(tmp_path)),
  492. observer_executor,
  493. )
  494. observer_executor.coordinator = coordinator_b
  495. await coordinator_b.stop_operation("root", operation.operation_id)
  496. executor.release.set()
  497. await coordinator_a.runtime.wait("root", operation.operation_id)
  498. ledger = await store.load("root")
  499. attempt = ledger.attempts[ledger.operations[operation.operation_id].attempt_ids[0]]
  500. assert attempt.status == AttemptStatus.STOPPED
  501. assert attempt.execution_stats.failure_code == FailureCode.STOPPED
  502. assert attempt.error == "Operation stopped"
  503. assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
  504. @pytest.mark.asyncio
  505. async def test_cas_retry_does_not_overwrite_concurrent_attempt_submission(tmp_path):
  506. class ConcurrentSubmissionStore:
  507. def __init__(self, inner):
  508. self.inner = inner
  509. self.base_path = inner.base_path
  510. self.injected = False
  511. async def load(self, root_trace_id):
  512. return await self.inner.load(root_trace_id)
  513. async def commit(
  514. self,
  515. ledger,
  516. expected_revision,
  517. idempotency_key=None,
  518. event=None,
  519. ):
  520. if (
  521. event is not None
  522. and event.event_type == "worker_failed"
  523. and not self.injected
  524. ):
  525. self.injected = True
  526. winner = await self.inner.load(ledger.root_trace_id)
  527. attempt = next(
  528. item
  529. for item in winner.attempts.values()
  530. if item.status == AttemptStatus.RUNNING
  531. )
  532. attempt.status = AttemptStatus.SUBMITTED
  533. attempt.submission = AttemptSubmission(summary="concurrent winner")
  534. attempt.completed_at = utc_now()
  535. winner.tasks[attempt.task_id].status = TaskStatus.AWAITING_VALIDATION
  536. await self.inner.commit(winner, winner.revision)
  537. return await self.inner.commit(
  538. ledger,
  539. expected_revision,
  540. idempotency_key=idempotency_key,
  541. event=event,
  542. )
  543. async def list_events(self, *args, **kwargs):
  544. return await self.inner.list_events(*args, **kwargs)
  545. async def list_recoverable(self, *args, **kwargs):
  546. return await self.inner.list_recoverable(*args, **kwargs)
  547. task_store = ConcurrentSubmissionStore(FileSystemTaskStore(str(tmp_path)))
  548. executor = StatsExecutor([], submit_worker=False)
  549. coordinator, store, _ = await make_coordinator(
  550. tmp_path, executor, task_store=task_store
  551. )
  552. task_id = await create_task(coordinator, "CAS submit winner")
  553. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  554. ledger = await store.load("root")
  555. attempt = ledger.attempts[result.attempt_id]
  556. assert task_store.injected is True
  557. assert attempt.status == AttemptStatus.SUBMITTED
  558. assert attempt.submission.summary == "concurrent winner"
  559. assert attempt.error is None
  560. assert attempt.execution_stats is None
  561. assert ledger.tasks[task_id].status == TaskStatus.AWAITING_VALIDATION