跳转至

第 4 章:Agent Loop 状态机 + LLM 流式

对应 dsh 真实源码:packages/core/agent-loop + packages/llm/llm + packages/llm/llm-deepseek (docs/agent-lifecycle.md、docs/subsystems/llm-streaming.md) 前置:第 1~3 章。产出文件:miniharness/llm/、miniharness/core/agent_loop/ + tests/test_loop.py

早期简化形态

本章代码为教学简化形态,与当前实现存在以下差异(学习时以当前实现为准,见 00-setup §0.6 简化表):

  • FakeLlmAdapter finish reason:本章为字符串("stop" / "tool-calls");实现为对象 {"kind": "stop"} / {"kind": "tool-calls"}(llm/fake.py:58,70)。
  • DeepSeekAdapter SSE:实现走 Anthropic 兼容 Messages 协议(上游 llm-deepseek 已删除 Chat Completions),以 message_stop 为完成点(EOF 未到 message_stop 抛 STREAM_CLOSED)、畸形 SSE 载荷抛 MALFORMED_RESPONSE、带内 error 事件即 provider 失败、HTTP 错误映射完整(401/403→AUTH、quota 措辞→QUOTA、429→RATE_LIMIT、400 上下文→CONTEXT_WINDOW_EXCEEDED 否则 INVALID_REQUEST、≥500→SERVER、其余 HTTP_<status>)、usage 按 Anthropic 拼写归一为 TokenUsage(wire/翻译在 llm/deepseek_messages.py,适配器 llm/deepseek.py;见 llm/protocol.py 的 StreamChunk 判别字段 type)。本章的 AUTH_ERROR/REQUEST_ERROR 二元映射已过时、[DONE]/choices[].delta 形态亦已过时。
  • 空响应:实现已产出 EMPTY_RESPONSE 错误且默认可重试(llm/retry_policy.py 白名单),§4.4 教学正文与 §4.8/§4.9 现已一致。
  • loop 片段:本章 loop.py 的 _append 方法、字符串 reason、扁平 assistant/message 形态均已过时;实现是 ContentBlock 消息对象 + 显式编号 + request/header 事件 + 模型流压缩内嵌 assistant/message(失败 attempt 落 assistant/attempt)(core/agent_loop/agent.py)。
  • 重试接线:真实调用入口必须挂载 apply_retry_planner(llm/retry.py:298),否则 agent/request-error 瀑布不生效(本章 §4.6 真实 API 示例为教学简化、未挂载;真实装配必须先挂)。
  • 时序图:完整时序含 request/header 事件与内嵌 assistant/message 的压缩流写入磁盘(见 core/agent_loop/agent.py 的 requestHeaderLogged 语义)。
  • stream 约定已 async 化(httpx 异步传输):实现签名为 async def stream(self, messages, tools, signal=None) 异步迭代(llm/protocol.py,httpx 原生 asyncio 传输,abort 置位即关闭连接、_aiter_raced 竞速抛 StreamAborted);agent 循环为单一 async 驱动 + followup/steer 同步门面(经进程级常驻单事件循环驱动,core/agent_loop/resident_loop.py);本章的同步 def stream 与同步泵形态已过时。

4.1 这一章要做什么

前三章分别备好了日志、插件骨架和工具管线,这一章把它们接起来:模型在循环里转起来,一个回合完整跑完。两个部分:

  1. llm.py——统一流协议 StreamChunk + LlmAdapter 接口 + DeepSeek 官方 SSE 适配器(httpx 异步传输,不装官方 SDK)
  2. loop.py——turn/step 状态机:inbox 排队、pre-step 拒绝、工具回灌继续、turn 括号闭合

用 FakeLlmAdapter 不需要 API key 就能跑通"文本 + 工具调用"的完整回合——测试和演示都靠它。

4.2 概念:turn/step 时序

sequenceDiagram
  participant U as 用户
  participant A as Agent
  participant D as Driver
  participant S as Session 日志
  participant L as LLM 扩展口
  participant T as tools
  U->>A: followup(content)
  A->>S: turn/start [durable]
  D->>D: claim 输入
  D->>D: pre-step (waterfall)
  alt 拒绝
    D->>S: turn/end(零 step)
  else 进入
    D->>S: step/start → user/message [durable]
    D->>L: request → stream
    L-->>D: StreamChunk*
    D->>S: assistant/message [durable]
    D->>T: tool/call → 管线 → tool/result [durable]
    D->>S: step/end
    alt 还有工具请求
      D->>D: 同 turn 内下一步(回灌结果继续问模型)
    else 无未偿之责
      D->>S: turn/end [durable]
    end
  end

