Просмотр исходного кода

任务:实现动态规划、阶段边界和协作停止语义

创建无固定 Workflow、Round、Branch 的 Root TaskSpec,使用显式验证模式启动 Planner 并支持 objective revise。

Direction 接受后幂等发布并将 Root 以阶段能力边界阻塞、构建投影为 partial;停止流程先持久化 stopping,再终止 Operation 和 Trace。
SamLee 20 часов назад
Родитель
Сommit
4110b7092e

+ 151 - 0
script_build_host/src/script_build_host/application/mission_factory.py

@@ -0,0 +1,151 @@
+"""Pure creation of the phase-one Root task and Planner RunConfig."""
+
+from __future__ import annotations
+
+import json
+from math import isfinite
+
+from agent import CompletionPolicy, RunConfig
+from agent.tools.builtin.knowledge import KnowledgeConfig
+
+from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
+from script_build_host.domain.records import MissionBinding
+
+
+def normalize_topic_summary(value: object, *, max_length: int = 160) -> str:
+    """Extract a bounded topic summary without inventing a title field."""
+
+    if isinstance(value, str):
+        text = " ".join(value.split())
+    elif value is None:
+        text = ""
+    else:
+        text = json.dumps(value, ensure_ascii=False, sort_keys=True, default=str)
+    return (text or "未提供选题摘要")[:max_length]
+
+
+class ScriptMissionFactory:
+    def build_root_task_spec(self, snapshot: ScriptBuildInputSnapshotV1) -> dict[str, object]:
+        topic_row = snapshot.topic.get("topic")
+        if not isinstance(topic_row, dict):
+            topic_row = {}
+        topic_summary = normalize_topic_summary(
+            topic_row.get("result") or topic_row.get("topic_direction") or snapshot.topic_id
+        )
+        account_name = str(snapshot.account.get("account_name") or "未指定账号")
+        criteria: list[dict[str, object]] = [
+            {
+                "criterion_id": "input-consistency",
+                "description": "全部创作可追溯到冻结输入",
+                "hard": True,
+            },
+            {
+                "criterion_id": "direction-accepted",
+                "description": "使用已验证并接受的创作方向",
+                "hard": True,
+            },
+            {
+                "criterion_id": "structure-complete",
+                "description": "段落、元素、关系完整且引用有效",
+                "hard": True,
+            },
+            {
+                "criterion_id": "evidence-closure",
+                "description": "事实表达有证据且创作表达不伪装成事实",
+                "hard": True,
+            },
+            {
+                "criterion_id": "legacy-readable",
+                "description": "能规范化投影为旧详情合同",
+                "hard": True,
+            },
+        ]
+        if account_name != "未指定账号":
+            criteria.insert(
+                4,
+                {
+                    "criterion_id": "persona-consistency",
+                    "description": "语言和观点符合冻结人设",
+                    "hard": True,
+                },
+            )
+        return {
+            "objective": (
+                f"为账号「{account_name}」基于选题「{topic_summary}」构建一份满足"
+                "选题、人设、策略、证据与旧结构化脚本输出合同的正式脚本候选"
+            ),
+            "acceptance_criteria": criteria,
+            "context_refs": [f"script-build://inputs/{snapshot.snapshot_id}"],
+        }
+
+    def build_planner_message(
+        self,
+        binding: MissionBinding,
+        snapshot: ScriptBuildInputSnapshotV1,
+    ) -> str:
+        return json.dumps(
+            {
+                "mission": "plan_toward_root",
+                "script_build_id": binding.script_build_id,
+                "input_snapshot_ref": f"script-build://inputs/{snapshot.snapshot_id}",
+                "input_sha256": snapshot.canonical_sha256,
+                "phase": 1,
+                "phase_boundary": "PHASE_ONE_CAPABILITY_BOUNDARY",
+                "instruction": (
+                    "Inspect the ledger and dynamically plan retrieval and direction work. "
+                    "After direction ACCEPT, block Root at the phase-one boundary."
+                ),
+            },
+            ensure_ascii=False,
+            sort_keys=True,
+        )
+
+    def build_run_config(
+        self,
+        binding: MissionBinding,
+        snapshot: ScriptBuildInputSnapshotV1,
+    ) -> RunConfig:
+        planner_config = _planner_model_config(snapshot.model_manifest)
+        return RunConfig(
+            agent_type="script_planner",
+            model=planner_config["model"],
+            temperature=planner_config["temperature"],
+            max_iterations=planner_config["max_iterations"],
+            completion_policy=CompletionPolicy.EXPLICIT_VALIDATION,
+            new_trace_id=binding.root_trace_id,
+            root_task_spec=self.build_root_task_spec(snapshot),
+            tool_groups=None,
+            parallel_tool_execution=False,
+            enable_memory=False,
+            enable_research_flow=False,
+            knowledge=KnowledgeConfig(
+                enable_extraction=False,
+                enable_completion_extraction=False,
+                enable_injection=False,
+            ),
+            context={
+                "script_build_id": binding.script_build_id,
+                "input_snapshot_id": snapshot.snapshot_id,
+            },
+            name=f"Script build {binding.script_build_id}",
+        )
+
+
+def _planner_model_config(manifest: dict[str, object]) -> dict[str, object]:
+    raw_presets = manifest.get("presets", manifest)
+    raw = raw_presets.get("script_planner", {}) if isinstance(raw_presets, dict) else {}
+    if not isinstance(raw, dict):
+        raw = {}
+    model = str(raw.get("model") or "gpt-4o")
+    temperature = float(raw.get("temperature", 0.3))
+    max_iterations = int(raw.get("max_iterations", 100))
+    if not model.strip() or not isfinite(temperature) or temperature < 0 or max_iterations < 1:
+        raise ValueError("frozen script_planner model configuration is invalid")
+    return {
+        "model": model,
+        "temperature": temperature,
+        "max_iterations": max_iterations,
+    }
+
+
+__all__ = ["ScriptMissionFactory", "normalize_topic_summary"]

