| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705 |
- import asyncio
- import pytest
- from agent.orchestration.config import OrchestrationConfig
- from agent.orchestration.coordinator import TaskConflict, TaskCoordinator
- from agent.orchestration.models import (
- AgentRole,
- ArtifactRef,
- AttemptStatus,
- AttemptSubmission,
- BackgroundOperation,
- CriterionResult,
- OperationKind,
- OperationStatus,
- TaskStatus,
- ValidationVerdict,
- )
- from agent.orchestration.protocols import WorkerRunResult
- from agent.orchestration.store import FileSystemArtifactStore, TraceEventSink
- from test_coordinator_integration import FakeExecutor, create_task, make_coordinator
- class BlockingWorkerExecutor:
- def __init__(self):
- self.coordinator = None
- self.started = asyncio.Event()
- self.context = None
- self.stop_calls = []
- async def run_worker(self, context):
- self.context = context
- self.started.set()
- await asyncio.Event().wait()
- async def run_validator(self, context):
- raise AssertionError("validator must not start")
- async def stop(self, trace_id):
- self.stop_calls.append(trace_id)
- return True
- class SlowWorkerExecutor:
- def __init__(self):
- self.coordinator = None
- async def run_worker(self, context):
- await asyncio.sleep(1)
- return WorkerRunResult(context["worker_trace_id"], "completed")
- async def run_validator(self, context):
- raise AssertionError("validator must not start")
- async def stop(self, trace_id):
- return True
- class HangingStopExecutor(BlockingWorkerExecutor):
- async def stop(self, trace_id):
- self.stop_calls.append(trace_id)
- await asyncio.Event().wait()
- class CancellationResistantExecutor(BlockingWorkerExecutor):
- def __init__(self):
- super().__init__()
- self.release = asyncio.Event()
- async def run_worker(self, context):
- self.context = context
- self.started.set()
- try:
- await asyncio.Event().wait()
- except asyncio.CancelledError:
- await self.release.wait()
- return WorkerRunResult(context["worker_trace_id"], "completed")
- class BlockingValidatorExecutor(FakeExecutor):
- def __init__(self):
- super().__init__([])
- self.validator_started = asyncio.Event()
- self.validator_context = None
- async def run_validator(self, context):
- self.validator_context = context
- self.validator_started.set()
- await asyncio.Event().wait()
- class SubmittedBlockingExecutor(FakeExecutor):
- def __init__(self):
- super().__init__([])
- self.submitted = asyncio.Event()
- self.release = asyncio.Event()
- async def run_worker(self, context):
- self.worker_calls += 1
- 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="durable before worker return",
- artifact_refs=[ArtifactRef(uri="memory://submitted", version="1")],
- ),
- )
- self.submitted.set()
- await self.release.wait()
- return WorkerRunResult(context["worker_trace_id"], "completed")
- def second_coordinator(tmp_path, store, trace_store, executor):
- coordinator = TaskCoordinator(
- store,
- FileSystemArtifactStore(str(tmp_path)),
- trace_store,
- OrchestrationConfig(stop_grace_seconds=0),
- TraceEventSink(str(tmp_path)),
- executor,
- )
- executor.coordinator = coordinator
- return coordinator
- async def persist_operation(store, operation):
- ledger = await store.load(operation.root_trace_id)
- ledger.operations[operation.operation_id] = operation
- await store.commit(ledger, ledger.revision)
- @pytest.mark.asyncio
- async def test_advance_normalizes_stop_requested_stages_atomically(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "normalize requested stop")
- operation = BackgroundOperation(
- operation_id="stop-requested-operation",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- )
- await persist_operation(store, operation)
- reserved = await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- await coordinator._mark_stage_started(
- "root",
- operation.operation_id,
- 1,
- attempt_id=reserved["attempt_id"],
- )
- ledger = await store.load("root")
- ledger.operations[operation.operation_id].status = OperationStatus.STOP_REQUESTED
- ledger.operations[operation.operation_id].execution_epoch = 2
- await store.commit(ledger, ledger.revision)
- stopped = await coordinator.advance_operation("root", operation.operation_id)
- ledger = await store.load("root")
- assert stopped.status == OperationStatus.STOPPED
- assert ledger.attempts[reserved["attempt_id"]].status == AttemptStatus.STOPPED
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- @pytest.mark.asyncio
- async def test_stop_marks_started_attempt_and_rejects_late_epoch_submission(tmp_path):
- executor = BlockingWorkerExecutor()
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- coordinator.config = OrchestrationConfig(stop_grace_seconds=0)
- task_id = await create_task(coordinator, "stop running worker")
- operation = await coordinator.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- await asyncio.wait_for(executor.started.wait(), timeout=1)
- stopped = await coordinator.stop_operation(
- "root", operation.operation_id, idempotency_key="stop-once"
- )
- ledger = await store.load("root")
- attempt = ledger.attempts[ledger.tasks[task_id].attempt_ids[-1]]
- assert stopped.status == OperationStatus.STOPPED
- assert attempt.status == AttemptStatus.STOPPED
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- assert executor.stop_calls == [attempt.worker_trace_id]
- with pytest.raises(TaskConflict, match="epoch|stopping"):
- await coordinator.submit_attempt(
- {
- **executor.context,
- "role": AgentRole.WORKER.value,
- "trace_id": attempt.worker_trace_id,
- "tool_call_id": "late-submit",
- },
- AttemptSubmission(
- summary="late",
- artifact_refs=[ArtifactRef(uri="memory://late", version="1")],
- ),
- )
- revision = ledger.revision
- replay = await coordinator.stop_operation(
- "root", operation.operation_id, idempotency_key="stop-once"
- )
- assert replay.status == OperationStatus.STOPPED
- assert (await store.load("root")).revision == revision
- with pytest.raises(TaskConflict, match="resume_not_safe"):
- await coordinator.resume_operation("root", operation.operation_id)
- @pytest.mark.asyncio
- async def test_hanging_executor_stop_is_bounded_by_grace_period(tmp_path):
- executor = HangingStopExecutor()
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- coordinator.config = OrchestrationConfig(stop_grace_seconds=0.01)
- task_id = await create_task(coordinator, "bounded stop")
- operation = await coordinator.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- await asyncio.wait_for(executor.started.wait(), timeout=1)
- stopped = await asyncio.wait_for(
- coordinator.stop_operation("root", operation.operation_id), timeout=0.5
- )
- assert stopped.status == OperationStatus.STOPPED
- assert executor.stop_calls
- assert (await store.load("root")).operations[operation.operation_id].status == OperationStatus.STOPPED
- @pytest.mark.asyncio
- async def test_cancellation_resistant_runtime_cannot_block_stop_or_commit_late(tmp_path):
- executor = CancellationResistantExecutor()
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- coordinator.config = OrchestrationConfig(stop_grace_seconds=0.01)
- task_id = await create_task(coordinator, "resists cancellation")
- operation = await coordinator.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- await asyncio.wait_for(executor.started.wait(), timeout=1)
- stopped = await asyncio.wait_for(
- coordinator.stop_operation("root", operation.operation_id), timeout=0.5
- )
- assert stopped.status == OperationStatus.STOPPED
- executor.release.set()
- await asyncio.wait_for(
- coordinator.runtime.wait("root", operation.operation_id), timeout=0.5
- )
- ledger = await store.load("root")
- assert ledger.operations[operation.operation_id].status == OperationStatus.STOPPED
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- @pytest.mark.asyncio
- async def test_worker_timeout_expires_attempt_without_starting_validation(tmp_path):
- executor = SlowWorkerExecutor()
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- coordinator.config = OrchestrationConfig(worker_timeout_seconds=0.001)
- task_id = await create_task(coordinator, "timeout")
- result = (await coordinator.dispatch_tasks("root", [task_id]))[0]
- ledger = await store.load("root")
- assert result.task_status == TaskStatus.NEEDS_REPLAN
- assert ledger.attempts[result.attempt_id].status == AttemptStatus.EXPIRED
- assert ledger.tasks[task_id].validation_ids == []
- @pytest.mark.asyncio
- async def test_recover_then_resume_operation_before_any_attempt_is_safe(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "resume before reservation")
- operation = BackgroundOperation(
- operation_id="crashed-before-reservation",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- started_at="2026-07-18T00:00:00+00:00",
- )
- await persist_operation(store, operation)
- recovered = await coordinator.recover_operations("root")
- assert recovered[0].status == OperationStatus.STOPPED
- resumed = await coordinator.resume_operation(
- "root", operation.operation_id, idempotency_key="resume-safe"
- )
- completed = await coordinator.await_operation("root", resumed.operation_id)
- ledger = await store.load("root")
- assert completed.status == OperationStatus.COMPLETED
- assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION
- assert executor.worker_calls == executor.validator_calls == 1
- @pytest.mark.asyncio
- async def test_recover_submitted_attempt_resumes_at_validation_only(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "resume submitted")
- operation = BackgroundOperation(
- operation_id="crashed-after-submit",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- started_at="2026-07-18T00:00:00+00:00",
- )
- await persist_operation(store, operation)
- reserved = await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- 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,
- "operation_id": operation.operation_id,
- "execution_epoch": 1,
- "tool_call_id": "submitted-before-crash",
- },
- AttemptSubmission(
- summary="submitted",
- artifact_refs=[ArtifactRef(uri="memory://submitted", version="1")],
- ),
- )
- await coordinator.recover_operations("root")
- resumed = await coordinator.resume_operation("root", operation.operation_id)
- await coordinator.await_operation("root", resumed.operation_id)
- ledger = await store.load("root")
- assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION
- assert executor.worker_calls == 0
- assert executor.validator_calls == 1
- assert len(ledger.tasks[task_id].attempt_ids) == 1
- @pytest.mark.asyncio
- async def test_recover_started_attempt_is_not_resumable(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "unsafe crash")
- operation = BackgroundOperation(
- operation_id="crashed-during-worker",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- started_at="2026-07-18T00:00:00+00:00",
- )
- await persist_operation(store, operation)
- reserved = await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- ledger = await store.load("root")
- ledger.attempts[reserved["attempt_id"]].started_at = "2026-07-18T00:00:01+00:00"
- await store.commit(ledger, ledger.revision)
- await coordinator.recover_operations("root")
- ledger = await store.load("root")
- assert ledger.attempts[reserved["attempt_id"]].status == AttemptStatus.STOPPED
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- with pytest.raises(TaskConflict, match="resume_not_safe"):
- await coordinator.resume_operation("root", operation.operation_id)
- @pytest.mark.asyncio
- async def test_remote_stop_cannot_be_overwritten_and_submitted_attempt_resumes_validation(
- tmp_path,
- ):
- running_executor = SubmittedBlockingExecutor()
- coordinator_a, store, trace_store = await make_coordinator(
- tmp_path, running_executor
- )
- task_id = await create_task(coordinator_a, "remote stop after submit")
- operation = await coordinator_a.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- await asyncio.wait_for(running_executor.submitted.wait(), timeout=1)
- resume_executor = FakeExecutor([])
- coordinator_b = second_coordinator(
- tmp_path, store, trace_store, resume_executor
- )
- stopped = await coordinator_b.stop_operation("root", operation.operation_id)
- running_executor.release.set()
- await coordinator_a.runtime.wait("root", operation.operation_id)
- ledger = await store.load("root")
- assert stopped.status == OperationStatus.STOPPED
- assert ledger.operations[operation.operation_id].status == OperationStatus.STOPPED
- assert ledger.tasks[task_id].status == TaskStatus.AWAITING_VALIDATION
- assert ledger.tasks[task_id].validation_ids == []
- resumed = await coordinator_b.resume_operation("root", operation.operation_id)
- completed = await coordinator_b.await_operation("root", resumed.operation_id)
- ledger = await store.load("root")
- assert completed.status == OperationStatus.COMPLETED
- assert ledger.tasks[task_id].status == TaskStatus.AWAITING_DECISION
- assert resume_executor.worker_calls == 0
- assert resume_executor.validator_calls == 1
- @pytest.mark.asyncio
- async def test_past_deadline_never_starts_agents_and_is_part_of_idempotency(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "already expired")
- operation = await coordinator.start_operation(
- "root",
- OperationKind.DISPATCH,
- task_ids=[task_id],
- deadline_at="2000-01-01T00:00:00Z",
- idempotency_key="deadline-key",
- )
- completed = await coordinator.await_operation("root", operation.operation_id)
- assert completed.status == OperationStatus.FAILED
- assert "deadline" in (completed.error or "").lower()
- assert executor.worker_calls == executor.validator_calls == 0
- with pytest.raises(TaskConflict, match="different operation or request"):
- await coordinator.start_operation(
- "root",
- OperationKind.DISPATCH,
- task_ids=[task_id],
- deadline_at="2040-01-01T00:00:00+00:00",
- idempotency_key="deadline-key",
- )
- ledger = await store.load("root")
- command = ledger.command_records["root:operation.dispatch:deadline-key"]
- assert command.operation_id == operation.operation_id
- events = await store.list_events("root", limit=1000)
- started = [item for item in events.events if item.event_type == "operation_started"]
- assert started[-1].operation_id == operation.operation_id
- @pytest.mark.asyncio
- async def test_await_operation_polls_shared_store_when_other_runtime_owns_task(tmp_path):
- executor = FakeExecutor([], delay=0.05)
- coordinator_a, store, trace_store = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator_a, "cross coordinator await")
- operation = await coordinator_a.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- while (await coordinator_a.get_operation("root", operation.operation_id)).status != OperationStatus.RUNNING:
- await asyncio.sleep(0)
- observer = second_coordinator(tmp_path, store, trace_store, FakeExecutor([]))
- completed = await asyncio.wait_for(
- observer.await_operation("root", operation.operation_id), timeout=2
- )
- assert completed.status == OperationStatus.COMPLETED
- @pytest.mark.asyncio
- async def test_unstarted_revalidation_is_reused_after_remote_stop_and_resume(tmp_path):
- initial_executor = FakeExecutor([], submit_validator=False)
- coordinator_a, store, trace_store = await make_coordinator(
- tmp_path, initial_executor
- )
- task_id = await create_task(coordinator_a, "resume revalidation")
- first = (await coordinator_a.dispatch_tasks("root", [task_id]))[0]
- assert first.validation_id
- assert (await store.load("root")).tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- entered = asyncio.Event()
- release = asyncio.Event()
- original_advance = coordinator_a.advance_validation
- async def pause_before_validation(*args, **kwargs):
- entered.set()
- await release.wait()
- return await original_advance(*args, **kwargs)
- coordinator_a.advance_validation = pause_before_validation
- operation = await coordinator_a.start_operation(
- "root",
- OperationKind.REVALIDATE,
- task_id=task_id,
- attempt_id=first.attempt_id,
- )
- await asyncio.wait_for(entered.wait(), timeout=1)
- resume_executor = FakeExecutor([])
- coordinator_b = second_coordinator(
- tmp_path, store, trace_store, resume_executor
- )
- await coordinator_b.stop_operation("root", operation.operation_id)
- release.set()
- await coordinator_a.runtime.wait("root", operation.operation_id)
- ledger = await store.load("root")
- validation_id = ledger.operations[operation.operation_id].validation_ids[-1]
- assert ledger.validations[validation_id].started_at is None
- assert ledger.operations[operation.operation_id].status == OperationStatus.STOPPED
- resumed = await coordinator_b.resume_operation("root", operation.operation_id)
- completed = await coordinator_b.await_operation("root", resumed.operation_id)
- ledger = await store.load("root")
- assert completed.status == OperationStatus.COMPLETED
- assert ledger.operations[operation.operation_id].validation_ids == [validation_id]
- assert ledger.validations[validation_id].status.value == "completed"
- assert resume_executor.validator_calls == 1
- @pytest.mark.asyncio
- async def test_recovery_expires_reserved_unstarted_attempt_after_deadline(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "expired reserved attempt")
- operation = BackgroundOperation(
- operation_id="expired-attempt-operation",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- deadline_at="2000-01-01T00:00:00+00:00",
- )
- await persist_operation(store, operation)
- reserved = await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- await coordinator.recover_operations("root")
- ledger = await store.load("root")
- assert ledger.operations[operation.operation_id].status == OperationStatus.FAILED
- assert ledger.attempts[reserved["attempt_id"]].status == AttemptStatus.EXPIRED
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- @pytest.mark.asyncio
- async def test_recovery_expires_reserved_unstarted_validation_after_deadline(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "expired reserved validation")
- operation = BackgroundOperation(
- operation_id="expired-validation-operation",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.RUNNING,
- task_ids=[task_id],
- execution_epoch=1,
- deadline_at="2000-01-01T00:00:00+00:00",
- )
- await persist_operation(store, operation)
- reserved = await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- 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,
- "operation_id": operation.operation_id,
- "execution_epoch": 1,
- "tool_call_id": "submit-before-expiry",
- },
- AttemptSubmission(
- summary="submitted",
- artifact_refs=[ArtifactRef(uri="memory://submitted", version="1")],
- ),
- )
- validation_id = await coordinator._start_validation(
- "root",
- task_id,
- reserved["attempt_id"],
- operation_id=operation.operation_id,
- execution_epoch=1,
- )
- await coordinator.recover_operations("root")
- ledger = await store.load("root")
- assert ledger.operations[operation.operation_id].status == OperationStatus.FAILED
- assert ledger.validations[validation_id].status.value == "expired"
- assert ledger.tasks[task_id].status == TaskStatus.NEEDS_REPLAN
- @pytest.mark.asyncio
- async def test_terminal_operation_cannot_reserve_a_late_attempt(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "reject late reservation")
- operation = BackgroundOperation(
- operation_id="stopped-operation",
- root_trace_id="root",
- kind=OperationKind.DISPATCH,
- request={"kind": "dispatch", "task_ids": [task_id], "worker_presets": ["worker"]},
- request_fingerprint="fingerprint",
- status=OperationStatus.STOPPED,
- task_ids=[task_id],
- execution_epoch=4,
- )
- await persist_operation(store, operation)
- with pytest.raises(TaskConflict, match="stale|stopping"):
- await coordinator._create_attempt(
- "root",
- task_id,
- "worker",
- None,
- operation_id=operation.operation_id,
- execution_epoch=4,
- )
- ledger = await store.load("root")
- assert ledger.tasks[task_id].attempt_ids == []
- assert ledger.tasks[task_id].status == TaskStatus.PENDING
- @pytest.mark.asyncio
- @pytest.mark.parametrize(
- "criterion_ids",
- [["c1", "c1"], ["c1", "unknown"]],
- )
- async def test_validation_report_rejects_duplicate_or_unknown_criteria(
- tmp_path,
- criterion_ids,
- ):
- executor = BlockingValidatorExecutor()
- coordinator, _, _ = await make_coordinator(tmp_path, executor)
- task_id = await create_task(coordinator, "strict report schema")
- operation = await coordinator.start_operation(
- "root", OperationKind.DISPATCH, task_ids=[task_id]
- )
- await asyncio.wait_for(executor.validator_started.wait(), timeout=1)
- context = executor.validator_context
- with pytest.raises(ValueError, match="duplicate|unknown"):
- await coordinator.submit_validation(
- {
- **context,
- "role": AgentRole.VALIDATOR.value,
- "trace_id": context["validator_trace_id"],
- "tool_call_id": f"bad-report-{criterion_ids[-1]}",
- },
- ValidationVerdict.PASSED,
- [
- CriterionResult(item, ValidationVerdict.PASSED, "checked")
- for item in criterion_ids
- ],
- "bad report",
- [],
- [],
- [],
- "reject",
- )
- await coordinator.stop_operation("root", operation.operation_id)
|