三个要点:

  1. turn 打开于认领输入之前。"被拒绝的尝试"也留下 turn/start + turn/end 的持久化记录——审计要看到"发生过一次尝试",即使它什么都没做。
  2. step = 一次模型请求 + 它调用的工具。工具结果回灌后,同一 turn 内自动再问一次模型(_continue)。所以"一次对话回合"可能包含多次模型请求,这是 agent 循环和普通聊天 API 的本质区别。
  3. 模型可见 ⟺ 已记录(第 1 章那句话在这里体现):user/message 在 pre-step 通过后才 append,模型永远看不到没进日志的输入。

逐箭头走读(对应上图从上到下,编号 1 起、step 每 turn 内重置为 1):

  1. U->>A: followup(content) —— 用户把一段输入投进 inbox(followup 队列,归下一 turn)。
  2. A->>S: turn/start [durable] —— Agent 认领输入前先开 turn,turn/start 立即落日志(durable,即使后面什么都不做也留下括号记录)。
  3. D->>D: claim 输入 —— 从 inbox 的 followup 队列取走用户输入。
  4. D->>D: pre-step (waterfall) —— 在派发前过 pre-step 决策瀑布(干预面/权限的扩展点)。
  5. alt 拒绝:pre-step 返回 reject → 直接 D->>S: turn/end {kind:'blocked'},不落 step/start——这就是"零 step turn":被拒绝的尝试也留痕。
  6. else 进入:pre-step 放行后才 D->>S: step/start → user/message [durable](step 从 1 开始计数)。
  7. D->>L: request → stream —— 向 LLM 扩展口发起流式请求(此步传入的是已持久化的消息投影)。
  8. L-->>D: StreamChunk* —— 模型逐块回流(text-delta / reasoning-delta / tool-call-delta / finish)。
  9. D->>S: assistant/message [durable] —— 收齐后把完整 assistant 消息落日志。
  10. D->>T: tool/call → 管线 → tool/result [durable] —— 若含工具调用:tool/call 先落日志,再走第 3 章管线,tool/result 落日志。
  11. D->>S: step/end —— 本 step 闭合,reason 为 {kind:'stop'|'tool-calls'|'max-tokens'|...}。
  12. alt 还有工具请求:若本 step 产出工具调用,D->>D 同 turn 内进入下一步(step 递增),把结果回灌给模型继续问(_continue)。
  13. else 无未偿之责:无未匹配工具调用 → D->>S: turn/end [durable] 闭合整个回合。

注意:step/end、turn/end 都落在 finally 中(失败也必落日志),这与第 1 章"括号平衡"基座一脉相承。

4.3 代码 step-by-step(llm.py)

步骤 1:StreamChunk 统一协议

STREAM_CHUNK_KINDS = frozenset({
    "block-start", "text-delta", "reasoning-delta",
    "tool-call-delta", "block-end", "usage", "finish",
})

class StreamChunk(dict):
    def __init__(self, kind, **payload):
        if kind not in STREAM_CHUNK_KINDS:
            raise ValueError(f"未知 chunk kind: {kind}")
        super().__init__({"kind": kind, **payload})

为什么要统一协议?因为不同模型厂商的流式格式各不相同(OpenAI 系、Anthropic 系、原生 SSE……),loop 不该关心厂商差异。所有适配器都吐同一种 StreamChunk,loop 只认这一种。

协议硬性规定(真实 dsh 逐条遵守,docs/subsystems/llm-streaming.md 有完整规范):

  • block-end 携带完整块;usage 必须在 finish 之前;finish 之后不再有值
  • 块索引关联交错增量:多块并行时用 index 区分

步骤 2:接口 + 错误收口

class LlmFailure(Exception):
    """统一错误收口:授权 / 请求 / 上下文溢出。"""
    def __init__(self, code, message):
        super().__init__(message)
        self.code = code

class LlmAdapter:
    """Service Definition:Consumer(agent-loop)只依赖这个协议。"""
    provider = "base"
    def stream(self, messages, tools):
        raise NotImplementedError

常规做法的错误处理是"哪个 SDK 抛什么就 catch 什么",不同厂商的报错对象还不一样。dsh 统一为 LlmFailure(code, message):授权失败 AUTH_ERROR、网络/HTTP REQUEST_ERROR、上下文溢出 CONTEXT_WINDOW_EXCEEDED,全部一个异常类型。