+ 456 - 0
script_build_host/src/script_build_host/application/mission_service.py

@@ -0,0 +1,456 @@
+"""Phase-one mission lifecycle and accepted-direction reconciliation."""
+
+from __future__ import annotations
+
+import asyncio
+from collections.abc import AsyncIterator
+from contextlib import asynccontextmanager
+from dataclasses import dataclass, field
+from typing import Any
+from uuid import uuid4
+
+from agent.orchestration import DecisionAction, OperationStatus, TaskStatus, ValidationVerdict
+
+from script_build_host.agents.validation import task_kind
+from script_build_host.domain.artifacts import ArtifactKind, ScriptDirectionArtifactV1
+from script_build_host.domain.errors import ProtocolViolation
+from script_build_host.domain.ports import (
+    BuildAuthorizer,
+    LegacyBuildStateRepository,
+    MissionBindingRepository,
+    PromptRequest,
+    PublicationRepository,
+    RuntimeManifestProvider,
+    ScriptBusinessArtifactRepository,
+)
+from script_build_host.domain.records import (
+    BuildStatus,
+    Principal,
+    PublicationState,
+    PublicationType,
+)
+from script_build_host.infrastructure.redaction import redact
+
+from .input_snapshot_service import AssembleInputRequest, ScriptInputSnapshotService
+from .mission_factory import ScriptMissionFactory
+
+PHASE_ONE_CAPABILITY_BOUNDARY = "PHASE_ONE_CAPABILITY_BOUNDARY"
+
+
+class BuildTransitionGate:
+    """Single-process exclusion between durable stop intent and publication."""
+
+    def __init__(self) -> None:
+        self._locks: dict[int, asyncio.Lock] = {}
+
+    @asynccontextmanager
+    async def hold(self, script_build_id: int) -> AsyncIterator[None]:
+        lock = self._locks.setdefault(script_build_id, asyncio.Lock())
+        async with lock:
+            yield
+
+
+@dataclass(frozen=True, slots=True)
+class StartScriptBuildCommand:
+    execution_id: int
+    topic_build_id: int
+    topic_id: int
+    principal: Principal
+    agent_type: str | None = None
+    agent_config: dict[str, Any] | None = None
+    data_source_url: str | None = None
+    strategies_always_on: tuple[Any, ...] = ()
+    strategies_on_demand: tuple[Any, ...] = ()
+    prompt_requests: tuple[PromptRequest, ...] = ()
+    runtime_prompt_manifest: tuple[dict[str, Any], ...] = ()
+    datasource_manifest: dict[str, Any] = field(default_factory=dict)
+    model_manifest: dict[str, Any] = field(default_factory=dict)
+
+
+@dataclass(frozen=True, slots=True)
+class ScriptMissionStartResult:
+    script_build_id: int
+    status: BuildStatus
+    root_trace_id: str
+    input_snapshot_id: str
+
+
+@dataclass(frozen=True, slots=True)
+class StopResult:
+    script_build_id: int
+    status: BuildStatus
+    stopped_operation_ids: tuple[str, ...] = ()
+
+
+class DirectionReconciler:
+    """Idempotently project one passed and accepted Direction artifact."""
+
+    def __init__(
+        self,
+        *,
+        coordinator: Any,
+        bindings: MissionBindingRepository,
+        artifacts: ScriptBusinessArtifactRepository,
+        publications: PublicationRepository,
+        legacy_state: LegacyBuildStateRepository,
+        transition_gate: BuildTransitionGate | None = None,
+    ) -> None:
+        self.coordinator = coordinator
+        self.bindings = bindings
+        self.artifacts = artifacts
+        self.publications = publications
+        self.legacy_state = legacy_state
+        self.transition_gate = transition_gate or BuildTransitionGate()
+
+    async def reconcile(self, script_build_id: int, root_trace_id: str) -> int:
+        async with self.transition_gate.hold(script_build_id):
+            return await self._reconcile_locked(script_build_id, root_trace_id)
+
+    async def _reconcile_locked(self, script_build_id: int, root_trace_id: str) -> int:
+        status = await self.legacy_state.get_status(script_build_id)
+        if status in {BuildStatus.STOPPING, BuildStatus.STOPPED}:
+            raise ProtocolViolation("direction publication is forbidden while stopping")
+        ledger = await self.coordinator.task_store.load(root_trace_id)
+        candidates = [
+            task
+            for task in ledger.tasks.values()
+            if task_kind(task.current_spec.context_refs) == "direction"
+            and task.status == TaskStatus.COMPLETED
+        ]
+        if len(candidates) != 1:
+            raise ProtocolViolation("phase one requires exactly one accepted direction task")
+        task = candidates[0]
+        decisions = [
+            ledger.decisions[item]
+            for item in task.decision_ids
+            if ledger.decisions[item].action == DecisionAction.ACCEPT
+        ]
+        if len(decisions) != 1:
+            raise ProtocolViolation("direction task requires exactly one ACCEPT decision")
+        decision = decisions[0]
+        if not decision.attempt_id or not decision.validation_id:
+            raise ProtocolViolation("direction ACCEPT is missing attempt or validation identity")
+        attempt = ledger.attempts[decision.attempt_id]
+        validation = ledger.validations[decision.validation_id]
+        if (
+            validation.attempt_id != attempt.attempt_id
+            or validation.verdict != ValidationVerdict.PASSED
+            or validation.snapshot_id != attempt.snapshot_id
+            or attempt.submission is None
+        ):
+            raise ProtocolViolation("direction ACCEPT is not closed over a passed snapshot")
+        refs = [
+            ref
+            for ref in attempt.submission.artifact_refs
+            if ref.kind == ArtifactKind.DIRECTION.value
+        ]
+        if len(refs) != 1:
+            raise ProtocolViolation("direction attempt must submit exactly one direction artifact")
+        ref = refs[0]
+        version = await self.artifacts.read_by_ref(
+            ref,
+            script_build_id=script_build_id,
+            task_id=task.task_id,
+            attempt_id=attempt.attempt_id,
+        )
+        if not isinstance(version.artifact, ScriptDirectionArtifactV1):
+            raise ProtocolViolation("accepted direction reference has the wrong artifact type")
+        publication = await self.publications.prepare(
+            script_build_id=script_build_id,
+            publication_type=PublicationType.DIRECTION,
+            accept_decision_id=decision.decision_id,
+            artifact_version_id=version.artifact_version_id,
+            expected_sha256=version.canonical_sha256,
+        )
+        if publication.state is PublicationState.PUBLISHED:
+            return version.artifact_version_id
+        try:
+            status = await self.legacy_state.get_status(script_build_id)
+            if status in {BuildStatus.STOPPING, BuildStatus.STOPPED}:
+                raise ProtocolViolation("direction publication is forbidden while stopping")
+            await self.legacy_state.project_direction(
+                script_build_id, version.artifact.legacy_markdown
+            )
+            await self.bindings.set_active_direction(
+                script_build_id=script_build_id,
+                artifact_version_id=version.artifact_version_id,
+            )
+            await self.publications.mark_published(publication.publication_id)
+        except Exception as exc:
+            await self.publications.mark_failed(
+                publication.publication_id,
+                error_code="DIRECTION_PROJECTION_FAILED",
+                error_summary=type(exc).__name__,
+            )
+            raise
+        return version.artifact_version_id
+
+
+class ScriptMissionService:
+    """Start, drive and cooperatively stop a single-process phase-one mission."""
+
+    def __init__(
+        self,
+        *,
+        runner: Any,
+        coordinator: Any,
+        factory: ScriptMissionFactory,
+        input_snapshots: ScriptInputSnapshotService,
+        bindings: MissionBindingRepository,
+        legacy_state: LegacyBuildStateRepository,
+        authorizer: BuildAuthorizer,
+        direction_reconciler: DirectionReconciler,
+        engine_version: str = "script-build-host/0.1",
+        schema_version: str = "script-build-input/v1",
+        stop_timeout_seconds: float = 5.0,
+        transition_gate: BuildTransitionGate | None = None,
+        runtime_manifest_provider: RuntimeManifestProvider | None = None,
+    ) -> None:
+        self.runner = runner
+        self.coordinator = coordinator
+        self.factory = factory
+        self.input_snapshots = input_snapshots
+        self.bindings = bindings
+        self.legacy_state = legacy_state
+        self.authorizer = authorizer
+        self.direction_reconciler = direction_reconciler
+        self.engine_version = engine_version
+        self.schema_version = schema_version
+        self.stop_timeout_seconds = stop_timeout_seconds
+        reconciler_gate = getattr(direction_reconciler, "transition_gate", None)
+        self.transition_gate: BuildTransitionGate = transition_gate or (
+            reconciler_gate
+            if isinstance(reconciler_gate, BuildTransitionGate)
+            else BuildTransitionGate()
+        )
+        self.runtime_manifest_provider = runtime_manifest_provider
+        self._runs: dict[int, asyncio.Task[None]] = {}
+
+    async def start(self, command: StartScriptBuildCommand) -> ScriptMissionStartResult:
+        await self.authorizer.require_source_access(
+            command.principal,
+            execution_id=command.execution_id,
+            topic_build_id=command.topic_build_id,
+            topic_id=command.topic_id,
+        )
+        await self.input_snapshots.validate_source(
+            execution_id=command.execution_id,
+            topic_build_id=command.topic_build_id,
+            topic_id=command.topic_id,
+        )
+        safe_config = redact(command.agent_config)
+        safe_strategies = redact(
+            {
+                "always_on": list(command.strategies_always_on),
+                "on_demand": list(command.strategies_on_demand),
+            }
+        )
+        if safe_config is not None and not isinstance(safe_config, dict):
+            raise ProtocolViolation("agent_config must remain an object after redaction")
+        if not isinstance(safe_strategies, dict):
+            raise ProtocolViolation("strategies_config must remain an object after redaction")
+        script_build_id = await self.legacy_state.create(
+            execution_id=command.execution_id,
+            topic_build_id=command.topic_build_id,
+            topic_id=command.topic_id,
+            agent_type=command.agent_type,
+            agent_config=safe_config,
+            data_source_url=command.data_source_url,
+            strategies_config=safe_strategies,
+        )
+        try:
+            runtime_datasources = (
+                await self.runtime_manifest_provider.load()
+                if self.runtime_manifest_provider is not None
+                else {}
+            )
+            assembled = await self.input_snapshots.assemble(
+                AssembleInputRequest(
+                    script_build_id=script_build_id,
+                    execution_id=command.execution_id,
+                    topic_build_id=command.topic_build_id,
+                    topic_id=command.topic_id,
+                    principal=command.principal,
+                    strategies_always_on=command.strategies_always_on,
+                    strategies_on_demand=command.strategies_on_demand,
+                    prompt_requests=command.prompt_requests,
+                    runtime_prompt_manifest=command.runtime_prompt_manifest,
+                    datasource_manifest={
+                        **runtime_datasources,
+                        **command.datasource_manifest,
+                    },
+                    model_manifest=command.model_manifest,
+                )
+            )
+            snapshot = await self.input_snapshots.freeze(assembled)
+            root_trace_id = str(uuid4())
+            binding = await self.bindings.create(
+                script_build_id=script_build_id,
+                root_trace_id=root_trace_id,
+                input_snapshot_id=int(snapshot.snapshot_id),
+                engine_version=self.engine_version,
+                schema_version=self.schema_version,
+            )
+            await self.legacy_state.set_status(script_build_id, BuildStatus.RUNNING)
+        except Exception as exc:
+            await self.legacy_state.set_status(
+                script_build_id,
+                BuildStatus.FAILED,
+                error_summary=type(exc).__name__,
+            )
+            raise
+        task = asyncio.create_task(
+            self.run(script_build_id), name=f"script-build:{script_build_id}"
+        )
+        self._runs[script_build_id] = task
+        task.add_done_callback(lambda done: self._discard_run(script_build_id, done))
+        return ScriptMissionStartResult(
+            script_build_id=script_build_id,
+            status=BuildStatus.RUNNING,
+            root_trace_id=binding.root_trace_id,
+            input_snapshot_id=snapshot.snapshot_id,
+        )
+
+    async def run(self, script_build_id: int) -> None:
+        binding = await self.bindings.get_by_build(script_build_id)
+        snapshot = await self.input_snapshots.get(
+            str(binding.input_snapshot_id), script_build_id=script_build_id
+        )
+        try:
+            await self.runner.run_result(
+                messages=[
+                    {
+                        "role": "user",
+                        "content": self.factory.build_planner_message(binding, snapshot),
+                    }
+                ],
+                config=self.factory.build_run_config(binding, snapshot),
+            )
+            completion = await self.coordinator.root_completion(binding.root_trace_id)
+            if completion["status"] == TaskStatus.COMPLETED.value:
+                raise ProtocolViolation("phase-one Root must not complete")
+            if (
+                completion["status"] != TaskStatus.BLOCKED.value
+                or completion.get("blocked_reason") != PHASE_ONE_CAPABILITY_BOUNDARY
+            ):
+                raise ProtocolViolation("mission did not stop at the phase-one boundary")
+            await self.direction_reconciler.reconcile(script_build_id, binding.root_trace_id)
+            async with self.transition_gate.hold(script_build_id):
+                if await self.legacy_state.get_status(script_build_id) == BuildStatus.STOPPING:
+                    return
+                await self.legacy_state.set_status(script_build_id, BuildStatus.PARTIAL)
+        except Exception as exc:
+            if await self.legacy_state.get_status(script_build_id) != BuildStatus.STOPPING:
+                await self.legacy_state.set_status(
+                    script_build_id,
+                    BuildStatus.FAILED,
+                    error_summary=type(exc).__name__,
+                )
+            raise
+
+    async def stop(self, script_build_id: int, principal: Principal) -> StopResult:
+        await self.authorizer.require_access(principal, script_build_id)
+        async with self.transition_gate.hold(script_build_id):
+            status = await self.legacy_state.get_status(script_build_id)
+            if status == BuildStatus.STOPPED:
+                return StopResult(script_build_id, BuildStatus.STOPPED)
+            await self.legacy_state.set_status(script_build_id, BuildStatus.STOPPING)
+        binding = await self.bindings.get_by_build(script_build_id)
+        ledger = await self.coordinator.task_store.load(binding.root_trace_id)
+        active = [
+            operation
+            for operation in ledger.operations.values()
+            if operation.status
+            in {
+                OperationStatus.PENDING,
+                OperationStatus.RUNNING,
+                OperationStatus.STOP_REQUESTED,
+            }
+        ]
+        for operation in active:
+            await self.coordinator.stop_operation(
+                binding.root_trace_id,
+                operation.operation_id,
+                idempotency_key=f"stop:{script_build_id}:{operation.operation_id}",
+            )
+        await self.runner.stop(binding.root_trace_id)
+        run_task = self._runs.get(script_build_id)
+        try:
+            await asyncio.wait_for(
+                self._wait_stopped(
+                    binding.root_trace_id,
+                    [item.operation_id for item in active],
+                    run_task,
+                ),
+                timeout=self.stop_timeout_seconds,
+            )
+        except TimeoutError:
+            return StopResult(
+                script_build_id,
+                BuildStatus.STOPPING,
+                tuple(item.operation_id for item in active),
+            )
+        await self.legacy_state.set_status(script_build_id, BuildStatus.STOPPED)
+        return StopResult(
+            script_build_id,
+            BuildStatus.STOPPED,
+            tuple(item.operation_id for item in active),
+        )
+
+    async def ensure_dispatch_allowed(self, script_build_id: int) -> None:
+        status = await self.legacy_state.get_status(script_build_id)
+        if status in {BuildStatus.STOPPING, BuildStatus.STOPPED}:
+            raise ProtocolViolation("new dispatch is forbidden after stop intent")
+
+    async def _wait_stopped(
+        self,
+        root_trace_id: str,
+        operation_ids: list[str],
+        run_task: asyncio.Task[None] | None,
+    ) -> None:
+        while True:
+            states = [
+                (await self.coordinator.get_operation(root_trace_id, item)).status
+                for item in operation_ids
+            ]
+            operations_terminal = all(
+                state
+                in {
+                    OperationStatus.STOPPED,
+                    OperationStatus.COMPLETED,
+                    OperationStatus.FAILED,
+                }
+                for state in states
+            )
+            planner_terminal = run_task is None or run_task.done()
+            trace = await self.runner.trace_store.get_trace(root_trace_id)
+            trace_terminal = trace is None or trace.status in {
+                "completed",
+                "failed",
+                "stopped",
+            }
+            if operations_terminal and planner_terminal and trace_terminal:
+                if run_task is not None and not run_task.cancelled():
+                    try:
+                        run_task.exception()
+                    except asyncio.InvalidStateError:
+                        pass
+                return
+            await asyncio.sleep(0.02)
+
+    def _discard_run(self, script_build_id: int, task: asyncio.Task[None]) -> None:
+        if self._runs.get(script_build_id) is task:
+            self._runs.pop(script_build_id, None)
+        if not task.cancelled():
+            task.exception()
+
+
+__all__ = [
+    "PHASE_ONE_CAPABILITY_BOUNDARY",
+    "BuildTransitionGate",
+    "DirectionReconciler",
+    "ScriptMissionService",
+    "ScriptMissionStartResult",
+    "StartScriptBuildCommand",
+    "StopResult",
+]

