| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168 |
- from __future__ import annotations
- from agent.orchestration import (
- AgentRole,
- ArtifactRef,
- CriterionResult,
- DeterministicWorkerContext,
- DeterministicWorkerResult,
- OrchestrationConfig,
- RoleContextRequest,
- TaskCoordinator,
- ValidationVerdict,
- )
- from agent.orchestration.protocols import ValidatorRunResult
- from agent.orchestration.store import FileSystemArtifactStore, FileSystemTaskStore
- from agent.trace.models import Trace
- from agent.trace.store import FileSystemTraceStore
- import pytest
- class _TaskContext:
- async def render(self, root_trace_id: str, ledger: object) -> str:
- assert root_trace_id == "root"
- return '{"state_revision":"ledger:compact"}'
- class _RoleContext:
- def __init__(self) -> None:
- self.requests: list[RoleContextRequest] = []
- async def build(self, request: RoleContextRequest) -> dict[str, object]:
- self.requests.append(request)
- return {
- "state_revision": f"ledger:{request.ledger_revision}",
- "role": request.role.value,
- }
- class _DeterministicWorker:
- def __init__(self) -> None:
- self.calls: list[DeterministicWorkerContext] = []
- async def supports(self, context: DeterministicWorkerContext) -> bool:
- return True
- async def execute(
- self, context: DeterministicWorkerContext
- ) -> DeterministicWorkerResult:
- self.calls.append(context)
- return DeterministicWorkerResult(
- summary="Host assembled the immutable artifact",
- artifact_refs=(ArtifactRef(uri="memory://deterministic", version="1"),),
- )
- class _Executor:
- def __init__(self) -> None:
- self.coordinator: TaskCoordinator | None = None
- self.worker_calls = 0
- self.validator_calls = 0
- async def run_worker(self, context: dict[str, object]) -> object:
- self.worker_calls += 1
- raise AssertionError(
- "LLM Worker must not run when deterministic execution supports Task"
- )
- async def run_validator(self, context: dict[str, object]) -> ValidatorRunResult:
- self.validator_calls += 1
- assert context["role_context"]
- coordinator = self.coordinator
- assert coordinator is not None
- task_spec = context["task_spec"]
- assert isinstance(task_spec, dict)
- criteria = [
- CriterionResult(
- criterion_id=str(item["criterion_id"]),
- verdict=ValidationVerdict.PASSED,
- reason="independently checked",
- )
- for item in task_spec["acceptance_criteria"]
- ]
- await coordinator.submit_validation(
- {
- **context,
- "role": AgentRole.VALIDATOR.value,
- "trace_id": context["validator_trace_id"],
- "tool_call_id": "validator-submit",
- },
- ValidationVerdict.PASSED,
- criteria,
- "passed",
- [],
- [],
- [],
- "accept",
- )
- return ValidatorRunResult(str(context["validator_trace_id"]), "completed")
- @pytest.mark.asyncio
- async def test_context_ports_and_deterministic_worker_preserve_full_audit_chain(
- tmp_path,
- ) -> None:
- trace_store = FileSystemTraceStore(str(tmp_path / "traces"))
- await trace_store.create_trace(
- Trace(trace_id="root", mode="agent", task="mission", agent_role="planner")
- )
- role_context = _RoleContext()
- deterministic = _DeterministicWorker()
- executor = _Executor()
- task_store = FileSystemTaskStore(str(tmp_path / "ledger"))
- coordinator = TaskCoordinator(
- task_store,
- FileSystemArtifactStore(str(tmp_path / "artifacts")),
- trace_store,
- OrchestrationConfig(),
- executor=executor,
- task_context_provider=_TaskContext(),
- role_context_provider=role_context,
- deterministic_worker=deterministic,
- )
- executor.coordinator = coordinator
- await coordinator.ensure_ledger(
- "root",
- {
- "objective": "mission",
- "acceptance_criteria": [
- {"criterion_id": "root", "description": "complete", "hard": True}
- ],
- },
- )
- created = await coordinator.create_tasks(
- "root",
- [
- {
- "objective": "mechanical assembly",
- "acceptance_criteria": [
- {"criterion_id": "closed", "description": "closed", "hard": True}
- ],
- }
- ],
- )
- task_id = created["tasks"][0]["task_id"]
- cycle = (await coordinator.dispatch_tasks("root", [task_id]))[0]
- assert cycle.validation is not None
- assert cycle.validation.verdict is ValidationVerdict.PASSED
- assert executor.worker_calls == 0
- assert executor.validator_calls == 1
- assert len(deterministic.calls) == 1
- assert [item.role for item in role_context.requests] == [
- AgentRole.WORKER,
- AgentRole.VALIDATOR,
- ]
- assert (
- await coordinator.task_context("root") == '{"state_revision":"ledger:compact"}'
- )
- ledger = await task_store.load("root")
- attempt = ledger.attempts[cycle.attempt_id]
- assert attempt.submission is not None
- worker_trace = await trace_store.get_trace(attempt.worker_trace_id)
- assert worker_trace is not None
- assert worker_trace.status == "completed"
- assert worker_trace.total_tokens == 0
- assert worker_trace.context["deterministic_worker"] is True
- assert await trace_store.get_trace_messages(attempt.worker_trace_id) == []
|