| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846 |
- 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.goal_models import GoalTree
- 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
- 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"))
- await trace_store.update_goal_tree("root", GoalTree(mission="mission"))
- 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", "mission")
- 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", "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 == "1"
- assert ledger.tasks[second].display_path == "2"
- assert ledger.tasks[third].display_path == "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", "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] == ["1.1", "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"))
- await trace_store.update_goal_tree("root", GoalTree(mission="mission"))
- 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", "mission")
- 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 OneShotGoalProjectionFailureStore(FileSystemTraceStore):
- def __init__(self, base_path):
- super().__init__(base_path)
- self.fail_next_goal_update = False
- async def update_goal(self, trace_id, goal_id, cascade_completion=True, **updates):
- if self.fail_next_goal_update:
- self.fail_next_goal_update = False
- raise RuntimeError("injected goal projection failure")
- return await super().update_goal(
- trace_id,
- goal_id,
- cascade_completion=cascade_completion,
- **updates,
- )
- class OneShotGoalTreeProjectionFailureStore(FileSystemTraceStore):
- def __init__(self, base_path):
- super().__init__(base_path)
- self.fail_next_goal_tree_update = False
- async def update_goal_tree(self, trace_id, goal_tree):
- if self.fail_next_goal_tree_update:
- self.fail_next_goal_tree_update = False
- raise RuntimeError("injected goal tree projection failure")
- return await super().update_goal_tree(trace_id, goal_tree)
- 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_goal_projection_failure_does_not_abort_committed_attempt(tmp_path):
- executor = FakeExecutor([ValidationVerdict.PASSED])
- trace_store = OneShotGoalProjectionFailureStore(str(tmp_path))
- coordinator, store, _ = await make_coordinator(
- tmp_path,
- executor,
- trace_store=trace_store,
- )
- task_id = await create_task(coordinator, "reservation-projection-failure")
- trace_store.fail_next_goal_update = True
- result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
- assert result.task_id == task_id
- assert result.task_status == TaskStatus.AWAITING_DECISION
- assert result.attempt_id is not None
- assert result.error is None
- ledger = await store.load("root")
- assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION
- assert ledger.attempts[result.attempt_id].status.value == "submitted"
- @pytest.mark.asyncio
- async def test_task_creation_survives_goal_tree_projection_failure_and_reconciles(tmp_path):
- executor = FakeExecutor([])
- trace_store = OneShotGoalTreeProjectionFailureStore(str(tmp_path))
- coordinator, store, _ = await make_coordinator(
- tmp_path,
- executor,
- trace_store=trace_store,
- )
- trace_store.fail_next_goal_tree_update = True
- task_id = await create_task(coordinator, "projection-outage")
- ledger = await store.load("root")
- assert ledger.tasks[task_id].status == TaskStatus.PENDING
- assert ledger.tasks[task_id].goal_id is None
- result = await coordinator.reconcile_goal_tree("root")
- ledger = await store.load("root")
- assert result["reconciled_tasks"] == 1
- assert ledger.tasks[task_id].goal_id is not None
- @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
|