+ 74 - 0
script_build_host/tests/test_mission_factory.py

@@ -0,0 +1,74 @@
+import json
+from datetime import UTC, datetime
+
+from script_build_host.application.mission_factory import ScriptMissionFactory
+from script_build_host.domain.input_snapshot import ScriptBuildInputSnapshotV1
+from script_build_host.domain.records import MissionBinding
+
+
+def _snapshot() -> ScriptBuildInputSnapshotV1:
+    return ScriptBuildInputSnapshotV1(
+        snapshot_id="11",
+        script_build_id=4,
+        execution_id=1,
+        topic_build_id=2,
+        topic_id=3,
+        topic={
+            "topic": {
+                "id": 3,
+                "result": "真实嵌套选题结果",
+                "topic_direction": "fallback",
+            }
+        },
+        account={"account_name": "每天心理学"},
+        persona_points=(),
+        section_patterns=(),
+        strategies=(),
+        prompt_manifest=(),
+        datasource_manifest={},
+        model_manifest={
+            "presets": {
+                "script_planner": {
+                    "model": "planner-model",
+                    "temperature": 0.1,
+                    "max_iterations": 42,
+                }
+            }
+        },
+        canonical_sha256="sha256:" + "b" * 64,
+        created_at=datetime.now(UTC),
+    )
+
+
+def _binding() -> MissionBinding:
+    now = datetime.now(UTC)
+    return MissionBinding(
+        binding_id=1,
+        script_build_id=4,
+        root_trace_id="root",
+        input_snapshot_id=11,
+        active_direction_artifact_version_id=None,
+        accepted_root_artifact_version_id=None,
+        engine_version="test",
+        schema_version="v1",
+        created_at=now,
+        updated_at=now,
+    )
+
+
+def test_factory_reads_nested_real_topic_and_does_not_double_prefix_digest() -> None:
+    factory = ScriptMissionFactory()
+    spec = factory.build_root_task_spec(_snapshot())
+    assert "真实嵌套选题结果" in str(spec["objective"])
+    assert "title" not in str(spec)
+    assert "workflow" not in str(spec).lower()
+    message = json.loads(factory.build_planner_message(_binding(), _snapshot()))
+    assert message["input_sha256"] == "sha256:" + "b" * 64
+    config = factory.build_run_config(_binding(), _snapshot())
+    assert config.agent_type == "script_planner"
+    assert config.model == "planner-model"
+    assert config.temperature == 0.1
+    assert config.max_iterations == 42
+    assert config.new_trace_id == "root"
+    assert config.enable_memory is False
+    assert config.enable_research_flow is False