真实 dsh 还有 EMPTY_RESPONSE 编码(空响应 = 可重试的规范错误)。实现也覆盖:空响应抛 LlmFailure(EMPTY_RESPONSE) 且默认可重试(见 4.8/4.9 差异表)。

步骤 3:FakeLlmAdapter —— 无 key 也能跑回合

class FakeLlmAdapter(LlmAdapter):
    provider = "fake"

    def __init__(self, tool_call=None, final_text="任务完成。"):
        self._tool = tool_call
        self._text = final_text
        self.calls = 0

    def stream(self, messages, tools):
        self.calls += 1
        if self._tool and self.calls == 1:
            arguments = self._tool.get("arguments", {})
            arguments_text = json.dumps(arguments, ensure_ascii=False)
            yield StreamChunk("block-start", index=0, blockType="tool-call")
            yield StreamChunk("tool-call-delta", index=0, id="call_0",
                              name=self._tool["name"], argumentsDelta=arguments_text)
            yield StreamChunk("block-end", index=0, block={
                "type": "tool-call", "id": "call_0", "name": self._tool["name"],
                "arguments": arguments_text,
            })
            yield StreamChunk("finish", reason="tool_calls")
        else:
            yield StreamChunk("block-start", index=0, blockType="text")
            yield StreamChunk("text-delta", index=0, text=self._text)
            yield StreamChunk("block-end", index=0, block={
                "type": "text", "text": self._text,
            })
            yield StreamChunk("finish", reason="stop")

行为规则:第一次调用返回一次工具调用,之后返回最终文本。这样就能驱动"模型调工具 → 拿到结果 → 再回答"的完整回合,全程不需要真实模型。注意 tool-call 块的 arguments 是 JSON 字符串,与上游 ContentBlock 一致——模型侧协议里参数就是字符串,不是对象。

步骤 4:DeepSeekAdapter —— 官方 SSE(httpx 异步传输)

import httpx

class DeepSeekAdapter(LlmAdapter):
    provider = "deepseek-official"
    CONNECT_TIMEOUT_S = 30.0
    READ_TIMEOUT_S = 300.0   # per-read idle watchdog(同上游 fetch 300s)

    def __init__(self, api_key=None, base_url=None, model="deepseek-chat"):
        self._key = api_key if api_key is not None else os.environ.get("DEEPSEEK_API_KEY", "")
        self._base = (base_url or os.environ.get("DEEPSEEK_BASE_URL", "https://api.deepseek.com/anthropic")).rstrip("/")
        self._model = model

    async def stream(self, messages, tools, signal=None):
        """async 迭代器:httpx 异步传输 + SSE 解析,逐 chunk 产出。
        signal.aborted 置位 → _aiter_raced 在下次取块前抛 StreamAborted,
        退出 async-with 即关闭连接(真取消,无遗留线程)。"""
        abort_event = getattr(signal, "event", None) if signal is not None else None
        body = {"model": self._model, "messages": serialize_messages(messages),
                "max_tokens": 256000, "stream": True}
        if tools:
            body["tools"] = [
                {"name": t["name"], "description": t.get("description", ""),
                 "input_schema": t.get("parameters", {})} for t in tools
            ]
        async with httpx.AsyncClient(timeout=httpx.Timeout(
                self.CONNECT_TIMEOUT_S, read=self.READ_TIMEOUT_S)) as client:
            async with client.stream("POST", self._base + "/chat/completions",
                                     json=body,
                                     headers={"Content-Type": "application/json",
                                              "Authorization": "Bearer " + self._key}) as resp:
                if resp.status_code >= 400:
                    detail = (await resp.aread()).decode("utf-8", "replace")[:500]
                    raise LlmFailure(_http_error_code(resp.status_code, detail),
                                     f"HTTP {resp.status_code}: {detail}")
                async for chunk in self._parse_sse(resp.aiter_lines(), abort_event):
                    yield chunk

    async def _parse_sse(self, aiter_lines, abort_event=None):
        """SSE spec-strict:事件只在空行终结时派发;EOF 处未终止尾部是截断,
        丢弃;[DONE] 之前 EOF → STREAM_CLOSED;畸形载荷 → MALFORMED_RESPONSE。"""
        texts, reasonings, pending = {}, {}, {}
        data_lines = []
        async for line in aiter_lines:
            if line == "":
                if data_lines:
                    data = "\n".join(data_lines)
                    data_lines = []
                    if data == "[DONE]":
                        break
                    piece = json.loads(data)
                    for choice in piece.get("choices", []):
                        delta = choice.get("delta", {})
                        if delta.get("reasoning_content"):
                            reasonings[choice["index"]] = reasonings.get(choice["index"], "") + delta["reasoning_content"]
                        if delta.get("content"):
                            texts[choice["index"]] = texts.get(choice["index"], "") + delta["content"]
                        for tc in delta.get("tool_calls") or []:
                            slot = pending.setdefault(tc["index"], {"id": "", "name": "", "arguments": ""})
                            fn = tc.get("function", {})
                            slot["id"] = tc.get("id") or slot["id"]
                            slot["name"] += fn.get("name", "")
                            slot["arguments"] += fn.get("arguments", "")
                continue
            if line.startswith("data:"):
                payload = line[5:]
                if payload.startswith(" "):
                    payload = payload[1:]
                data_lines.append(payload)

        for idx in sorted(texts):
            yield StreamChunk("block-start", index=idx, blockType="text")
            yield StreamChunk("text-delta", index=idx, text=texts[idx])
            yield StreamChunk("block-end", index=idx, block={"type": "text", "text": texts[idx]})
        for idx in sorted(reasonings):
            yield StreamChunk("block-start", index=idx, blockType="reasoning")
            yield StreamChunk("reasoning-delta", index=idx, text=reasonings[idx])
            yield StreamChunk("block-end", index=idx, block={"type": "reasoning", "text": reasonings[idx]})
        if pending:
            for idx, slot in sorted(pending.items()):
                call_id = slot["id"] or f"call_{idx}"
                yield StreamChunk("block-start", index=idx, blockType="tool-call")
                yield StreamChunk("tool-call-delta", index=idx, id=call_id,
                                  name=slot["name"], argumentsDelta=slot["arguments"])
                yield StreamChunk("block-end", index=idx, block={
                    "type": "tool-call", "id": call_id, "name": slot["name"], "arguments": slot["arguments"],
                })
        yield StreamChunk("finish", reason="tool_calls" if pending else "stop")

