import asyncio from collections import deque import pytest from agent.orchestration.config import OrchestrationConfig from agent.orchestration.coordinator import TaskConflict, TaskCoordinator from agent.orchestration.models import ( AgentRole, ArtifactRef, AttemptSubmission, CriterionResult, DecisionAction, OperationKind, OperationStatus, TaskStatus, ValidationVerdict, ) from agent.orchestration.protocols import ValidatorRunResult, WorkerRunResult from agent.orchestration.store import FileSystemArtifactStore, FileSystemTaskStore, TraceEventSink from agent.trace.models import Trace from agent.trace.store import FileSystemTraceStore from agent.core.runner import AgentRunner from agent.orchestration.wiring import wire_orchestration ROOT_TASK_SPEC = { "objective": "mission", "acceptance_criteria": [ {"criterion_id": "mission-done", "description": "mission is complete", "hard": True} ], } class FakeExecutor: def __init__(self, verdicts, submit_worker=True, submit_validator=True, delay=0): self.verdicts = deque(verdicts) self.submit_worker = submit_worker self.submit_validator = submit_validator self.delay = delay self.coordinator = None self.active = 0 self.max_active = 0 self.worker_calls = 0 self.validator_calls = 0 async def run_worker(self, context): self.worker_calls += 1 self.active += 1 self.max_active = max(self.max_active, self.active) if self.delay: await asyncio.sleep(self.delay) try: if self.submit_worker: await self.coordinator.submit_attempt( { **context, "role": AgentRole.WORKER.value, "trace_id": context["worker_trace_id"], "tool_call_id": f"submit-{context['attempt_id']}", }, AttemptSubmission( summary="done", artifact_refs=[ArtifactRef(uri=f"memory://{context['attempt_id']}", version="1")], ), ) return WorkerRunResult(context["worker_trace_id"], "completed") finally: self.active -= 1 async def run_validator(self, context): self.validator_calls += 1 verdict = self.verdicts.popleft() if self.verdicts else ValidationVerdict.PASSED if self.submit_validator: criteria = [ CriterionResult( criterion_id=item["criterion_id"], verdict=verdict, reason="checked", ) for item in context["task_spec"]["acceptance_criteria"] ] await self.coordinator.submit_validation( { **context, "role": AgentRole.VALIDATOR.value, "trace_id": context["validator_trace_id"], "tool_call_id": f"validate-{context['validation_id']}", }, verdict, criteria, "independent report", [], [], [], "accept" if verdict == ValidationVerdict.PASSED else "replan", ) return ValidatorRunResult(context["validator_trace_id"], "completed") async def make_coordinator( tmp_path, executor, artifact_store=None, trace_store=None, task_store=None, ): trace_store = trace_store or FileSystemTraceStore(str(tmp_path)) await trace_store.create_trace(Trace(trace_id="root", mode="agent", task="mission", agent_role="planner")) task_store = task_store or FileSystemTaskStore(str(tmp_path)) coordinator = TaskCoordinator( task_store, artifact_store or FileSystemArtifactStore(str(tmp_path)), trace_store, OrchestrationConfig(max_parallel_tasks=4), TraceEventSink(str(tmp_path)), executor, ) executor.coordinator = coordinator await coordinator.ensure_ledger("root", ROOT_TASK_SPEC) return coordinator, task_store, trace_store async def create_task(coordinator, objective="task", parent_task_id=None): result = await coordinator.create_tasks( "root", [{ "objective": objective, "acceptance_criteria": [{"criterion_id": "c1", "description": "must pass", "hard": True}], }], parent_task_id=parent_task_id, ) return result["tasks"][0]["task_id"] @pytest.mark.asyncio async def test_passed_validation_requires_planner_accept(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0] assert cycle.task_status == TaskStatus.AWAITING_DECISION ledger = await store.load("root") assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION await coordinator.decide_task( "root", task_id, cycle.validation_id, DecisionAction.ACCEPT, {"reason": "all hard criteria passed"}, "decision-1", ) assert (await store.load("root")).tasks[task_id].status == TaskStatus.COMPLETED @pytest.mark.asyncio async def test_failed_or_inconclusive_cannot_be_accepted(tmp_path): executor = FakeExecutor([ValidationVerdict.FAILED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0] with pytest.raises(TaskConflict, match="passed"): await coordinator.decide_task( "root", task_id, cycle.validation_id, DecisionAction.ACCEPT, {"reason": "override"}, "bad-accept", ) assert (await store.load("root")).tasks[task_id].status == TaskStatus.AWAITING_DECISION @pytest.mark.asyncio async def test_repair_is_limited_to_once_per_spec_version(tmp_path): executor = FakeExecutor([ValidationVerdict.FAILED, ValidationVerdict.FAILED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) first = (await coordinator.dispatch_tasks("root", [task_id]))[0] await coordinator.decide_task( "root", task_id, first.validation_id, DecisionAction.REPAIR, {"reason": "small local correction"}, "repair-1", ) second = (await coordinator.dispatch_tasks("root", [task_id]))[0] ledger = await store.load("root") attempts = [ledger.attempts[x] for x in ledger.tasks[task_id].attempt_ids] assert attempts[0].worker_trace_id == attempts[1].worker_trace_id with pytest.raises(TaskConflict, match="limit"): await coordinator.decide_task( "root", task_id, second.validation_id, DecisionAction.REPAIR, {"reason": "try again"}, "repair-2", ) @pytest.mark.asyncio async def test_worker_without_submit_attempt_needs_replan(tmp_path): executor = FakeExecutor([], submit_worker=False) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) result = (await coordinator.dispatch_tasks("root", [task_id]))[0] ledger = await store.load("root") assert result.task_status == TaskStatus.NEEDS_REPLAN assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN assert ledger.attempts[result.attempt_id].status.value == "failed" @pytest.mark.asyncio async def test_validator_without_submit_is_error_then_revalidates_with_new_trace(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED], submit_validator=False) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) first = (await coordinator.dispatch_tasks("root", [task_id]))[0] ledger = await store.load("root") first_validation = ledger.validations[first.validation_id] assert first.task_status == TaskStatus.NEEDS_REPLAN assert first_validation.status.value == "error" assert first_validation.verdict is None executor.submit_validator = True executor.verdicts.append(ValidationVerdict.PASSED) second = await coordinator.revalidate_attempt( "root", task_id, first.attempt_id, "revalidate-1" ) assert second.task_status == TaskStatus.AWAITING_DECISION assert second.validation_id != first.validation_id assert second.validation.validator_trace_id != first_validation.validator_trace_id @pytest.mark.asyncio async def test_retry_uses_new_trace_and_revise_invalidates_old_validation(tmp_path): executor = FakeExecutor([ ValidationVerdict.FAILED, ValidationVerdict.INCONCLUSIVE, ValidationVerdict.PASSED, ]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator) first = (await coordinator.dispatch_tasks("root", [task_id]))[0] await coordinator.decide_task( "root", task_id, first.validation_id, DecisionAction.RETRY, {"reason": "use a fresh worker"}, "retry-1", ) second = (await coordinator.dispatch_tasks("root", [task_id]))[0] ledger = await store.load("root") first_attempt = ledger.attempts[first.attempt_id] second_attempt = ledger.attempts[second.attempt_id] assert first_attempt.worker_trace_id != second_attempt.worker_trace_id await coordinator.decide_task( "root", task_id, second.validation_id, DecisionAction.REVISE, { "reason": "clarify criterion", "objective": "revised task", "acceptance_criteria": [{"criterion_id": "c1", "description": "revised", "hard": True}], }, "revise-1", ) ledger = await store.load("root") assert ledger.tasks[task_id].current_spec_version == 2 third = (await coordinator.dispatch_tasks("root", [task_id]))[0] with pytest.raises(TaskConflict, match="obsolete"): await coordinator.decide_task( "root", task_id, first.validation_id, DecisionAction.ACCEPT, {"reason": "stale"}, "stale-accept", ) await coordinator.decide_task( "root", task_id, third.validation_id, DecisionAction.ACCEPT, {"reason": "current validation passed"}, "current-accept", ) @pytest.mark.asyncio async def test_block_unblock_cancel_and_supersede(tmp_path): executor = FakeExecutor([]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator, "blocked") await coordinator.decide_task( "root", task_id, None, DecisionAction.BLOCK, {"reason": "external dependency"}, "block-1", ) assert (await store.load("root")).tasks[task_id].blocked_reason == "external dependency" await coordinator.decide_task( "root", task_id, None, DecisionAction.UNBLOCK, {}, "unblock-1" ) assert (await store.load("root")).tasks[task_id].status == TaskStatus.NEEDS_REPLAN await coordinator.decide_task( "root", task_id, None, DecisionAction.CANCEL, {"reason": "no longer needed"}, "cancel-1", ) assert (await store.load("root")).tasks[task_id].status == TaskStatus.CANCELLED old_id = await create_task(coordinator, "old") result = await coordinator.decide_task( "root", old_id, None, DecisionAction.SUPERSEDE, { "reason": "replace spec", "replacement": { "objective": "replacement", "acceptance_criteria": [ {"criterion_id": "replacement-c1", "description": "replacement passes"} ], "context_refs": ["replacement-context"], }, }, "supersede-1", ) replacement_id = result["payload"]["replacement_task_id"] ledger = await store.load("root") assert ledger.tasks[old_id].status == TaskStatus.SUPERSEDED assert ledger.tasks[replacement_id].status == TaskStatus.PENDING replacement_spec = ledger.tasks[replacement_id].current_spec assert replacement_spec.objective == "replacement" assert replacement_spec.acceptance_criteria[0].criterion_id == "replacement-c1" assert replacement_spec.context_refs == ("replacement-context",) @pytest.mark.asyncio async def test_insert_after_keeps_stable_ids_and_updates_display_order(tmp_path): executor = FakeExecutor([]) coordinator, store, _ = await make_coordinator(tmp_path, executor) first = await create_task(coordinator, "first") third = await create_task(coordinator, "third") inserted = await coordinator.create_tasks( "root", [{ "objective": "second", "acceptance_criteria": [ {"criterion_id": "second-c1", "description": "second passes"} ], "context_refs": ["second-context"], }], placement={"after_task_id": first, "focus": True}, idempotency_key="insert", ) second = inserted["tasks"][0]["task_id"] ledger = await store.load("root") assert ledger.tasks[first].display_path == "0.1" assert ledger.tasks[second].display_path == "0.2" assert ledger.tasks[third].display_path == "0.3" assert ledger.tasks[second].current_spec.context_refs == ("second-context",) assert ledger.focused_task_id == second assert len({first, second, third}) == 3 repeated = await coordinator.create_tasks( "root", [{ "objective": "second", "acceptance_criteria": [ {"criterion_id": "second-c1", "description": "second passes"} ], "context_refs": ["second-context"], }], placement={"after_task_id": first, "focus": True}, idempotency_key="insert", ) assert repeated["tasks"][0]["task_id"] == second @pytest.mark.asyncio async def test_split_children_do_not_auto_complete_parent(tmp_path): executor = FakeExecutor([ValidationVerdict.FAILED, ValidationVerdict.PASSED, ValidationVerdict.PASSED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) parent_id = await create_task(coordinator, "parent") first = (await coordinator.dispatch_tasks("root", [parent_id]))[0] decision = await coordinator.decide_task( "root", parent_id, first.validation_id, DecisionAction.SPLIT, { "reason": "split work", "tasks": [ { "objective": "child one", "acceptance_criteria": [{"criterion_id": "c1", "description": "pass"}], "context_refs": ["child-one-context"], }, { "objective": "child two", "acceptance_criteria": [{"criterion_id": "c2", "description": "pass"}], "context_refs": ["child-two-context"], }, ], }, "split-1", ) child_ids = decision["payload"]["child_task_ids"] ledger = await store.load("root") assert ledger.tasks[parent_id].child_task_ids == child_ids assert [ledger.tasks[x].display_path for x in child_ids] == ["0.1.1", "0.1.2"] assert [ledger.tasks[x].current_spec.context_refs for x in child_ids] == [ ("child-one-context",), ("child-two-context",), ] assert [ ledger.tasks[x].current_spec.acceptance_criteria[0].criterion_id for x in child_ids ] == ["c1", "c2"] cycles = await coordinator.dispatch_tasks("root", child_ids) for child_id, cycle in zip(child_ids, cycles): await coordinator.decide_task( "root", child_id, cycle.validation_id, DecisionAction.ACCEPT, {"reason": "passed"}, f"accept-{child_id}", ) ledger = await store.load("root") assert all(ledger.tasks[x].status == TaskStatus.COMPLETED for x in child_ids) assert ledger.tasks[parent_id].status == TaskStatus.NEEDS_REPLAN @pytest.mark.asyncio async def test_parallel_tasks_are_isolated_and_bounded(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED] * 4, delay=0.02) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_ids = [await create_task(coordinator, f"task-{i}") for i in range(4)] results = await coordinator.dispatch_tasks("root", task_ids, idempotency_key="batch") assert [x.task_id for x in results] == task_ids assert executor.max_active <= 4 assert all(x.task_status == TaskStatus.AWAITING_DECISION for x in results) ledger = await store.load("root") assert len({ledger.attempts[ledger.tasks[x].attempt_ids[-1]].worker_trace_id for x in task_ids}) == 4 repeated = await coordinator.dispatch_tasks("root", task_ids, idempotency_key="batch") assert [x.attempt_id for x in repeated] == [x.attempt_id for x in results] @pytest.mark.asyncio async def test_real_runner_local_executor_creates_independent_terminal_traces(tmp_path): import json async def fake_llm(messages, tools, **kwargs): names = {item["function"]["name"] for item in tools or []} if "submit_attempt" in names: arguments = { "summary": "implemented", "artifact_refs": [{"uri": "memory://result", "version": "1"}], "evidence_refs": [], } tool_name = "submit_attempt" elif "submit_validation" in names: arguments = { "verdict": "passed", "criterion_results": [{"criterion_id": "c1", "verdict": "passed", "reason": "verified"}], "summary": "independent pass", "evidence_refs": [], "unverified_claims": [], "risks": [], "recommendation": "accept", } tool_name = "submit_validation" else: raise AssertionError(f"unexpected tools: {names}") return { "content": "", "tool_calls": [{ "id": f"call-{tool_name}", "type": "function", "function": {"name": tool_name, "arguments": json.dumps(arguments)}, }], "finish_reason": "tool_calls", } trace_store = FileSystemTraceStore(str(tmp_path)) await trace_store.create_trace(Trace(trace_id="root", mode="agent", task="mission", agent_role="planner")) runner = AgentRunner(trace_store=trace_store, llm_call=fake_llm) coordinator = wire_orchestration( runner, FileSystemTaskStore(str(tmp_path)), FileSystemArtifactStore(str(tmp_path)), ) await coordinator.ensure_ledger("root", ROOT_TASK_SPEC) task_id = await create_task(coordinator) cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0] ledger = await coordinator.task_store.load("root") worker = await trace_store.get_trace(ledger.attempts[cycle.attempt_id].worker_trace_id) validator = await trace_store.get_trace(cycle.validation.validator_trace_id) assert cycle.task_status == TaskStatus.AWAITING_DECISION assert worker.agent_role == "worker" and worker.result_summary assert validator.agent_role == "validator" and validator.result_summary assert worker.trace_id != validator.trace_id class OneShotGetFailureArtifactStore(FileSystemArtifactStore): def __init__(self, base_path): super().__init__(base_path) self.get_calls = 0 def for_root(self, root_trace_id): self.root_trace_id = root_trace_id return self async def get(self, root_trace_id, snapshot_id=None): self.get_calls += 1 if self.get_calls == 2: raise RuntimeError("injected artifact read failure") return await super().get(root_trace_id, snapshot_id) class FailingDispatchCompletionTaskStore(FileSystemTaskStore): def __init__(self, base_path): super().__init__(base_path) self.fail_next_dispatch_completion = True async def commit(self, ledger, expected_revision, idempotency_key=None, event=None): if ( self.fail_next_dispatch_completion and idempotency_key == "root:batch" ): self.fail_next_dispatch_completion = False raise RuntimeError("injected dispatch completion persistence failure") return await super().commit( ledger, expected_revision, idempotency_key=idempotency_key, event=event, ) class OneShotLoadFailureTaskStore(FileSystemTaskStore): def __init__(self, base_path): super().__init__(base_path) self.fail_next_load = False self._active_operation_loads_before_failure = 1 async def load(self, root_trace_id): ledger = await super().load(root_trace_id) if self.fail_next_load and any( operation.status.value == "running" for operation in ledger.operations.values() ): if self._active_operation_loads_before_failure: self._active_operation_loads_before_failure -= 1 else: self.fail_next_load = False raise RuntimeError("injected reservation load failure") return ledger class FailingEventSink: async def emit(self, root_trace_id, event_type, payload): raise RuntimeError("injected event sink failure") @pytest.mark.asyncio async def test_parallel_unexpected_branch_error_does_not_cancel_siblings(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED] * 4, delay=0.01) artifact_store = OneShotGetFailureArtifactStore(str(tmp_path)) coordinator, store, _ = await make_coordinator( tmp_path, executor, artifact_store=artifact_store ) task_ids = [await create_task(coordinator, f"isolated-{index}") for index in range(4)] results = await coordinator.dispatch_tasks("root", task_ids) assert [result.task_id for result in results] == task_ids failures = [result for result in results if result.error] assert len(failures) == 1 assert "artifact read failure" in failures[0].error assert failures[0].task_status == TaskStatus.NEEDS_REPLAN successes = [result for result in results if not result.error] assert len(successes) == 3 assert all(result.task_status == TaskStatus.AWAITING_DECISION for result in successes) ledger = await store.load("root") assert all( ledger.tasks[task_id].status not in {TaskStatus.RUNNING, TaskStatus.VALIDATING} for task_id in task_ids ) @pytest.mark.asyncio async def test_batch_reservation_conflict_does_not_strand_other_tasks(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_ids = [await create_task(coordinator, f"reserve-{index}") for index in range(3)] existing = await coordinator._create_attempt( "root", task_ids[1], "worker", "existing-reservation" ) results = await coordinator.dispatch_tasks("root", task_ids) assert [result.task_id for result in results] == task_ids assert results[0].task_status == TaskStatus.AWAITING_DECISION assert results[1].attempt_id is None assert "Dispatch conflict" in results[1].error assert results[2].task_status == TaskStatus.AWAITING_DECISION ledger = await store.load("root") assert ledger.tasks[task_ids[0]].status == TaskStatus.AWAITING_DECISION assert ledger.tasks[task_ids[2]].status == TaskStatus.AWAITING_DECISION assert ledger.attempts[existing["attempt_id"]].status.value == "running" @pytest.mark.asyncio async def test_explicit_lifecycle_never_creates_goal_tree(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) coordinator, store, trace_store = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator, "task-ledger-only") result = (await coordinator.dispatch_tasks("root", [task_id]))[0] assert result.task_status == TaskStatus.AWAITING_DECISION assert (await store.load("root")).tasks[task_id].status == TaskStatus.AWAITING_DECISION assert await trace_store.get_goal_tree("root") is None @pytest.mark.asyncio async def test_task_context_is_rendered_from_ledger(tmp_path): coordinator, _, _ = await make_coordinator(tmp_path, FakeExecutor([])) await create_task(coordinator, "draft the answer") context = await coordinator.task_context("root") assert "## Current Task Plan" in context assert "- 0 [waiting_children] mission" in context assert "- 0.1 [pending] draft the answer" in context @pytest.mark.asyncio async def test_dispatch_returns_results_when_batch_result_persistence_fails(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED]) task_store = FailingDispatchCompletionTaskStore(str(tmp_path)) coordinator, store, _ = await make_coordinator( tmp_path, executor, task_store=task_store, ) task_ids = [await create_task(coordinator, f"batch-{index}") for index in range(2)] first = await coordinator.dispatch_tasks( "root", task_ids, idempotency_key="batch", ) assert [result.task_status for result in first] == [ TaskStatus.AWAITING_DECISION, TaskStatus.AWAITING_DECISION, ] assert executor.worker_calls == 2 assert executor.validator_calls == 2 repeated = await coordinator.dispatch_tasks( "root", task_ids, idempotency_key="batch", ) assert [result.attempt_id for result in repeated] == [ result.attempt_id for result in first ] assert executor.worker_calls == 2 assert executor.validator_calls == 2 ledger = await store.load("root") assert all(len(ledger.tasks[task_id].attempt_ids) == 1 for task_id in task_ids) assert all(len(ledger.tasks[task_id].validation_ids) == 1 for task_id in task_ids) @pytest.mark.asyncio async def test_dispatch_idempotency_key_is_bound_to_tasks_and_presets(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) coordinator, _, _ = await make_coordinator(tmp_path, executor) first_task = await create_task(coordinator, "first-binding") second_task = await create_task(coordinator, "second-binding") await coordinator.dispatch_tasks( "root", [first_task], worker_presets=["worker"], idempotency_key="bound-batch", ) with pytest.raises(TaskConflict, match="different task_ids"): await coordinator.dispatch_tasks( "root", [second_task], worker_presets=["worker"], idempotency_key="bound-batch", ) with pytest.raises(TaskConflict, match="different worker_presets"): await coordinator.dispatch_tasks( "root", [first_task], worker_presets=["alternate-worker"], idempotency_key="bound-batch", ) assert executor.worker_calls == 1 assert executor.validator_calls == 1 @pytest.mark.asyncio async def test_concurrent_dispatch_replay_does_not_duplicate_agent_execution(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED], delay=0.02) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator, "concurrent-idempotency") first, replay = await asyncio.gather( coordinator.dispatch_tasks("root", [task_id], idempotency_key="concurrent"), coordinator.dispatch_tasks("root", [task_id], idempotency_key="concurrent"), ) assert executor.worker_calls == 1 assert executor.validator_calls == 1 assert any( result[0].task_status == TaskStatus.AWAITING_DECISION for result in (first, replay) ) assert first[0].attempt_id == replay[0].attempt_id assert all(result[0].task_status == TaskStatus.AWAITING_DECISION for result in (first, replay)) final = await coordinator.dispatch_tasks( "root", [task_id], idempotency_key="concurrent", ) ledger = await store.load("root") assert final[0].task_status == TaskStatus.AWAITING_DECISION assert len(ledger.tasks[task_id].attempt_ids) == 1 assert len(ledger.tasks[task_id].validation_ids) == 1 @pytest.mark.asyncio async def test_reservation_load_failure_does_not_corrupt_existing_attempt(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) task_store = OneShotLoadFailureTaskStore(str(tmp_path)) coordinator, store, _ = await make_coordinator( tmp_path, executor, task_store=task_store, ) running_task = await create_task(coordinator, "already-running") existing = await coordinator._create_attempt( "root", running_task, "worker", "existing-running-attempt" ) sibling_task = await create_task(coordinator, "unrelated-sibling") task_store.fail_next_load = True results = await coordinator.dispatch_tasks( "root", [running_task, sibling_task], ) assert [result.task_id for result in results] == [running_task, sibling_task] assert results[0].attempt_id is None assert results[0].task_status == TaskStatus.RUNNING assert "reservation load failure" in results[0].error assert results[1].task_status == TaskStatus.AWAITING_DECISION ledger = await store.load("root") assert ledger.tasks[running_task].status == TaskStatus.RUNNING assert ledger.attempts[existing["attempt_id"]].status.value == "running" @pytest.mark.asyncio async def test_event_sink_failure_does_not_rollback_committed_ledger(tmp_path): executor = FakeExecutor([]) coordinator, store, _ = await make_coordinator(tmp_path, executor) coordinator.event_sink = FailingEventSink() task_id = await create_task(coordinator, "event-outage") ledger = await store.load("root") assert ledger.tasks[task_id].status == TaskStatus.PENDING @pytest.mark.asyncio async def test_background_operation_is_durable_and_idempotently_bound(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator, "durable-operation") operation = await coordinator.start_operation( "root", OperationKind.DISPATCH, task_ids=[task_id], idempotency_key="operation-key", ) completed = await coordinator.await_operation("root", operation.operation_id) replay = await coordinator.start_operation( "root", OperationKind.DISPATCH, task_ids=[task_id], idempotency_key="operation-key", ) ledger = await store.load("root") assert completed.status == OperationStatus.COMPLETED assert replay.operation_id == completed.operation_id assert ledger.operations[completed.operation_id].attempt_ids assert ledger.operations[completed.operation_id].validation_ids assert executor.worker_calls == executor.validator_calls == 1 other = await create_task(coordinator, "different-operation") with pytest.raises(TaskConflict, match="different task_ids"): await coordinator.start_operation( "root", OperationKind.DISPATCH, task_ids=[other], idempotency_key="operation-key", ) @pytest.mark.asyncio async def test_submitted_attempt_advances_to_validation_without_worker_replay(tmp_path): executor = FakeExecutor([ValidationVerdict.PASSED]) coordinator, store, _ = await make_coordinator(tmp_path, executor) task_id = await create_task(coordinator, "resume-after-submit") reserved = await coordinator._create_attempt( "root", task_id, "worker", "resume-reservation" ) await coordinator.submit_attempt( { "role": AgentRole.WORKER.value, "root_trace_id": "root", "task_id": task_id, "attempt_id": reserved["attempt_id"], "trace_id": reserved["worker_trace_id"], "spec_version": 1, "tool_call_id": "resume-submission", }, AttemptSubmission( summary="already executed", artifact_refs=[ArtifactRef(uri="memory://submitted", version="1")], ), ) result = await coordinator.advance_cycle("root", task_id, reserved["attempt_id"]) assert result.task_status == TaskStatus.AWAITING_DECISION assert executor.worker_calls == 0 assert executor.validator_calls == 1 assert len((await store.load("root")).tasks[task_id].attempt_ids) == 1