test_orchestration_v2_control.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666
  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. assert runner.config.context == {
  153. "root_trace_id": "root",
  154. "task_id": "task",
  155. "spec_version": 2,
  156. "attempt_id": "attempt",
  157. "completion_policy": "explicit_validation",
  158. "operation_id": "operation",
  159. "execution_epoch": 7,
  160. "deadline": "2030-01-01T00:00:00+00:00",
  161. }
  162. @pytest.mark.asyncio
  163. async def test_validator_omits_absent_v2_context_and_normalizes_empty_result():
  164. class Runner:
  165. async def run_result(self, *, messages, config):
  166. self.config = config
  167. return {}
  168. runner = Runner()
  169. result = await LocalAgentExecutor(runner).run_validator({
  170. "validator_preset": "validator",
  171. "validator_trace_id": "validator-trace",
  172. "root_trace_id": "root",
  173. "task_id": "task",
  174. "spec_version": 1,
  175. "attempt_id": "attempt",
  176. "snapshot_id": "snapshot",
  177. "validation_id": "validation",
  178. "task_spec": {},
  179. "artifact_snapshot": {},
  180. })
  181. assert result.trace_id == "validator-trace"
  182. assert result.status == "failed"
  183. assert result.summary == ""
  184. assert result.error is None
  185. assert runner.config.trace_id is None
  186. assert runner.config.new_trace_id == "validator-trace"
  187. assert "operation_id" not in runner.config.context
  188. assert "execution_epoch" not in runner.config.context
  189. assert "deadline" not in runner.config.context
  190. @pytest.mark.parametrize(
  191. "kwargs",
  192. [
  193. {"primary_model": " "},
  194. {"primary_model": 123},
  195. {"total_tokens": -1},
  196. {"total_tokens": True},
  197. {"total_cost": -0.1},
  198. {"total_cost": float("nan")},
  199. {"total_cost": float("inf")},
  200. {"total_cost": True},
  201. {"failure_code": "unknown"},
  202. {"failure_code": 123},
  203. ],
  204. )
  205. def test_execution_stats_reject_invalid_values(kwargs):
  206. with pytest.raises(ValueError, match="ExecutionStats"):
  207. ExecutionStats(**kwargs)
  208. def test_execution_stats_normalizes_known_failure_code_string():
  209. assert ExecutionStats(failure_code="timeout").failure_code == FailureCode.TIMEOUT
  210. def test_old_stage_without_stats_roundtrips_and_duration_is_derived():
  211. attempt = TaskAttempt.from_dict({
  212. "attempt_id": "attempt",
  213. "task_id": "task",
  214. "spec_version": 1,
  215. "worker_trace_id": "trace",
  216. "worker_preset": "worker",
  217. "execution_mode": "new",
  218. "started_at": "2026-07-18T00:00:00Z",
  219. "completed_at": "2026-07-18T00:00:01.250+00:00",
  220. })
  221. assert attempt.execution_stats is None
  222. assert attempt.duration_ms == 1250
  223. assert replace(attempt, completed_at="2026-07-17T23:59:59Z").duration_ms == 0
  224. assert replace(attempt, completed_at=None).duration_ms is None
  225. @pytest.mark.asyncio
  226. async def test_repair_continuation_records_only_trace_usage_delta():
  227. trace = SimpleNamespace(
  228. context={}, total_tokens=100, total_cost=1.0, model="repair-model"
  229. )
  230. class Coordinator:
  231. async def validate_continue_from(self, *args):
  232. return "trace"
  233. class Store:
  234. async def get_trace(self, trace_id):
  235. return trace
  236. async def update_trace(self, trace_id, **updates):
  237. trace.context = updates["context"]
  238. class Runner:
  239. task_coordinator = Coordinator()
  240. trace_store = Store()
  241. async def run_result(self, **kwargs):
  242. trace.total_tokens = 130
  243. trace.total_cost = 1.4
  244. return {
  245. "trace_id": "trace",
  246. "status": "completed",
  247. "stats": {"total_tokens": 130, "total_cost": 1.4},
  248. }
  249. result = await LocalAgentExecutor(Runner()).run_worker({
  250. "worker_preset": "worker",
  251. "worker_trace_id": "trace",
  252. "root_trace_id": "root",
  253. "task_id": "task",
  254. "spec_version": 1,
  255. "attempt_id": "new-attempt",
  256. "prior_attempt_id": "old-attempt",
  257. "task_spec": {},
  258. "continue_trace_id": "trace",
  259. })
  260. assert result.execution_stats.primary_model == "repair-model"
  261. assert result.execution_stats.total_tokens == 30
  262. assert result.execution_stats.total_cost == pytest.approx(0.4)
  263. class StatsExecutor(FakeExecutor):
  264. async def run_worker(self, context):
  265. result = await super().run_worker(context)
  266. return replace(
  267. result,
  268. execution_stats=ExecutionStats(
  269. primary_model="worker-model", total_tokens=12, total_cost=0.12
  270. ),
  271. )
  272. async def run_validator(self, context):
  273. result = await super().run_validator(context)
  274. return replace(
  275. result,
  276. execution_stats=ExecutionStats(
  277. primary_model="validator-model", total_tokens=7, total_cost=0.07
  278. ),
  279. )
  280. @pytest.mark.asyncio
  281. async def test_successful_stage_stats_persist_once_and_failed_verdict_is_not_run_failure(
  282. tmp_path,
  283. ):
  284. executor = StatsExecutor([ValidationVerdict.FAILED])
  285. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  286. task_id = await create_task(coordinator, "persist execution stats")
  287. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  288. ledger = await store.load("root")
  289. attempt = ledger.attempts[result.attempt_id]
  290. validation = ledger.validations[result.validation_id]
  291. assert attempt.execution_stats.primary_model == "worker-model"
  292. assert attempt.execution_stats.total_tokens == 12
  293. assert validation.verdict == ValidationVerdict.FAILED
  294. assert validation.execution_stats.primary_model == "validator-model"
  295. assert validation.execution_stats.failure_code is None
  296. assert attempt.duration_ms is not None
  297. assert validation.duration_ms is not None
  298. @pytest.mark.asyncio
  299. async def test_completed_worker_without_submit_is_protocol_failure(tmp_path):
  300. executor = StatsExecutor([], submit_worker=False)
  301. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  302. task_id = await create_task(coordinator, "protocol failure")
  303. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  304. attempt = (await store.load("root")).attempts[result.attempt_id]
  305. assert attempt.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
  306. assert attempt.execution_stats.total_tokens == 12
  307. @pytest.mark.asyncio
  308. async def test_protocol_failure_overrides_executor_supplied_failure_code(tmp_path):
  309. class MisclassifiedExecutor(StatsExecutor):
  310. async def run_worker(self, context):
  311. result = await super().run_worker(context)
  312. return replace(
  313. result,
  314. execution_stats=replace(
  315. result.execution_stats,
  316. failure_code=FailureCode.AGENT_FAILED,
  317. ),
  318. )
  319. executor = MisclassifiedExecutor([], submit_worker=False)
  320. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  321. task_id = await create_task(coordinator, "protocol code wins")
  322. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  323. attempt = (await store.load("root")).attempts[result.attempt_id]
  324. assert attempt.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
  325. @pytest.mark.asyncio
  326. async def test_structured_worker_failure_persists_and_feeds_next_attempt(tmp_path):
  327. failure = FailureDetail(
  328. code="INPUT_SCOPE_MISMATCH",
  329. message="Paragraph has no covering Structure",
  330. disposition=FailureDisposition.REPLAN_TASK,
  331. source_tool="save_structured_script_candidate",
  332. details={"scope": "paragraph/3"},
  333. )
  334. class StructuredFailureExecutor:
  335. def __init__(self):
  336. self.contexts = []
  337. self.coordinator = None
  338. async def run_worker(self, context):
  339. self.contexts.append(context)
  340. return WorkerRunResult(
  341. trace_id=context["worker_trace_id"],
  342. status="failed",
  343. summary="save failed before submit",
  344. error=f"{failure.code}: {failure.message}",
  345. failure=failure,
  346. execution_stats=ExecutionStats(
  347. primary_model="worker-model",
  348. total_tokens=9,
  349. failure_code=FailureCode.TOOL_FAILURE,
  350. ),
  351. )
  352. async def run_validator(self, _context):
  353. raise AssertionError("validator must not run")
  354. async def stop(self, _trace_id):
  355. return True
  356. executor = StructuredFailureExecutor()
  357. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  358. task_id = await create_task(coordinator, "structured failure")
  359. first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  360. ledger = await store.load("root")
  361. first_attempt = ledger.attempts[first.attempt_id]
  362. assert first.failure == failure
  363. assert first_attempt.failure == failure
  364. assert first_attempt.worker_summary == "save failed before submit"
  365. assert first_attempt.execution_stats.failure_code == FailureCode.TOOL_FAILURE
  366. await coordinator.decide_task(
  367. "root",
  368. task_id,
  369. None,
  370. DecisionAction.RETRY,
  371. {"reason": "retry after inspecting failure"},
  372. "retry-structured-failure",
  373. )
  374. await coordinator.dispatch_tasks("root", [task_id])
  375. feedback = executor.contexts[1]["prior_feedback"]
  376. assert feedback["attempt"]["failure"] == failure.to_dict()
  377. assert feedback["attempt"]["worker_summary"] == "save failed before submit"
  378. assert feedback["attempt"]["artifact_submitted"] is False
  379. assert feedback.get("validation") is None
  380. @pytest.mark.asyncio
  381. async def test_local_executor_uses_structured_failure_code_instead_of_agent_failed():
  382. failure = FailureDetail(
  383. code="INPUT_SCOPE_MISMATCH",
  384. message="revise contract",
  385. disposition=FailureDisposition.REPLAN_TASK,
  386. )
  387. class Runner:
  388. trace_store = None
  389. async def run_result(self, **_kwargs):
  390. return {
  391. "status": "failed",
  392. "summary": "",
  393. "error": "protocol_violation",
  394. "failure": failure.to_dict(),
  395. }
  396. result = await LocalAgentExecutor(Runner()).run_worker(
  397. {
  398. "worker_preset": "worker",
  399. "worker_trace_id": "worker-trace",
  400. "root_trace_id": "root",
  401. "task_id": "task",
  402. "spec_version": 1,
  403. "attempt_id": "attempt",
  404. "task_spec": {},
  405. "continue_trace_id": None,
  406. }
  407. )
  408. assert result.failure == failure
  409. assert result.execution_stats.failure_code == FailureCode.TOOL_FAILURE
  410. @pytest.mark.asyncio
  411. async def test_late_failure_cannot_overwrite_concurrent_submission(tmp_path):
  412. executor = FakeExecutor([])
  413. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  414. task_id = await create_task(coordinator, "concurrent submission")
  415. operation = await coordinator.start_operation(
  416. "root", "dispatch", task_ids=[task_id]
  417. )
  418. completed = await coordinator.await_operation("root", operation.operation_id)
  419. ledger = await store.load("root")
  420. attempt_id = completed.attempt_ids[0]
  421. original = ledger.attempts[attempt_id]
  422. await coordinator._mark_worker_failure(
  423. "root",
  424. task_id,
  425. attempt_id,
  426. "failed",
  427. "late failure",
  428. ExecutionStats(failure_code=FailureCode.AGENT_FAILED),
  429. operation.operation_id,
  430. completed.execution_epoch,
  431. )
  432. after = (await store.load("root")).attempts[attempt_id]
  433. assert after.status == original.status
  434. assert after.error == original.error
  435. assert after.execution_stats == original.execution_stats
  436. @pytest.mark.asyncio
  437. async def test_timeout_persists_machine_failure_code(tmp_path):
  438. class SlowExecutor(FakeExecutor):
  439. async def run_worker(self, context):
  440. await asyncio.sleep(1)
  441. return WorkerRunResult(context["worker_trace_id"], "completed")
  442. executor = SlowExecutor([])
  443. coordinator, store, _ = await make_coordinator(tmp_path, executor)
  444. coordinator.config = OrchestrationConfig(worker_timeout_seconds=0.001)
  445. task_id = await create_task(coordinator, "timeout stats")
  446. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  447. attempt = (await store.load("root")).attempts[result.attempt_id]
  448. assert attempt.status == AttemptStatus.EXPIRED
  449. assert attempt.execution_stats.failure_code == FailureCode.TIMEOUT
  450. @pytest.mark.asyncio
  451. async def test_remote_stop_wins_over_late_worker_failure(tmp_path):
  452. class LateFailureExecutor:
  453. def __init__(self):
  454. self.started = asyncio.Event()
  455. self.release = asyncio.Event()
  456. async def run_worker(self, context):
  457. self.started.set()
  458. await self.release.wait()
  459. return WorkerRunResult(
  460. context["worker_trace_id"],
  461. "failed",
  462. error="late failure",
  463. execution_stats=ExecutionStats(
  464. total_tokens=3, failure_code=FailureCode.AGENT_FAILED
  465. ),
  466. )
  467. async def run_validator(self, context):
  468. raise AssertionError("validator must not run")
  469. async def stop(self, trace_id):
  470. return True
  471. executor = LateFailureExecutor()
  472. coordinator_a, store, trace_store = await make_coordinator(tmp_path, executor)
  473. task_id = await create_task(coordinator_a, "remote stop stats")
  474. operation = await coordinator_a.start_operation("root", "dispatch", task_ids=[task_id])
  475. await asyncio.wait_for(executor.started.wait(), timeout=1)
  476. observer_executor = FakeExecutor([])
  477. coordinator_b = TaskCoordinator(
  478. store,
  479. FileSystemArtifactStore(str(tmp_path)),
  480. trace_store,
  481. OrchestrationConfig(stop_grace_seconds=0),
  482. TraceEventSink(str(tmp_path)),
  483. observer_executor,
  484. )
  485. observer_executor.coordinator = coordinator_b
  486. await coordinator_b.stop_operation("root", operation.operation_id)
  487. executor.release.set()
  488. await coordinator_a.runtime.wait("root", operation.operation_id)
  489. ledger = await store.load("root")
  490. attempt = ledger.attempts[ledger.operations[operation.operation_id].attempt_ids[0]]
  491. assert attempt.status == AttemptStatus.STOPPED
  492. assert attempt.execution_stats.failure_code == FailureCode.STOPPED
  493. assert attempt.error == "Operation stopped"
  494. assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
  495. @pytest.mark.asyncio
  496. async def test_cas_retry_does_not_overwrite_concurrent_attempt_submission(tmp_path):
  497. class ConcurrentSubmissionStore:
  498. def __init__(self, inner):
  499. self.inner = inner
  500. self.base_path = inner.base_path
  501. self.injected = False
  502. async def load(self, root_trace_id):
  503. return await self.inner.load(root_trace_id)
  504. async def commit(
  505. self,
  506. ledger,
  507. expected_revision,
  508. idempotency_key=None,
  509. event=None,
  510. ):
  511. if (
  512. event is not None
  513. and event.event_type == "worker_failed"
  514. and not self.injected
  515. ):
  516. self.injected = True
  517. winner = await self.inner.load(ledger.root_trace_id)
  518. attempt = next(
  519. item
  520. for item in winner.attempts.values()
  521. if item.status == AttemptStatus.RUNNING
  522. )
  523. attempt.status = AttemptStatus.SUBMITTED
  524. attempt.submission = AttemptSubmission(summary="concurrent winner")
  525. attempt.completed_at = utc_now()
  526. winner.tasks[attempt.task_id].status = TaskStatus.AWAITING_VALIDATION
  527. await self.inner.commit(winner, winner.revision)
  528. return await self.inner.commit(
  529. ledger,
  530. expected_revision,
  531. idempotency_key=idempotency_key,
  532. event=event,
  533. )
  534. async def list_events(self, *args, **kwargs):
  535. return await self.inner.list_events(*args, **kwargs)
  536. async def list_recoverable(self, *args, **kwargs):
  537. return await self.inner.list_recoverable(*args, **kwargs)
  538. task_store = ConcurrentSubmissionStore(FileSystemTaskStore(str(tmp_path)))
  539. executor = StatsExecutor([], submit_worker=False)
  540. coordinator, store, _ = await make_coordinator(
  541. tmp_path, executor, task_store=task_store
  542. )
  543. task_id = await create_task(coordinator, "CAS submit winner")
  544. result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
  545. ledger = await store.load("root")
  546. attempt = ledger.attempts[result.attempt_id]
  547. assert task_store.injected is True
  548. assert attempt.status == AttemptStatus.SUBMITTED
  549. assert attempt.submission.summary == "concurrent winner"
  550. assert attempt.error is None
  551. assert attempt.execution_stats is None
  552. assert ledger.tasks[task_id].status == TaskStatus.AWAITING_VALIDATION