核心逻辑在 SSE 循环里:reasoning_content、content、tool_calls 三类增量各自累积,其中 tool-call 的 name 和 arguments 是分片到达的,必须攒齐。攒齐之后按 ContentBlock 重放为 block-start → delta* → block-end——这样 Consumer(loop)看到的永远是完整的块结构,而不是碎片。

两个细节:

  • baseURL / apiKey 从环境变量读取,代码里只存引用。凭据的完整处理在第 6 章(凭据扩展口)。
  • 错误在源头就分类(_http_error_code,同上游 adapter.ts):401/403→AUTH、quota 措辞→QUOTA、429→RATE_LIMIT、400 上下文超限→CONTEXT_WINDOW_EXCEEDED、500+→SERVER、其余→HTTP_<status>,并携带 status / providerRetryAfterMs / requestId 事实——调用方不用猜。

传输层为什么用 httpx 而不是 urllib?urllib 的阻塞读无法被真正取消,只能靠固定 120s 超时兜底(producer 线程滞留),与上游"per-read 300s watchdog + fetch 流中断可取消"的语义有差距。httpx 是原生 asyncio 传输:abort 置位即抛 StreamAborted 并关闭连接,READ_TIMEOUT_S=300 同上游每读间隙超时,且测试可用 httpx.MockTransport 注入、无需打真实网络。这是本手册"成熟开源库优先,无语义等价库处才手写"的典型取舍——存在维护良好的异步 HTTP 库时,不手写传输层。

4.4 代码 step-by-step(loop.py)

步骤 1:状态与入口

class AgentLoop:
    def __init__(self, session, adapter, tools, ctx, system_prompt="你是一个助手。", max_steps=None):
        self.session = session
        self.adapter = adapter
        self.tools = tools
        self.ctx = ctx
        self.system_prompt = system_prompt
        # 死循环守卫:未显式传入时取 DEFAULT_MAX_STEPS(默认 50),
        # 可被环境变量 MINIHARNESS_MAX_STEPS 覆盖(极大值即等效关闭)
        self.max_steps = max_steps or DEFAULT_MAX_STEPS
        self.status = "idle"
        self.inbox = deque()        # 排队输入(唯一入口)
        self._turn_open = False
        self._continue = False      # 工具回灌后是否继续

    def followup(self, content, source="user"):
        """用户输入:先进 inbox,pre-step 通过后才 append 进日志。"""
        self.inbox.append({"role": "user", "content": content, "source": source})
        self._pump()

    def run(self, content):
        self.followup(content)
        return self.last_response()

inbox 是唯一入口:用户的输入先进队列,由 _pump 消费。为什么排队而不是直接处理?因为一次 followup 可能触发多轮工具调用,期间再来新输入必须排队,不能打断当前回合。

