| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148 |
- import pytest
- from agent.orchestration.coordinator import OrchestrationError
- from agent.orchestration.models import (
- DecisionAction,
- FailureCode,
- TaskStatus,
- ValidationVerdict,
- )
- from test_coordinator_integration import FakeExecutor, create_task, make_coordinator
- async def _accept_child(coordinator, task_id, key):
- cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
- await coordinator.decide_task(
- "root",
- task_id,
- cycle.validation_id,
- DecisionAction.ACCEPT,
- {"reason": "accepted"},
- key,
- )
- return cycle
- @pytest.mark.asyncio
- async def test_new_attempt_distinguishes_empty_and_ordered_child_bindings(tmp_path):
- executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- first = await create_task(coordinator, "first")
- second = await create_task(coordinator, "second")
- await _accept_child(coordinator, second, "accept-second")
- await _accept_child(coordinator, first, "accept-first")
- ledger = await store.load("root")
- root = ledger.tasks[ledger.root_task_id]
- reserved = await coordinator._create_attempt("root", root.task_id, "worker", None)
- frozen = (await store.load("root")).attempts[reserved["attempt_id"]]
- assert frozen.accepted_child_decision_ids == (
- ledger.tasks[first].decision_ids[-1],
- ledger.tasks[second].decision_ids[-1],
- )
- first_attempt = ledger.attempts[ledger.tasks[first].attempt_ids[-1]]
- assert first_attempt.accepted_child_decision_ids == ()
- @pytest.mark.asyncio
- @pytest.mark.parametrize(
- "corruption", ["duplicate", "reverse", "missing", "cross", "non_accept"]
- )
- async def test_frozen_child_decision_binding_rejects_corruption(tmp_path, corruption):
- executor = FakeExecutor([ValidationVerdict.PASSED, ValidationVerdict.PASSED])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- first = await create_task(coordinator, "first")
- second = await create_task(coordinator, "second")
- await _accept_child(coordinator, first, "accept-first")
- await _accept_child(coordinator, second, "accept-second")
- ledger = await store.load("root")
- root = ledger.tasks[ledger.root_task_id]
- reserved = await coordinator._create_attempt("root", root.task_id, "worker", None)
- ledger = await store.load("root")
- attempt = ledger.attempts[reserved["attempt_id"]]
- ids = list(attempt.accepted_child_decision_ids)
- if corruption == "duplicate":
- attempt.accepted_child_decision_ids = (ids[0], ids[0])
- elif corruption == "reverse":
- attempt.accepted_child_decision_ids = tuple(reversed(ids))
- elif corruption == "missing":
- attempt.accepted_child_decision_ids = ("missing-decision",)
- elif corruption == "cross":
- decision = ledger.decisions[ids[0]]
- object.__setattr__(decision, "task_id", root.task_id)
- else:
- decision = ledger.decisions[ids[0]]
- object.__setattr__(decision, "action", DecisionAction.RETRY)
- with pytest.raises(OrchestrationError):
- coordinator._accepted_child_results(ledger, root, attempt)
- @pytest.mark.asyncio
- @pytest.mark.parametrize("missing", ["attempt", "validation", "snapshot"])
- async def test_child_accept_binding_requires_complete_references(tmp_path, missing):
- executor = FakeExecutor([ValidationVerdict.PASSED])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- child = await create_task(coordinator, "child")
- await _accept_child(coordinator, child, "accept-child")
- ledger = await store.load("root")
- root = ledger.tasks[ledger.root_task_id]
- reserved = await coordinator._create_attempt("root", root.task_id, "worker", None)
- ledger = await store.load("root")
- parent_attempt = ledger.attempts[reserved["attempt_id"]]
- decision = ledger.decisions[parent_attempt.accepted_child_decision_ids[0]]
- child_attempt = ledger.attempts[decision.attempt_id]
- if missing == "attempt":
- del ledger.attempts[decision.attempt_id]
- elif missing == "validation":
- del ledger.validations[decision.validation_id]
- else:
- child_attempt.snapshot_id = None
- with pytest.raises(OrchestrationError):
- coordinator._accepted_child_results(ledger, root, parent_attempt)
- @pytest.mark.asyncio
- async def test_missing_artifact_snapshot_fails_parent_as_protocol_violation(tmp_path):
- executor = FakeExecutor([ValidationVerdict.PASSED])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- child = await create_task(coordinator, "child")
- await _accept_child(coordinator, child, "accept-child")
- ledger = await store.load("root")
- root = ledger.tasks[ledger.root_task_id]
- reserved = await coordinator._create_attempt("root", root.task_id, "worker", None)
- decision = ledger.decisions[ledger.tasks[child].decision_ids[-1]]
- child_attempt = ledger.attempts[decision.attempt_id]
- artifact_path = (
- tmp_path
- / "root"
- / "orchestration"
- / "artifacts"
- / f"{child_attempt.snapshot_id}.json"
- )
- artifact_path.unlink()
- result = await coordinator.advance_cycle("root", root.task_id, reserved["attempt_id"])
- failed = (await store.load("root")).attempts[reserved["attempt_id"]]
- assert "Invalid frozen child input binding" in result.error
- assert failed.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
- assert executor.worker_calls == 1 # child only; the parent Worker was never invoked
- @pytest.mark.asyncio
- async def test_legacy_running_attempt_never_guesses_child_inputs(tmp_path):
- executor = FakeExecutor([])
- coordinator, store, _ = await make_coordinator(tmp_path, executor)
- root_id = (await store.load("root")).root_task_id
- reserved = await coordinator._create_attempt("root", root_id, "worker", None)
- ledger = await store.load("root")
- ledger.attempts[reserved["attempt_id"]].accepted_child_decision_ids = None
- await store.commit(ledger, expected_revision=ledger.revision)
- result = await coordinator.advance_cycle("root", root_id, reserved["attempt_id"])
- recovered = await store.load("root")
- attempt = recovered.attempts[reserved["attempt_id"]]
- assert result.error == "Legacy running attempt has no frozen child input binding"
- assert executor.worker_calls == 0
- assert recovered.tasks[root_id].status == TaskStatus.NEEDS_REPLAN
- assert attempt.execution_stats.failure_code == FailureCode.PROTOCOL_VIOLATION
|