+ 213 - 0
script_build_host/tests/test_mission_service_boundaries.py

@@ -0,0 +1,213 @@
+from __future__ import annotations
+
+import asyncio
+from types import SimpleNamespace
+
+import pytest
+
+from script_build_host.application.mission_service import (
+    BuildTransitionGate,
+    DirectionReconciler,
+    ScriptMissionService,
+    StartScriptBuildCommand,
+)
+from script_build_host.domain.errors import InputRelationMismatch, ProtocolViolation
+from script_build_host.domain.records import BuildStatus, Principal
+
+
+class _RejectingInputs:
+    async def validate_source(self, **_: object) -> None:
+        raise InputRelationMismatch()
+
+
+class _SourceAuthorizer:
+    async def require_source_access(self, *_args: object, **_kwargs: object) -> None:
+        return None
+
+    async def require_access(self, *_args: object, **_kwargs: object) -> None:
+        return None
+
+
+class _LegacyState:
+    def __init__(self) -> None:
+        self.created = 0
+
+    async def create(self, **_: object) -> int:
+        self.created += 1
+        return 1
+
+
+@pytest.mark.asyncio
+async def test_relation_mismatch_produces_no_build_or_binding_write() -> None:
+    legacy = _LegacyState()
+    bindings = SimpleNamespace(created=0)
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=SimpleNamespace(),
+        input_snapshots=_RejectingInputs(),  # type: ignore[arg-type]
+        bindings=bindings,
+        legacy_state=legacy,  # type: ignore[arg-type]
+        authorizer=_SourceAuthorizer(),
+        direction_reconciler=SimpleNamespace(),
+    )
+    with pytest.raises(InputRelationMismatch):
+        await service.start(
+            StartScriptBuildCommand(
+                execution_id=1,
+                topic_build_id=2,
+                topic_id=3,
+                principal=Principal("tester"),
+            )
+        )
+    assert legacy.created == 0
+    assert bindings.created == 0
+
+
+@pytest.mark.asyncio
+async def test_source_authorization_runs_before_input_or_build_access() -> None:
+    legacy = _LegacyState()
+    inputs = SimpleNamespace(validated=False)
+
+    class Denied:
+        async def require_source_access(self, *_args: object, **_kwargs: object) -> None:
+            raise PermissionError("source forbidden")
+
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=SimpleNamespace(),
+        input_snapshots=inputs,
+        bindings=SimpleNamespace(),
+        legacy_state=legacy,  # type: ignore[arg-type]
+        authorizer=Denied(),  # type: ignore[arg-type]
+        direction_reconciler=SimpleNamespace(),
+    )
+    with pytest.raises(PermissionError, match="source forbidden"):
+        await service.start(StartScriptBuildCommand(1, 2, 3, Principal("intruder")))
+    assert legacy.created == 0
+    assert inputs.validated is False
+
+
+class _StopState:
+    def __init__(self, status: BuildStatus) -> None:
+        self.status = status
+
+    async def get_status(self, _build: int) -> BuildStatus:
+        return self.status
+
+    async def set_status(self, _build: int, status: BuildStatus, **_: object) -> None:
+        self.status = status
+
+
+@pytest.mark.asyncio
+async def test_stop_waits_for_planner_even_when_there_are_no_operations() -> None:
+    state = _StopState(BuildStatus.RUNNING)
+    pending: asyncio.Future[None] = asyncio.get_running_loop().create_future()
+    trace = SimpleNamespace(status="running")
+    runner = SimpleNamespace(
+        stop=lambda _root: _async_value(True),
+        trace_store=SimpleNamespace(get_trace=lambda _root: _async_value(trace)),
+    )
+    coordinator = SimpleNamespace(
+        task_store=SimpleNamespace(load=lambda _root: _async_value(SimpleNamespace(operations={})))
+    )
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=coordinator,
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(
+            get_by_build=lambda _build: _async_value(SimpleNamespace(root_trace_id="root"))
+        ),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+        stop_timeout_seconds=0.01,
+    )
+    service._runs[7] = pending  # type: ignore[assignment]
+    result = await service.stop(7, Principal("owner"))
+    assert result.status == BuildStatus.STOPPING
+    assert state.status == BuildStatus.STOPPING
+    pending.cancel()
+
+
+@pytest.mark.asyncio
+async def test_stop_is_idempotent_and_publication_is_gated_by_durable_status() -> None:
+    state = _StopState(BuildStatus.STOPPED)
+    service = ScriptMissionService(
+        runner=SimpleNamespace(),
+        coordinator=SimpleNamespace(),
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=SimpleNamespace(),
+    )
+    result = await service.stop(7, Principal("owner"))
+    assert result.status == BuildStatus.STOPPED
+
+    state.status = BuildStatus.STOPPING
+    publications = SimpleNamespace(prepared=False)
+    reconciler = DirectionReconciler(
+        coordinator=SimpleNamespace(),
+        bindings=SimpleNamespace(),
+        artifacts=SimpleNamespace(),
+        publications=publications,
+        legacy_state=state,
+    )
+    with pytest.raises(ProtocolViolation, match="forbidden while stopping"):
+        await reconciler.reconcile(7, "root")
+    assert publications.prepared is False
+
+
+@pytest.mark.asyncio
+async def test_stop_intent_and_direction_reconcile_are_serialized() -> None:
+    gate = BuildTransitionGate()
+    entered = asyncio.Event()
+    release = asyncio.Event()
+    state = _StopState(BuildStatus.RUNNING)
+
+    class BlockingReconciler:
+        transition_gate = gate
+
+        async def reconcile(self, script_build_id: int, _root: str) -> int:
+            async with gate.hold(script_build_id):
+                entered.set()
+                await release.wait()
+                assert state.status is BuildStatus.RUNNING
+                return 1
+
+    runner = SimpleNamespace(
+        stop=lambda _root: _async_value(True),
+        trace_store=SimpleNamespace(get_trace=lambda _root: _async_value(None)),
+    )
+    coordinator = SimpleNamespace(
+        task_store=SimpleNamespace(load=lambda _root: _async_value(SimpleNamespace(operations={})))
+    )
+    service = ScriptMissionService(
+        runner=runner,
+        coordinator=coordinator,
+        factory=SimpleNamespace(),
+        input_snapshots=SimpleNamespace(),
+        bindings=SimpleNamespace(
+            get_by_build=lambda _build: _async_value(SimpleNamespace(root_trace_id="root"))
+        ),
+        legacy_state=state,
+        authorizer=SimpleNamespace(require_access=lambda _principal, _build: _async_value(None)),
+        direction_reconciler=BlockingReconciler(),  # type: ignore[arg-type]
+        transition_gate=gate,
+    )
+    reconcile = asyncio.create_task(service.direction_reconciler.reconcile(7, "root"))
+    await entered.wait()
+    stop = asyncio.create_task(service.stop(7, Principal("owner")))
+    await asyncio.sleep(0)
+    assert state.status is BuildStatus.RUNNING
+    release.set()
+    assert await reconcile == 1
+    assert (await stop).status is BuildStatus.STOPPED
+
+
+async def _async_value(value):
+    return value