步骤 2:turn 生命周期

    def _open_turn(self):
        if self._turn_open:
            return
        self.status = "running"
        self.session.append({"type": "turn/start"})
        self._turn_open = True

    def _close_turn(self, reason="completed"):
        if not self._turn_open:
            return
        self.session.append({"type": "turn/end", "reason": reason})
        self._turn_open = False
        self.status = "idle"

turn/start / turn/end 是第 1 章的括号。reason 默认为 "completed",被拒绝的尝试、被打断的回合会写别的值。

步骤 3:主循环

    def _pump(self):
        steps = 0
        while self.inbox or self._continue:
            steps += 1
            if steps > self.max_steps:
                raise RuntimeError(f"超过最大 step 数 {self.max_steps}(DEFAULT_MAX_STEPS,可经 MINIHARNESS_MAX_STEPS 或构造参数覆盖),疑似死循环")
            self._open_turn()
            claimed = self.inbox.popleft() if self.inbox else None
            self._run_step(claimed)
            if not self.inbox and not self._continue:
                self._close_turn()

循环条件:有排队输入,或有工具回灌待继续(_continue)。max_steps 是死循环守卫:模型如果永远调工具不结束,会在上限步数时报错而不是挂死(测试 test_max_steps_guard 固定)。守卫默认 DEFAULT_MAX_STEPS = 50,可由构造参数或环境变量 MINIHARNESS_MAX_STEPS 覆盖——上游 AgentLoop 无硬性 step 上限,mini 保留该保守守卫作安全网。

步骤 4:一个 step

    def _run_step(self, claimed):
        if claimed is not None:
            decision = self.ctx.waterfall("agent/pre-step", {"messages": [claimed]})
            if isinstance(decision, dict) and decision.get("verdict") == "reject":
                return   # 零 step 尝试:turn 照常闭合
            self._step += 1
            self._append("step/start")
            self._append("user/message", content=claimed["content"],
                         surfaceOp="append", source=claimed.get("source", "user"))
        else:
            self._step += 1
            self._append("step/start")   # 工具回灌后的继续

        history = derive_messages(self.session.events)
        messages = [{"role": "system", "content": self.system_prompt}] + history
        chunks = list(self.adapter.stream(messages, self._tool_definitions()))
        text = "".join(c.get("text", "") for c in chunks if c["kind"] == "text-delta")

        # tool-call-delta 是增量分片:按 (id) 累积 name 与 argumentsDelta
        pending_calls = {}
        for c in chunks:
            if c["kind"] != "tool-call-delta":
                continue
            key = c.get("id") or str(c.get("index"))
            slot = pending_calls.setdefault(key, {"name": "", "argumentsDelta": ""})
            slot["name"] += c.get("name", "")
            slot["argumentsDelta"] += c.get("argumentsDelta", "")
        tool_calls = [{"name": s["name"], "arguments": s["argumentsDelta"]} for s in pending_calls.values()]

        self._append("assistant/message", content=text, surfaceOp="append", toolCalls=tool_calls)
        for call in tool_calls:
            self._run_tool(call["name"], call["arguments"])
        self._append("step/end")
        self._continue = bool(tool_calls)

四个要点:

  • _append 是回合事件的统一入口,自动注入 turn / step 编号——与上游一致,从 1 起(session/invariant.ts nextTurn: 1, nextStep: 1,每 turn 内 step 重置为 1)。这些字段与上游完全一致。
  • pre-step 拒绝:waterfall 返回 {"verdict": "reject"} → 不落 step/start,turn 直接闭合。这就是"零 step turn":被拒绝的尝试也留下括号痕迹(4.2 要点 1)。
  • 历史 = derive_messages(日志) + system prompt,绝不另存——第 1 章的投影在这里消费。
  • tool-call-delta 是增量分片,loop 按 id 累积 name 与 argumentsDelta,组装成完整的 toolCalls 再落日志。所以日志里存的是完整参数(JSON 字符串),而不是碎片。
  • _continue = bool(tool_calls):有工具调用 → 同 turn 内再问模型。第二次 adapter 调用时,消息历史里已经多了 tool/result。

步骤 5:工具执行与落日志

    def _run_tool(self, name, arguments):
        tool = self.tools.resolve(name)
        if tool is None:
            result = ToolResult(ok=False, is_error=True, error=f"未知工具: {name}")
        else:
            if isinstance(arguments, str):
                try:
                    arguments = json.loads(arguments)
                except json.JSONDecodeError:
                    arguments = {}
            self.session.append({"type": "tool/call", "name": name, "arguments": arguments})  # 执行前先记录
            result = run_pipeline(self.ctx, tool, arguments)
        self.session.append({"type": "tool/result", "name": name, "content": result.content,
                             "isError": result.is_error, "error": result.error, "surfaceOp": "append"})

