| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192 |
- 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", "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 == "cross":
- attempt.accepted_child_decision_ids = ("missing-decision",)
- 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
- 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
|