| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932 |
- from __future__ import annotations
- import asyncio
- from dataclasses import replace
- from datetime import UTC, datetime
- from types import SimpleNamespace
- from unittest.mock import AsyncMock
- import pytest
- from agent.orchestration import DecisionAction, OperationStatus, TaskStatus
- from agent.trace.models import Message, Trace
- from agent.trace.store import FileSystemTraceStore
- from script_build_host.application.mission_factory import ScriptMissionFactory
- from script_build_host.application.mission_service import (
- PHASE_ONE_CAPABILITY_BOUNDARY,
- PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
- ScriptMissionService,
- )
- from script_build_host.application.observation_views import mission_snapshot_view
- from script_build_host.domain.errors import (
- MissionRecoveryRequired,
- PhaseTwoBoundaryNotReady,
- PlannerPolicyMigrationRequired,
- )
- from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
- from script_build_host.domain.records import (
- BuildStatus,
- MissionBinding,
- Principal,
- PublicationState,
- )
- from script_build_host.infrastructure.canonical_json import canonical_sha256
- class _TraceStore:
- def __init__(self) -> None:
- self.events: list[dict[str, object]] = []
- self.messages: list[SimpleNamespace] = []
- async def get_events(self, _root: str, _since: int):
- return list(self.events)
- async def append_event(self, _root: str, event: str, payload: dict[str, object]):
- self.events.append({"event": event, **payload})
- return len(self.events)
- async def get_trace_messages(self, _root: str):
- return list(self.messages)
- class _Runner:
- def __init__(self, root: SimpleNamespace) -> None:
- self.root = root
- self.trace_store = _TraceStore()
- self.tools = SimpleNamespace(
- get_tool_names=lambda **_kwargs: ["plan_script_tasks", "decide_script_task"]
- )
- self.observed_config = None
- async def run(self, *, messages, config):
- self.observed_config = config
- assert [item["role"] for item in messages] == ["system", "user"]
- yield SimpleNamespace(trace_id="root")
- yield SimpleNamespace(role="system", sequence=11)
- yield SimpleNamespace(role="user", sequence=12)
- assert self.root.status is TaskStatus.NEEDS_REPLAN
- self.root.status = TaskStatus.BLOCKED
- self.root.blocked_reason = PHASE_TWO_CANDIDATE_PORTFOLIO_READY
- class _Coordinator:
- def __init__(self, root: SimpleNamespace) -> None:
- self.root = root
- self.unblock_calls = 0
- self.task_store = SimpleNamespace(
- load=lambda _root: _async(
- SimpleNamespace(root_task_id="root-task", tasks={"root-task": root})
- )
- )
- async def decide_task(
- self,
- _root_trace_id,
- task_id,
- _validation_id,
- action,
- _payload,
- idempotency_key,
- ):
- assert task_id == "root-task"
- assert action is DecisionAction.UNBLOCK
- assert idempotency_key == "phase-two-unblock:7"
- self.unblock_calls += 1
- self.root.status = TaskStatus.NEEDS_REPLAN
- self.root.blocked_reason = None
- async def root_completion(self, _root: str):
- return {
- "status": self.root.status.value,
- "blocked_reason": self.root.blocked_reason,
- }
- class _State:
- def __init__(self) -> None:
- self.status = BuildStatus.RUNNING
- self.checkpoint = None
- async def get_status(self, _build: int):
- return self.status
- async def set_status(self, _build: int, status: BuildStatus, **_kwargs):
- self.status = status
- async def set_checkpoint(self, _build: int, **values):
- self.status = BuildStatus.PARTIAL
- self.checkpoint = values
- def _snapshot() -> ScriptBuildInputSnapshotV1:
- return ScriptBuildInputSnapshotV1(
- snapshot_id="11",
- script_build_id=7,
- execution_id=1,
- topic_build_id=2,
- topic_id=3,
- topic={"topic": {"id": 3, "result": "topic"}},
- account={"account_name": "acct"},
- persona_points=(),
- section_patterns=(),
- strategies=(),
- prompt_manifest=(),
- datasource_manifest={},
- model_manifest={
- "presets": {
- "script_planner": {
- "model": "fake",
- "temperature": 0,
- "max_iterations": 10,
- }
- }
- },
- canonical_sha256="sha256:" + "1" * 64,
- created_at=datetime.now(UTC),
- )
- def _binding() -> MissionBinding:
- now = datetime.now(UTC)
- return MissionBinding(
- binding_id=1,
- script_build_id=7,
- root_trace_id="root",
- input_snapshot_id=11,
- active_direction_artifact_version_id=3,
- accepted_root_artifact_version_id=None,
- engine_version="test",
- schema_version="v1",
- created_at=now,
- updated_at=now,
- )
- @pytest.mark.asyncio
- async def test_phase_two_policy_is_persisted_before_root_unblock_and_same_trace_continues() -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.BLOCKED,
- blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
- )
- runner = _Runner(root)
- coordinator = _Coordinator(root)
- state = _State()
- service = ScriptMissionService(
- runner=runner,
- coordinator=coordinator,
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert runner.observed_config.trace_id == "root"
- assert runner.observed_config.new_trace_id is None
- assert runner.trace_store.events[0]["event"] == "planner_policy_migrated"
- assert state.status is BuildStatus.PARTIAL
- assert state.checkpoint["checkpoint_code"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
- async def _async(value):
- return value
- class _ReentryRunner(_Runner):
- def __init__(self, root: SimpleNamespace) -> None:
- super().__init__(root)
- self.run_result_calls = 0
- async def run_result(self, *, messages, config):
- self.run_result_calls += 1
- self.observed_config = config
- assert messages == []
- assert self.root.status is TaskStatus.NEEDS_REPLAN
- self.root.status = TaskStatus.BLOCKED
- self.root.blocked_reason = PHASE_TWO_CANDIDATE_PORTFOLIO_READY
- return SimpleNamespace()
- def _migration_event(runner: _Runner, snapshot: ScriptBuildInputSnapshotV1) -> dict[str, object]:
- factory = ScriptMissionFactory()
- policy = factory.build_phase_two_policy(snapshot)
- continuation = factory.build_phase_two_message(
- _binding(), snapshot, direction_artifact_version_id=3
- )
- runner.trace_store.messages = [
- SimpleNamespace(sequence=11, role="system", content=policy, parent_sequence=10),
- SimpleNamespace(sequence=12, role="user", content=continuation, parent_sequence=11),
- ]
- return {
- "event": "planner_policy_migrated",
- "policy_version": "script-build-phase-two/v1",
- "policy_digest": canonical_sha256(policy).wire,
- "toolset_digest": canonical_sha256(
- sorted(runner.tools.get_tool_names(groups=["script_build"]))
- ).wire,
- "input_snapshot_id": snapshot.snapshot_id,
- "system_sequence": 11,
- "user_sequence": 12,
- }
- @pytest.mark.asyncio
- @pytest.mark.parametrize("root_status", [TaskStatus.BLOCKED, TaskStatus.NEEDS_REPLAN])
- async def test_phase_two_policy_reentry_resumes_before_or_after_unblock_without_rewriting_policy(
- root_status: TaskStatus,
- ) -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=root_status,
- blocked_reason=(
- PHASE_ONE_CAPABILITY_BOUNDARY if root_status is TaskStatus.BLOCKED else None
- ),
- )
- runner = _ReentryRunner(root)
- runner.trace_store.events = [_migration_event(runner, _snapshot())]
- coordinator = _Coordinator(root)
- state = _State()
- state.status = BuildStatus.PARTIAL
- service = ScriptMissionService(
- runner=runner,
- coordinator=coordinator,
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert runner.run_result_calls == 1
- assert len(runner.trace_store.events) == 1
- assert coordinator.unblock_calls == (1 if root_status is TaskStatus.BLOCKED else 0)
- assert state.status is BuildStatus.PARTIAL
- assert state.checkpoint["checkpoint_code"] == PHASE_TWO_CANDIDATE_PORTFOLIO_READY
- @pytest.mark.asyncio
- async def test_phase_two_reentry_with_output_requires_recovery_without_terminal_failure() -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.NEEDS_REPLAN,
- blocked_reason=None,
- )
- runner = _ReentryRunner(root)
- runner.trace_store.events = [_migration_event(runner, _snapshot())]
- runner.trace_store.messages.append(
- SimpleNamespace(sequence=13, role="assistant", content="started")
- )
- state = _State()
- state.status = BuildStatus.RUNNING
- service = ScriptMissionService(
- runner=runner,
- coordinator=_Coordinator(root),
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- with pytest.raises(MissionRecoveryRequired):
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert state.status is BuildStatus.RUNNING
- assert runner.run_result_calls == 0
- @pytest.mark.asyncio
- @pytest.mark.parametrize(
- ("root_status", "has_model_output"),
- [
- (TaskStatus.BLOCKED, False),
- (TaskStatus.NEEDS_REPLAN, False),
- (TaskStatus.NEEDS_REPLAN, True),
- ],
- )
- async def test_historical_snapshot_v2_policy_reentry_survives_filestore_reload(
- tmp_path,
- root_status: TaskStatus,
- has_model_output: bool,
- ) -> None:
- parent = _snapshot()
- business_digest = canonical_sha256(parent.to_input().business_payload()).wire
- current = replace(
- parent,
- snapshot_id="12",
- parent_snapshot_id=parent.snapshot_id,
- parent_snapshot_sha256=parent.canonical_sha256,
- business_input_sha256=business_digest,
- )
- binding = replace(_binding(), input_snapshot_id=12)
- policy = ScriptMissionFactory().build_phase_two_policy(current)
- continuation = ScriptMissionFactory().build_phase_two_message(
- binding,
- current,
- direction_artifact_version_id=3,
- )
- store_path = tmp_path / "phase-two-reentry"
- writer = FileSystemTraceStore(str(store_path))
- await writer.create_trace(
- Trace(
- trace_id="root",
- mode="agent",
- agent_type="script_planner",
- agent_role="planner",
- )
- )
- await writer.add_message(
- Message.create(
- trace_id="root",
- role="system",
- sequence=11,
- parent_sequence=10,
- content=policy,
- )
- )
- await writer.add_message(
- Message.create(
- trace_id="root",
- role="user",
- sequence=12,
- parent_sequence=11,
- content=continuation,
- )
- )
- toolset_digest = canonical_sha256(["decide_script_task", "plan_script_tasks"]).wire
- await writer.append_event(
- "root",
- "planner_policy_migrated",
- {
- "policy_version": "script-build-phase-two/v1",
- "policy_digest": canonical_sha256(policy).wire,
- "toolset_digest": toolset_digest,
- "input_snapshot_id": current.snapshot_id,
- "system_sequence": 11,
- "user_sequence": 12,
- },
- )
- if has_model_output:
- await writer.add_message(
- Message.create(
- trace_id="root",
- role="assistant",
- sequence=13,
- parent_sequence=12,
- content={"text": "phase two already started"},
- )
- )
- root = SimpleNamespace(
- task_id="root-task",
- status=root_status,
- blocked_reason=(
- PHASE_ONE_CAPABILITY_BOUNDARY if root_status is TaskStatus.BLOCKED else None
- ),
- )
- runner = _ReentryRunner(root)
- runner.trace_store = FileSystemTraceStore(str(store_path))
- coordinator = _Coordinator(root)
- state = _State()
- state.status = BuildStatus.PARTIAL if root_status is TaskStatus.BLOCKED else BuildStatus.RUNNING
- service = ScriptMissionService(
- runner=runner,
- coordinator=coordinator,
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- if has_model_output:
- with pytest.raises(MissionRecoveryRequired):
- await service.run_phase_two(
- 7,
- binding=binding,
- snapshot=current,
- direction_artifact_version_id=3,
- )
- assert runner.run_result_calls == 0
- assert state.status is BuildStatus.RUNNING
- return
- await service.run_phase_two(
- 7,
- binding=binding,
- snapshot=current,
- direction_artifact_version_id=3,
- )
- reloaded = FileSystemTraceStore(str(store_path))
- messages = await reloaded.get_trace_messages("root")
- events = await reloaded.get_events("root", 0)
- assert [(item.sequence, item.role) for item in messages if item.sequence >= 11] == [
- (11, "system"),
- (12, "user"),
- ]
- assert len([item for item in events if item.get("event") == "planner_policy_migrated"]) == 1
- assert coordinator.unblock_calls == (1 if root_status is TaskStatus.BLOCKED else 0)
- assert state.status is BuildStatus.PARTIAL
- @pytest.mark.asyncio
- async def test_phase_two_reentry_rejects_unbacked_policy_event() -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.BLOCKED,
- blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
- )
- runner = _ReentryRunner(root)
- runner.trace_store.events = [_migration_event(runner, _snapshot())]
- runner.trace_store.messages[0].content = "tampered policy"
- state = _State()
- state.status = BuildStatus.PARTIAL
- service = ScriptMissionService(
- runner=runner,
- coordinator=_Coordinator(root),
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- with pytest.raises(PlannerPolicyMigrationRequired):
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert runner.run_result_calls == 0
- assert state.status is BuildStatus.FAILED
- @pytest.mark.asyncio
- async def test_phase_two_reentry_recovers_messages_written_before_migration_event() -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.BLOCKED,
- blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
- )
- runner = _ReentryRunner(root)
- event = _migration_event(runner, _snapshot())
- runner.trace_store.events = []
- runner.trace_store.messages[0].parent_sequence = 10
- runner.trace_store.messages[1].parent_sequence = 11
- state = _State()
- state.status = BuildStatus.PARTIAL
- service = ScriptMissionService(
- runner=runner,
- coordinator=_Coordinator(root),
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert runner.run_result_calls == 1
- assert len(runner.trace_store.events) == 1
- assert runner.trace_store.events[0]["policy_digest"] == event["policy_digest"]
- assert state.status is BuildStatus.PARTIAL
- @pytest.mark.asyncio
- @pytest.mark.parametrize("failure", ["duplicate", "toolset"])
- async def test_phase_two_reentry_rejects_ambiguous_policy_identity(failure: str) -> None:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.BLOCKED,
- blocked_reason=PHASE_ONE_CAPABILITY_BOUNDARY,
- )
- runner = _ReentryRunner(root)
- event = _migration_event(runner, _snapshot())
- if failure == "duplicate":
- runner.trace_store.events = [event, dict(event)]
- else:
- runner.trace_store.events = [{**event, "toolset_digest": "sha256:" + "0" * 64}]
- state = _State()
- state.status = BuildStatus.PARTIAL
- service = ScriptMissionService(
- runner=runner,
- coordinator=_Coordinator(root),
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(),
- legacy_state=state,
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- with pytest.raises(PlannerPolicyMigrationRequired):
- await service.run_phase_two(
- 7,
- binding=_binding(),
- snapshot=_snapshot(),
- direction_artifact_version_id=3,
- )
- assert runner.run_result_calls == 0
- assert state.status is BuildStatus.FAILED
- @pytest.mark.asyncio
- async def test_direction_accept_may_cross_only_one_verified_prompt_snapshot_lineage() -> None:
- parent = _snapshot()
- business_digest = canonical_sha256(parent.to_input().business_payload()).wire
- current = replace(
- parent,
- snapshot_id="12",
- parent_snapshot_id=parent.snapshot_id,
- parent_snapshot_sha256=parent.canonical_sha256,
- business_input_sha256=business_digest,
- )
- spec = SimpleNamespace(
- version=1,
- context_refs=(
- "script-build://task-kinds/direction",
- f"script-build://inputs/{parent.snapshot_id}",
- ),
- )
- direction = SimpleNamespace(
- status=TaskStatus.COMPLETED,
- current_spec=spec,
- specs=[spec],
- decision_ids=["direction-accept"],
- )
- ledger = SimpleNamespace(
- tasks={"direction": direction},
- decisions={
- "direction-accept": SimpleNamespace(
- action=DecisionAction.ACCEPT,
- attempt_id="direction-attempt",
- )
- },
- attempts={"direction-attempt": SimpleNamespace(spec_version=1)},
- )
- class Snapshots:
- async def get(self, snapshot_id: str, *, script_build_id: int):
- assert (snapshot_id, script_build_id) == ("11", 7)
- return parent
- service = ScriptMissionService(
- runner=SimpleNamespace(),
- coordinator=SimpleNamespace(),
- factory=ScriptMissionFactory(),
- input_snapshots=Snapshots(),
- bindings=SimpleNamespace(),
- legacy_state=SimpleNamespace(),
- authorizer=SimpleNamespace(),
- direction_reconciler=SimpleNamespace(),
- )
- await service._verify_direction_snapshot_lineage(
- script_build_id=7,
- ledger=ledger,
- current_snapshot=current,
- )
- tampered = replace(current, business_input_sha256="sha256:" + "9" * 64)
- with pytest.raises(PhaseTwoBoundaryNotReady, match="immutable lineage"):
- await service._verify_direction_snapshot_lineage(
- script_build_id=7,
- ledger=ledger,
- current_snapshot=tampered,
- )
- class _AdvanceSnapshots:
- def __init__(self) -> None:
- self.extend_calls = 0
- async def get(self, _snapshot_id: str, *, script_build_id: int):
- assert script_build_id == 7
- return _snapshot()
- async def extend_prompt_lineage(self, *_args, **_kwargs):
- self.extend_calls += 1
- return _snapshot()
- class _AdvanceState(_State):
- def __init__(self, status: BuildStatus) -> None:
- super().__init__()
- self.status = status
- def _advance_service(
- *,
- status: BuildStatus = BuildStatus.PARTIAL,
- blocked_reason: str = PHASE_ONE_CAPABILITY_BOUNDARY,
- operation_status: OperationStatus | None = None,
- publication_closed: bool = True,
- ) -> tuple[ScriptMissionService, _AdvanceSnapshots]:
- root = SimpleNamespace(
- task_id="root-task",
- status=TaskStatus.BLOCKED,
- blocked_reason=blocked_reason,
- decision_ids=["root-decision"],
- current_spec=SimpleNamespace(context_refs=()),
- )
- operation = (
- {"operation": SimpleNamespace(status=operation_status)}
- if operation_status is not None
- else {}
- )
- ledger = SimpleNamespace(
- root_task_id="root-task",
- tasks={"root-task": root},
- decisions={
- "root-decision": SimpleNamespace(
- action=DecisionAction.BLOCK,
- reason=blocked_reason,
- )
- },
- operations=operation,
- )
- async def root_completion(_root: str):
- return {"status": root.status.value, "blocked_reason": root.blocked_reason}
- coordinator = SimpleNamespace(
- root_completion=root_completion,
- task_store=SimpleNamespace(load=lambda _root: _async(ledger)),
- )
- snapshots = _AdvanceSnapshots()
- publication = (
- SimpleNamespace(
- state=PublicationState.PUBLISHED,
- artifact_version_id=3,
- )
- if publication_closed
- else None
- )
- reconciler = SimpleNamespace(
- reconcile=lambda _build, _root: _async(3),
- publications=SimpleNamespace(get_by_build=lambda *_args, **_kwargs: _async(publication)),
- )
- service = ScriptMissionService(
- runner=SimpleNamespace(
- trace_store=SimpleNamespace(
- get_trace=lambda _root: _async(SimpleNamespace(agent_type="script_planner"))
- )
- ),
- coordinator=coordinator,
- factory=ScriptMissionFactory(),
- input_snapshots=snapshots,
- bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
- legacy_state=_AdvanceState(status),
- authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
- direction_reconciler=reconciler,
- phase_two_prompt_requests=(SimpleNamespace(),),
- phase_two_required_presets=("missing-phase-two-preset",),
- )
- service.run_phase_two = AsyncMock() # type: ignore[method-assign]
- return service, snapshots
- @pytest.mark.asyncio
- @pytest.mark.parametrize(
- ("status", "blocked_reason", "operation_status", "publication_closed"),
- [
- (BuildStatus.STOPPED, PHASE_ONE_CAPABILITY_BOUNDARY, None, True),
- (BuildStatus.STOPPING, PHASE_ONE_CAPABILITY_BOUNDARY, None, True),
- (BuildStatus.PARTIAL, "WRONG_BOUNDARY", None, True),
- (BuildStatus.PARTIAL, PHASE_ONE_CAPABILITY_BOUNDARY, OperationStatus.RUNNING, True),
- (BuildStatus.PARTIAL, PHASE_ONE_CAPABILITY_BOUNDARY, None, False),
- ],
- )
- async def test_phase_two_advance_rejects_invalid_entry_without_extending_snapshot(
- status: BuildStatus,
- blocked_reason: str,
- operation_status: OperationStatus | None,
- publication_closed: bool,
- ) -> None:
- service, snapshots = _advance_service(
- status=status,
- blocked_reason=blocked_reason,
- operation_status=operation_status,
- publication_closed=publication_closed,
- )
- with pytest.raises(PhaseTwoBoundaryNotReady):
- await service.advance_to_phase_two(7, Principal("owner"))
- assert snapshots.extend_calls == 0
- @pytest.mark.asyncio
- async def test_concurrent_phase_two_advance_creates_only_one_active_run() -> None:
- state = _AdvanceState(BuildStatus.PARTIAL)
- release_prepare = asyncio.Event()
- release_run = asyncio.Event()
- prepare_calls = 0
- async def prepare(_build: int):
- nonlocal prepare_calls
- prepare_calls += 1
- await release_prepare.wait()
- return _binding(), _snapshot(), 3
- async def run_phase_two(*_args, **_kwargs):
- await release_run.wait()
- service = ScriptMissionService(
- runner=SimpleNamespace(),
- coordinator=SimpleNamespace(
- root_completion=lambda _root: _async(
- {
- "status": TaskStatus.BLOCKED.value,
- "blocked_reason": PHASE_ONE_CAPABILITY_BOUNDARY,
- }
- )
- ),
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
- legacy_state=state,
- authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
- direction_reconciler=SimpleNamespace(),
- )
- service._prepare_phase_two_transition = prepare # type: ignore[method-assign]
- service.run_phase_two = run_phase_two # type: ignore[method-assign]
- first = asyncio.create_task(service.advance_to_phase_two(7, Principal("owner")))
- second = asyncio.create_task(service.advance_to_phase_two(7, Principal("owner")))
- await asyncio.sleep(0)
- release_prepare.set()
- first_result, second_result = await asyncio.gather(first, second)
- assert prepare_calls == 1
- assert first_result.status is BuildStatus.RUNNING
- assert second_result.status is BuildStatus.RUNNING
- assert first_result.root_trace_id == second_result.root_trace_id == "root"
- release_run.set()
- await asyncio.sleep(0)
- @pytest.mark.asyncio
- async def test_http_advance_uses_safe_reentry_after_policy_was_unblocked() -> None:
- state = _AdvanceState(BuildStatus.PARTIAL)
- release_run = asyncio.Event()
- coordinator = SimpleNamespace(
- root_completion=lambda _root: _async(
- {"status": TaskStatus.NEEDS_REPLAN.value, "blocked_reason": None}
- )
- )
- service = ScriptMissionService(
- runner=SimpleNamespace(),
- coordinator=coordinator,
- factory=ScriptMissionFactory(),
- input_snapshots=SimpleNamespace(),
- bindings=SimpleNamespace(get_by_build=lambda _build: _async(_binding())),
- legacy_state=state,
- authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async(None)),
- direction_reconciler=SimpleNamespace(),
- )
- service._prepare_phase_two_transition = AsyncMock( # type: ignore[method-assign]
- side_effect=AssertionError("normal Phase1 boundary path must not be used")
- )
- service._prepare_phase_two_reentry = AsyncMock( # type: ignore[method-assign]
- return_value=(_binding(), _snapshot(), 3)
- )
- async def run_phase_two(*_args, **_kwargs):
- await release_run.wait()
- service.run_phase_two = run_phase_two # type: ignore[method-assign]
- result = await service.advance_to_phase_two(7, Principal("owner"))
- assert result.status is BuildStatus.RUNNING
- service._prepare_phase_two_reentry.assert_awaited_once_with(7) # type: ignore[attr-defined]
- service._prepare_phase_two_transition.assert_not_awaited() # type: ignore[attr-defined]
- release_run.set()
- await asyncio.sleep(0)
- def test_mission_observation_dto_bounds_text_and_omits_internal_payloads() -> None:
- secret = "must-not-leak"
- artifact_ref = SimpleNamespace(
- uri="script-build://artifact-versions/1",
- kind="paragraph",
- version="1",
- digest="sha256:" + "a" * 64,
- )
- task = SimpleNamespace(
- task_id="task-1",
- parent_task_id="root",
- display_path="root/task-1",
- status=TaskStatus.COMPLETED,
- current_spec=SimpleNamespace(
- version=1,
- objective="candidate body " * 500,
- acceptance_criteria=[],
- context_refs=("protected://" + secret,),
- ),
- child_task_ids=[],
- attempt_ids=["attempt-1"],
- validation_ids=[],
- decision_ids=["decision-1"],
- blocked_reason=None,
- superseded_by=None,
- created_at="now",
- updated_at="now",
- )
- attempt = SimpleNamespace(
- attempt_id="attempt-1",
- task_id="task-1",
- spec_version=1,
- worker_trace_id="worker",
- worker_preset="script_paragraph_worker",
- status="submitted",
- operation_id="operation-1",
- snapshot_id="snapshot-1",
- accepted_child_decision_ids=[],
- submission=SimpleNamespace(
- summary="candidate body " * 500,
- artifact_refs=[artifact_ref],
- evidence_refs=[],
- ),
- error=None,
- created_at="now",
- updated_at="now",
- protected_context={"token": secret},
- )
- decision = SimpleNamespace(
- decision_id="decision-1",
- task_id="task-1",
- action=DecisionAction.ACCEPT,
- reason="accepted",
- from_status="awaiting_decision",
- to_status="completed",
- attempt_id="attempt-1",
- validation_id=None,
- created_at="now",
- payload={"tool_arguments": secret},
- )
- operation = SimpleNamespace(
- operation_id="operation-1",
- kind="dispatch",
- status=OperationStatus.COMPLETED,
- task_ids=["task-1"],
- attempt_ids=["attempt-1"],
- validation_ids=[],
- deadline_at=None,
- error=None,
- created_at="now",
- updated_at="now",
- request={"arguments": secret},
- )
- ledger = SimpleNamespace(
- revision=1,
- root_trace_id="root",
- root_task_id="root",
- root_objective="root",
- focused_task_id=None,
- tasks={"task-1": task},
- attempts={"attempt-1": attempt},
- validations={},
- decisions={"decision-1": decision},
- operations={"operation-1": operation},
- protected_context={"token": secret},
- )
- view = mission_snapshot_view(ledger)
- rendered = str(view)
- assert secret not in rendered
- assert view["tasks"][0]["current_spec"]["context_refs"] == ["[internal-ref]"]
- assert len(view["tasks"][0]["current_spec"]["objective"]) == 2_000
- assert len(view["attempts"][0]["submission"]["summary"]) == 500
|