tool/call 在执行前落日志(durable),tool/result 是唯一模型面向的结果——和第 3 章管线无缝对接。注意"未知工具"也被规范化成 ToolResult(is_error=True):模型收到错误消息并自己决定怎么办,而不是整个回合崩溃。

4.5 验收:硬性规定 + 测试

tests/test_loop.py 固定的规定:

  1. turn/start 与 turn/end 成对且 turn_balance == 0
  2. 拒绝的尝试:有 turn/start + turn/end,无 step/start
  3. 工具调用回合:tool/call → tool/result 相邻且都 durable;模型在同 turn 内被请求 ≥2 次
  4. 历史永远从日志派生;max_steps 防死循环
  5. StreamChunk 协议:finish 是最后一个 chunk
python -m unittest tests.test_loop -v

4.6 用真实 API 跑一次(可选)

from miniharness import Context, Session, ToolRegistry, Tool, AgentLoop, DeepSeekAdapter

session = Session("real-001")
ctx = Context()
reg = ToolRegistry(ctx)
reg.register(Tool(name="bash", description="Run a shell command.",
                  parameters={"type": "object", "properties": {"cmd": {"type": "string"}}, "required": ["cmd"]},
                  execute=lambda args, e: f"stdout: {args['cmd']}"))
loop = AgentLoop(session, DeepSeekAdapter(model="deepseek-chat"), reg, ctx)
print(loop.run("用 bash 执行 echo hello,然后告诉我结果。"))
print([e["type"] for e in session.events])

需要环境变量 DEEPSEEK_API_KEY(可选 DEEPSEEK_BASE_URL 指向兼容代理)。真实装配请在构造 loop 前调用 apply_retry_planner(ctx)(§4.9),本示例为教学简化省略。跑完后 session.events 里能看到完整回合:turn/step 括号、消息、工具调用与结果,全都在。

4.7 检查点练习

  1. 加拒绝理由:让 pre-step 的 reject 带 reason,reject 时把它写进 turn/end 的 reason 字段,并断言日志可审计。
  2. 流内嵌往返:构造 assistant/message 的内嵌 stream 记录,断言 expand_assistant_stream 能精确还原逐 chunk 序列(含每个 delta 边界与相对时间),且 derive_messages 不受流记录影响。
  3. 并发工具:给 _run_tool 加 is_concurrency_safe 并行执行(线程),非安全工具串行——跑通测试。

4.8 回到 dsh:真实源码对照

打开 deepseek-harness/packages/core/agent-loop/src:

  • 主 Driver 的 _run_step 对应我们的 _run_step——真实实现更复杂(交错工具批、屏障、回灌顺序)
  • docs/agent-lifecycle.md 顶部的 Mermaid 时序图:与我们 4.2 的图逐条对应

下面这些真实扩展点简化版没有实现,对照时不要找"上游为什么多这些东西",它们是刻意省略的:

上游事件/扩展点 用途 mini 对应
system-prompt/assemble waterfall 提示词按片段组装(hook 可注入上下文) 已实现(core/system_prompt.py assemble + contexts/tools/variables 提供器,第 13 章)
agent/request waterfall → llm/stream 请求构造拦截(steering) 已实现(_request_config:seed=路由 config,payload 附 turn/step/signal,provider 缺失 fail loud)
agent/request-error waterfall 规范错误(如上下文溢出)后的重试决策 已实现(§4.9 重试/退避)
agent/turn-stopping serial turn 结束前串行终点检查 已实现(serial/aserial;step/end 后 next-step 空才派发;pre-step 拒绝→blocked 终局不派发;max-tokens 粘滞不降级)
finish {kind:'error'\|'aborted'} 带内失败 流中途失败也可经协议传递 已实现(带内错误/中止与异常路径同走 agent/request-error waterfall,§4.9)
EMPTY_RESPONSE 编码 空响应 = 规范错误,可重试 已实现且默认可重试(§4.9)

4.9 重试/退避与上下文溢出降级

对应 dsh:packages/llm/llm/src/retry-policy.ts + packages/llm/llm-retry/src/index.ts + packages/core/agent/src/runtime-types.ts(agent/request-error)。

