Explorar el Código

语义辅助:加入可回退排序、持久配额与公共访问协议

新增可选的 Context Broker 语义排序器,只允许重排 Host 授权 handle;加入内容寻址缓存、超时回退、非法输出防护和跨进程持久的每 Mission 调用配额。公共 Context Access Protocol 以组合片段进入冻结 Prompt 和 digest,统一 Bootstrap、sufficiency、required handles、分页与 NotModified 规则。
SamLee hace 13 horas
padre
commit
c2c57b7a48

+ 1 - 1
script_build_host/src/script_build_host/agents/model_manifest.py

@@ -40,7 +40,7 @@ def resolve_script_role_model_manifest(
         role = str(item["role"])
         default: dict[str, Any] = {
             "model": model,
-            "temperature": 0.0 if role == "validator" else 0.2,
+            "temperature": 0.0 if role in {"validator", "utility"} else 0.2,
             "max_iterations": min(
                 max_iterations,
                 100 if role == "planner" else _VALIDATOR_CAPS.get(preset, 20),

+ 4 - 0
script_build_host/src/script_build_host/agents/prompts/__init__.py

@@ -3,6 +3,8 @@
 from .contracts import (
     COMPARISON_WORKER_PROMPT,
     COMPOSE_WORKER_PROMPT,
+    CONTEXT_ACCESS_PROTOCOL,
+    CONTEXT_BROKER_PROMPT,
     DECODE_RETRIEVAL_WORKER_PROMPT,
     DIRECTION_WORKER_PROMPT,
     ELEMENT_SET_WORKER_PROMPT,
@@ -25,6 +27,8 @@ from .contracts import (
 __all__ = [
     "COMPARISON_WORKER_PROMPT",
     "COMPOSE_WORKER_PROMPT",
+    "CONTEXT_ACCESS_PROTOCOL",
+    "CONTEXT_BROKER_PROMPT",
     "DECODE_RETRIEVAL_WORKER_PROMPT",
     "DIRECTION_WORKER_PROMPT",
     "ELEMENT_SET_WORKER_PROMPT",

+ 21 - 0
script_build_host/src/script_build_host/agents/prompts/context_access_protocol.md

@@ -0,0 +1,21 @@
+## Context Access Protocol
+
+Workbench Bootstrap 是 Host 按当前 Root、角色、Task、Attempt 和冻结闭包生成的受保护上下文。
+先使用 Bootstrap,再按缺口读取;不得猜测、伪造或改写 handle、Decision ID、scope 或冻结输入。
+
+- 先检查 `context_bundle.sufficiency`、`required_handles`、Card、`state_revision` 及数量/省略信息。
+  `sufficient` 时不要重复搜索;`ambiguous` 时使用建议查询或最小必要过滤条件补读;
+  `insufficient` 时明确报告缺口,不得用常识或标签补齐事实。
+- Card 是定位摘要,不是完整正文或充分证据。需要原文、精确字段或完整判断时,使用
+  `read_mission_context(handle=...)`。`outline` 只用于定位,不能代替 `content`
+  成为正文质量或事实证据。
+- `search_mission_context` 只搜索当前角有权闭包。优先用 Goal、Task kind、source type 缩小范围,
+  并在 query 中写明语义 scope;不要为了“确认一遍”遍历整个 Mission。
+- 搜索结果 `has_more=true` 时,仅在仍需后续结果时原样使用 `next_cursor`。不得更换
+  query/filter 后复用 cursor。`NotModified` 表示同一查询的数据未变,应继续使用已有判断。
+- 详情读取返回 `next_cursor` 时,必须在实际需要完整内容时按顺序读到
+  `next_cursor=null` 且 `exhausted=true`。不要重复读取已 exhausted 的 handle。
+- Validator 在提交 `passed` 前必须逐个读完全部 `required_handles` 的 `content`页;未读完就给出
+  PASS 会被 Host 以 `CONTEXT_NOT_EXHAUSTED` 拒绝。
+- Broker 详情分页只是读取已冻结数据,不算新的上游 Retrieval 调用。上游检索次数仍以当前
+  Retrieval Task 的角色规则为准。

+ 4 - 0
script_build_host/src/script_build_host/agents/prompts/context_broker.md

@@ -0,0 +1,4 @@
+你是 Context Broker 的可选语义排序器。
+
+你只能在 Host 已经通过权限、任务闭包和硬过滤的候选 handle 中排序。
+返回严格 JSON:`{"handles":["..."]}`。不得创造 handle,不得改写内容,不得扩大权限。

+ 29 - 14
script_build_host/src/script_build_host/agents/prompts/contracts.py

@@ -14,28 +14,42 @@ def _read(filename: str) -> str:
     return (_ROOT / filename).read_text(encoding="utf-8")
 
 
-PLANNER_PROMPT = _read("script_planner.md")
-PATTERN_RETRIEVAL_WORKER_PROMPT = _read("pattern_retrieval_worker.md")
-DECODE_RETRIEVAL_WORKER_PROMPT = _read("decode_retrieval_worker.md")
-EXTERNAL_RETRIEVAL_WORKER_PROMPT = _read("external_retrieval_worker.md")
-KNOWLEDGE_RETRIEVAL_WORKER_PROMPT = _read("knowledge_retrieval_worker.md")
-DIRECTION_WORKER_PROMPT = _read("direction_worker.md")
-RETRIEVAL_VALIDATOR_PROMPT = _read("retrieval_validator.md")
-SCRIPT_CANDIDATE_VALIDATOR_PROMPT = _read("candidate_validator.md")
-STRUCTURE_WORKER_PROMPT = _read("structure_worker.md")
-PARAGRAPH_WORKER_PROMPT = _read("paragraph_worker.md")
-ELEMENT_SET_WORKER_PROMPT = _read("element_set_worker.md")
-COMPARISON_WORKER_PROMPT = _read("comparison_worker.md")
+CONTEXT_ACCESS_PROTOCOL = _read("context_access_protocol.md")
+
+
+def _role_prompt(filename: str) -> str:
+    return f"{CONTEXT_ACCESS_PROTOCOL}\n\n---\n\n{_read(filename)}"
+
+
+PLANNER_PROMPT = _role_prompt("script_planner.md")
+PATTERN_RETRIEVAL_WORKER_PROMPT = _role_prompt("pattern_retrieval_worker.md")
+DECODE_RETRIEVAL_WORKER_PROMPT = _role_prompt("decode_retrieval_worker.md")
+EXTERNAL_RETRIEVAL_WORKER_PROMPT = _role_prompt("external_retrieval_worker.md")
+KNOWLEDGE_RETRIEVAL_WORKER_PROMPT = _role_prompt("knowledge_retrieval_worker.md")
+DIRECTION_WORKER_PROMPT = _role_prompt("direction_worker.md")
+RETRIEVAL_VALIDATOR_PROMPT = _role_prompt("retrieval_validator.md")
+SCRIPT_CANDIDATE_VALIDATOR_PROMPT = _role_prompt("candidate_validator.md")
+STRUCTURE_WORKER_PROMPT = _role_prompt("structure_worker.md")
+PARAGRAPH_WORKER_PROMPT = _role_prompt("paragraph_worker.md")
+ELEMENT_SET_WORKER_PROMPT = _role_prompt("element_set_worker.md")
+COMPARISON_WORKER_PROMPT = _role_prompt("comparison_worker.md")
 COMPOSE_WORKER_PROMPT = _read("compose_worker.md")
 PORTFOLIO_WORKER_PROMPT = _read("portfolio_worker.md")
-ROOT_WORKER_PROMPT = _read("root_worker.md")
-ROOT_VALIDATOR_PROMPT = _read("root_validator.md")
+ROOT_WORKER_PROMPT = _role_prompt("root_worker.md")
+ROOT_VALIDATOR_PROMPT = _role_prompt("root_validator.md")
+CONTEXT_BROKER_PROMPT = _read("context_broker.md")
 
 
 def script_build_prompt_manifest() -> tuple[dict[str, object], ...]:
     """Return every Prompt a new Phase1+Phase2 InputSnapshot must freeze."""
 
     values = {
+        "script_context_broker": (
+            "utility",
+            "context_broker.md",
+            CONTEXT_BROKER_PROMPT,
+            1,
+        ),
         "script_planner": ("planner", "script_planner.md", PLANNER_PROMPT, 2),
         "script_direction_worker": (
             "worker",
@@ -160,6 +174,7 @@ def phase_one_prompt_manifest() -> tuple[dict[str, object], ...]:
     """Compatibility view used by older callers and Phase1 fixture tests."""
 
     phase_one = {
+        "script_context_broker",
         "script_planner",
         "script_direction_worker",
         "script_pattern_retrieval_worker",

+ 190 - 0
script_build_host/src/script_build_host/application/context_semantics.py

@@ -0,0 +1,190 @@
+"""Optional semantic assistance for Context Broker ranking.
+
+Correctness never depends on this component: it may only reorder handles that
+the deterministic Broker has already authorized.
+"""
+
+from __future__ import annotations
+
+import asyncio
+import json
+import os
+from collections.abc import Mapping, Sequence
+from dataclasses import dataclass
+from hashlib import sha256
+from pathlib import Path
+from tempfile import NamedTemporaryFile
+from typing import Any
+
+from script_build_host.agents.prompts import CONTEXT_BROKER_PROMPT
+from script_build_host.domain.context_broker import (
+    ContextCard,
+    canonical_digest,
+    estimate_context_tokens,
+)
+
+_SYSTEM_PROMPT = CONTEXT_BROKER_PROMPT
+
+
+@dataclass(frozen=True, slots=True)
+class SemanticSelection:
+    handles: tuple[str, ...]
+    cache_hit: bool = False
+    calls: int = 0
+    fallback_reason: str | None = None
+
+
+class AdaptiveSemanticSelector:
+    """Bounded, cached, fail-open reranking of an authorized candidate set."""
+
+    def __init__(
+        self,
+        *,
+        llm_call: Any | None,
+        cache_root: Path | None,
+        mission_call_limit: int = 24,
+        timeout_seconds: float = 15.0,
+    ) -> None:
+        self._llm_call = llm_call
+        self._cache_root = cache_root
+        self._mission_call_limit = mission_call_limit
+        self._timeout_seconds = timeout_seconds
+        self._calls_by_root: dict[str, int] = {}
+        self._limit_lock = asyncio.Lock()
+
+    async def select(
+        self,
+        *,
+        root_trace_id: str,
+        query: str,
+        cards: Sequence[ContextCard],
+        model_config: Mapping[str, Any],
+    ) -> SemanticSelection:
+        eligible = tuple(card.handle for card in cards[:20])
+        if not query.strip() or len(eligible) < 2:
+            return SemanticSelection(eligible, fallback_reason="semantic_not_needed")
+        if self._llm_call is None:
+            return SemanticSelection(eligible, fallback_reason="semantic_provider_unavailable")
+        request: dict[str, Any] = {"query": query[:2_000], "candidates": []}
+        for card in cards[:20]:
+            candidate = {
+                "handle": card.handle,
+                "source_type": card.source_type,
+                "summary": card.summary[:300],
+                "excerpt": card.excerpt[:300],
+            }
+            trial = {**request, "candidates": [*request["candidates"], candidate]}
+            if request["candidates"] and estimate_context_tokens(trial) > 7_000:
+                break
+            request = trial
+        allowed = tuple(item["handle"] for item in request["candidates"])
+        key = canonical_digest(
+            {
+                "request": request,
+                "prompt": sha256(_SYSTEM_PROMPT.encode()).hexdigest(),
+                "model": dict(model_config),
+            }
+        ).removeprefix("sha256:")
+        async with self._limit_lock:
+            cached = await self._read_cache(key)
+            if cached is not None and set(cached).issubset(allowed):
+                return SemanticSelection(tuple(cached), cache_hit=True)
+            if not await self._claim_mission_call(root_trace_id):
+                return SemanticSelection(allowed, fallback_reason="semantic_call_limit")
+        try:
+            result = await asyncio.wait_for(
+                self._llm_call(
+                    messages=[
+                        {"role": "system", "content": _SYSTEM_PROMPT},
+                        {"role": "user", "content": json.dumps(request, ensure_ascii=False)},
+                    ],
+                    model=str(model_config.get("model") or ""),
+                    tools=[],
+                    temperature=0.0,
+                    max_tokens=1_000,
+                ),
+                timeout=self._timeout_seconds,
+            )
+            parsed = json.loads(str(result.get("content") or ""))
+            handles = parsed.get("handles") if isinstance(parsed, Mapping) else None
+            if (
+                not isinstance(handles, list)
+                or not handles
+                or any(not isinstance(item, str) or item not in allowed for item in handles)
+            ):
+                raise ValueError("semantic selector returned unauthorized handles")
+            selected = tuple(dict.fromkeys(handles))
+            selected += tuple(item for item in allowed if item not in selected)
+            await self._write_cache(key, selected)
+            return SemanticSelection(selected, calls=1)
+        except Exception as exc:
+            return SemanticSelection(
+                allowed,
+                calls=1,
+                fallback_reason=f"{type(exc).__name__}:semantic_fallback"[:120],
+            )
+
+    async def _claim_mission_call(self, root_trace_id: str) -> bool:
+        if self._cache_root is None:
+            used = self._calls_by_root.get(root_trace_id, 0)
+            if used >= self._mission_call_limit:
+                return False
+            self._calls_by_root[root_trace_id] = used + 1
+            return True
+        return await asyncio.to_thread(self._claim_persistent_call_sync, root_trace_id)
+
+    def _claim_persistent_call_sync(self, root_trace_id: str) -> bool:
+        assert self._cache_root is not None
+        mission_key = sha256(root_trace_id.encode()).hexdigest()
+        quota_root = self._cache_root / "mission-call-quota" / mission_key
+        quota_root.mkdir(parents=True, exist_ok=True)
+        for index in range(self._mission_call_limit):
+            target = quota_root / f"{index:04d}.claim"
+            try:
+                descriptor = os.open(target, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
+            except FileExistsError:
+                continue
+            with os.fdopen(descriptor, "w", encoding="utf-8") as stream:
+                stream.write("claimed\n")
+                stream.flush()
+                os.fsync(stream.fileno())
+            return True
+        return False
+
+    async def _read_cache(self, key: str) -> tuple[str, ...] | None:
+        if self._cache_root is None:
+            return None
+        target = self._cache_root / f"{key}.json"
+        if not target.is_file():
+            return None
+        try:
+            raw = await asyncio.to_thread(target.read_text, encoding="utf-8")
+            value = json.loads(raw)
+            handles = value.get("handles") if isinstance(value, Mapping) else None
+            if isinstance(handles, list) and all(isinstance(item, str) for item in handles):
+                return tuple(handles)
+        except (OSError, json.JSONDecodeError):
+            return None
+        return None
+
+    async def _write_cache(self, key: str, handles: Sequence[str]) -> None:
+        if self._cache_root is None:
+            return
+        await asyncio.to_thread(self._write_cache_sync, key, handles)
+
+    def _write_cache_sync(self, key: str, handles: Sequence[str]) -> None:
+        assert self._cache_root is not None
+        self._cache_root.mkdir(parents=True, exist_ok=True)
+        target = self._cache_root / f"{key}.json"
+        with NamedTemporaryFile("w", dir=self._cache_root, encoding="utf-8", delete=False) as tmp:
+            json.dump({"handles": list(handles)}, tmp, ensure_ascii=False, sort_keys=True)
+            tmp.flush()
+            os.fsync(tmp.fileno())
+            temporary = Path(tmp.name)
+        try:
+            os.replace(temporary, target)
+        finally:
+            temporary.unlink(missing_ok=True)
+
+
+__all__ = ["AdaptiveSemanticSelector", "SemanticSelection"]