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)