扩展点:loop 在适配器抛 LlmFailure 时派发 agent/request-error waterfall, payload {agent, turn, step, provider, failure, retryPolicy, signal}(与上游逐字段 一致)。监听器返回 {kind:'retry'} 且不调 next() = 自己接管恢复;调 next() 委派; 默认 undefined 失败终局。重试规划器由装配方显式挂载(AgentLoop 构造无副作用, 同上游插件 apply 时挂载):headless / sessions / acp / sdk / demo / 示例在构造 loop 前调用 apply_retry_planner(ctx)(幂等,可重复调用)。

策略解析(llm/retry_policy.py,同 retry-policy.ts):

  • 两种模式:normal(maxRetries + retryableCodes 白名单)/ always(无限重试)
  • 默认:maxRetries 5、initialDelayMs 500、maxDelayMs 10000、jitterRatio 0.1、 可重试码 [EMPTY_RESPONSE, RATE_LIMIT, SERVER, TIMEOUT, TRANSPORT]
  • 严格校验:未知键拒绝、backoff 正有限且 initial ≤ max、jitter ∈ [0,1]、 maxRetries 非负整数、codes 非空无重复;解析结果冻结,provider 注册时捕获

恢复决策(llm/retry.py,同 llm-retry/index.ts):

  1. 策略 undefined → 直接委派(不重试)
  2. always:派发前检查熔合信号——已中止则失败终局;先委派下游——下游给出 retry 决策后复查熔合信号,中止胜过决策;下游监听器抛错经 logger.warn 容错、 按未接管处理;失败/未接管/被中止压制后自己无限重试(不判 code)
  3. normal:code 不在白名单 → 委派;同 turn/step/provider/policyKey 的 llm/retry 计数 ≥ maxRetries → 委派(放弃)
  4. 派发前检查熔合信号:请求 signal 与重试插件 lifetime 任一已中止 → 连 llm/retry 都不落、直接委派(上游 backoff 首行检查语义);否则先落 llm/retry(durable,含策略细节/retryId/序数/delayMs/failure 快照), 可取消等待结束后落 llm/retry-started 并返回 {kind:'retry'};retryId 同一对全程复用
  5. 延迟决议:providerRetryAfterMs(429 的 Retry-After,纯数字秒 ×1000 或 HTTP-date) 有效时优先——超过 maxDelayMs 则 normal 放弃 / always 改用本地延迟;否则本地退避 min(initial × 2^min(retry-1, 1024), max) × (1 - ratio + 2×ratio×rand) 再封顶 maxDelayMs
  6. 可取消:等待为事件驱动多信号竞速 asyncio.wait(等价上游 AbortSignal.any([signal, lifetime.signal])——请求 signal 是 loop 每 phase 新建的 取消 Event(agent.ts:325 每 phase 新建 AbortController 的载体对应),lifetime 是 重试插件自身的生命周期信号;任一 .event 置位即醒;无 .event 的裸测试替身信号 回退分片轮询 .aborted);normal 分支同样在派发前检查 abort(与上游一致)

生命周期与拆解:重试插件经 ctx.on("agent/request-error", ...) 挂载监听器, 并登记 effect teardown(label 'llm-retry: abort and drain active recovery'): 拆解时注销监听器 + lifetime.abort + 排干在途恢复(gather(..., return_exceptions=True) 即 allSettled)。已拆解后的迟到回调命中陈旧守卫直接返回(不再进入下游策略)。

接线语义:重试是同 step 内重新发起模型请求——messages 不变(失败 attempt 内嵌流落 assistant/attempt、不产生任何消息事件,derive_messages 不受 llm/retry 影响)、request/header 只落一次(上游仅在 header 变化时追加)、 assistant/message 内嵌成功 attempt 的压缩流且不带 sourceEventSeqs。LlmFailure 扩展 status / providerRetryAfterMs / requestId(x-request-id / x-deepseek-request-id)可选字段;socket 超时映射 TIMEOUT。

上下文溢出降级:CONTEXT_WINDOW_EXCEEDED(400 上下文超限)不在默认白名单 → 重试规划器不接管,委派下游。装配方在 apply_retry_planner(ctx) 之后挂载 install_compaction(ctx)(幂等):压缩引擎监听 agent/request-error,对 CONTEXT_WINDOW_EXCEEDED 强制减容(见 miniharness/compaction/ 与报告 04 §9.4),且仅当 surface replaceGeneration 前进(检查点真实写入磁盘)才返回 {kind:'retry'},计数上限 maxOverflowRetries,成功响应/回合结束边界复位。既无压缩也无接管 → 终局 turn/end reason 为 {kind:'error'}。

与压缩互补的还有一个可选、不送模型的 tool-result 裁剪服务: install_tool_result_pruner(ctx)(miniharness/compaction/tool_result_pruner.py,同 上游 compaction-tool-result-pruner)。它只是把 ToolResultPruner 注册为 ctx.toolResultPruner 服务;真正触发它的地方是压缩引擎(miniharness/compaction/ engine.py:139-167)——无论压力触发(step 边界)还是 context-overflow 触发,引擎都会在 做 token 摘要压缩之前先调用 prune.prune_session,按预算(默认 {thresholdChars 8192, headChars 4096, tailChars 1024},按 Unicode 码点计价)对当前 surface 上每个超预算的 tool/result 节点做"保留 head + 固定标记 + 保留 tail"的原地替换 (PRUNE_MARKER = "[... tool result middle pruned ...]"),重新测量后若已降到阈值以下 就无需再走模型摘要。每次替换紧邻一个 log-only 的 compaction/prune 影子计价事件 (经 tokenMeter 计价被遮蔽节点,供纯消费者无需逐节点状态即可扣减)+ 一个带 replace surfaceOp + sourceEventSeqs 的 tool/result 替换事件。装配顺序在 apply_retry_planner(ctx) → install_compaction(ctx) 之后(幂等),demo 装配链已覆盖。

验证:python -m unittest tests.test_retry -v(策略解析、退避边界、 Retry-After 解析、全部 recover 分支、lifetime 信号与竞速等待、插件 teardown 排干/陈旧守卫、loop 集成——重试成功/耗尽终局/非白名单终局); 压缩/溢出见 tests/test_compaction.py。

延伸:图像输入请求(DeepSeek Files API 执行簇)与图像卸载投影

文本回合走 §4.3 的 adapter;带图的回合多一层图像输入请求路由。真实上游 v0.1.6-alpha.1 A 组 #1-3 把"含图消息如何进模型"做成了完整文件管线,mini 对应 miniharness/llm/deepseek_files/ 十二模块(对照 docs/architecture.md 映射表)。

路由目标(ImageRequestTarget,request-image v6):先算"这张图以什么规格请求" ——按像素预算决策(detail 网格 + 4096 单边封顶),与文本 token 一起交给定价器 (request-pricing.py)得出总成本;这是"请求前先知道要付多少",不是请求后的记账。

文件管线(files_api.py + upload_index.py + file_store.py):

  1. file-id 优先:同一会话内已上传过的图片(durable files-v3.json 索引,内容寻址),直接复用 file-id,不重复上传。
  2. 新增则上传:单飞(并发闸)共享同一 upload 任务;配额用尽按 LRU 恢复;索引经 filelock(30s)+ os.replace 原子发布。
  3. base64 回退:file-id 不可用(如已知解码器限制)时退回消息内 base64,且回退决定在请求前(resolve_image_attachment_access),不在中途。
  4. stale-id 恰一次重试:请求期 file-id 已失效(stale)→ 有界重试换新 id,不无限重试。

载体差异:上游走官方 Files API(HTTP 语义有一套独立映射),mini 用 httpx 实现同一契约;上传状态索引上游无等价物,属 mini 生产就绪增量(已登记)。

配套的上下文压力出口——图像卸载投影(miniharness/core/session/projections.py + compaction/image_offload.py,对应上游 compaction-image-offload):当上下文压力过大需要"腾地方",逐字压缩不是唯一手段——可以把最旧的输入图片从历史里卸载,token 立刻省一大块。机制:

  • image/offload 是 durable 事件:记录要卸载的图片 occurrence(当前 surface 的 user/message 或 tool/result 节点 + 深度优先序号),offload_oldest_images 逐个挑选,选定集严格递增 + 拆分校验 fail-closed。
  • 读侧 derive_messages 经 fold_projections 把被选中 occurrence 投影为不可变 offloaded 副本(身份保留,只是内容变占位文本)——历史结点不删,模型侧只见占位。
  • 卸载请求若因上下文仍超限失败,在 agent/request-error 上带 IMAGE_OFFLOAD_REQUIRED 标记触发重新卸载尝试:不消耗重试预算、不落 retry 事件(与上一小节重试的分野)。

4.10 收尾

回合跑通的那一刻,前三章的积木全部就位:日志在写、插件在拦、工具在跑、模型在转。这一章最后要记住的是 turn/step 的分层——turn 是对话的括号,step 是括号里的每一轮"请求 + 工具"。下一章处理一个没解决的实际问题:这些日志怎么写入磁盘、崩溃怎么恢复、整个系统怎么组合启动。