第 04 章:流式输出与中断

非流式的 agent 有一个糟糕的体验:模型思考 30 秒,屏幕死寂 30 秒,然后一整段话砸出来。这一章我们切换到 SSE 流式,实现打字机效果,并接上 Ctrl+C 中断——而且中断后已收到的部分内容要保留。流式也是整个教程里最容易写错的一章,慢慢来。

SSE 协议五分钟速成

请求体加 "stream": true 后,响应变成 text/event-stream:一行一个 data: {json},每个 JSON 是一个增量 chunk,最后以 data: [DONE] 收尾:

data: {"choices":[{"index":0,"delta":{"role":"assistant","content":"你"},"finish_reason":null}]}
data: {"choices":[{"index":0,"delta":{"content":"好"},"finish_reason":null}]}
data: {"choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}
data: [DONE]

delta 是增量、finish_reason 是结局(stop 正常结束 / tool_calls 有工单 / length 撞输出上限)。解析就是:按行切 buffer,data: 开头就 JSON.parse,跨 TCP 分片的半行留到下一轮——代码里的双层 while 就是干这个的。

难点:tool_calls 是"流着来"的

文本增量好处理,拼起来就行。工具调用麻烦:它的 arguments 是一个 JSON 字符串,被切成任意长度的分片陆续到达,每个分片只带 index(第几个调用)和一小段字符串。pi 的 pi-ai 在流式过程中还会对半拉的 JSON 做渐进解析,以便 UI 实时预览(比如 diff 一边流一边渲染);mini-pi 只做必要的事——按 index 归并、把字符串拼完整,最后再一次性 JSON.parse:

src/main.ts(streamChat 的累积核心)
for (const tc of delta.tool_calls ?? []) {
  const slot = (acc.toolCalls[tc.index ?? 0] ??= { id: "", name: "", arguments: "" });
  if (tc.id) slot.id += tc.id;
  if (tc.function?.name) slot.name += tc.function.name;
  if (tc.function?.arguments) slot.arguments += tc.function.arguments;
}

中断:AbortController + 保留部分结果

pi 的设计里 abort 是一等公民:中止的请求带着 stopReason: "aborted"已收到的部分内容返回,部分回复照常进入会话历史——下次接着聊。我们照做:AbortController 同时接进 fetchbash 执行,Ctrl+C 触发,流式函数在中止时返回累积了一半的结果。

还有一个 pi 抄来的细节:finish_reason === "length" 意味着输出被 token 上限截断,此时流式拼出来的工具参数可能是不完整的 JSON。pi 的 agent-loop.ts 会把这一整批 toolCall 判失败("参数可能被截断,请重新发起"),让模型重试,而不是执行一个可能残缺的操作。我们也加上。

代码

用流式版替换上一章的 chatWithTools:

src/main.ts(streamChat 完整实现)
type FinishReason = "stop" | "tool_calls" | "length" | "aborted";

interface AssistantResult {
  content: string;
  toolCalls: ToolCall[];
  finishReason: FinishReason;
}

/**
 * 流式调用:解析 SSE,文本增量通过 onText 实时吐出;
 * tool_calls 的增量按 index 累积成完整调用。
 * signal 中止时返回已收到的部分结果,finishReason 为 "aborted"。
 */
async function streamChat(
  messages: Message[],
  signal: AbortSignal,
  onText: (delta: string) => void,
): Promise<AssistantResult> {
  const acc = {
    content: "",
    toolCalls: [] as { id: string; name: string; arguments: string }[],
    finishReason: "stop" as FinishReason,
  };

  const assemble = (a: typeof acc, reason: FinishReason): AssistantResult => ({
    content: a.content,
    toolCalls: a.toolCalls
      .filter(Boolean)
      .map((t, i) => ({
        id: t.id || `call_${i}`,
        type: "function" as const,
        function: { name: t.name, arguments: t.arguments },
      })),
    finishReason: reason,
  });

  let res: Response;
  try {
    res = await fetch(`${CONFIG.baseUrl}/chat/completions`, {
      method: "POST",
      headers: {
        "Content-Type": "application/json",
        Authorization: `Bearer ${CONFIG.apiKey}`,
      },
      body: JSON.stringify({ model: CONFIG.model, messages, tools: TOOL_DEFS, stream: true }),
      signal,
    });
  } catch (err) {
    if (signal.aborted) return assemble(acc, "aborted");
    throw err;
  }
  if (!res.ok) throw new Error(`HTTP ${res.status}: ${await res.text()}`);
  if (!res.body) throw new Error("Response has no body");

  const reader = res.body.getReader();
  const decoder = new TextDecoder();
  let buffer = "";
  try {
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      buffer += decoder.decode(value, { stream: true });
      let nl: number;
      while ((nl = buffer.indexOf("\n")) >= 0) {
        const line = buffer.slice(0, nl).trim();
        buffer = buffer.slice(nl + 1);
        if (!line.startsWith("data:")) continue;
        const data = line.slice(5).trim();
        if (data === "[DONE]") continue;
        let chunk: any;
        try {
          chunk = JSON.parse(data);
        } catch {
          continue; // 容忍不完整的一行,等下一个分片
        }
        const choice = chunk.choices?.[0];
        if (!choice) continue;
        const delta = choice.delta ?? {};
        if (typeof delta.content === "string" && delta.content) {
          acc.content += delta.content;
          onText(delta.content);
        }
        for (const tc of delta.tool_calls ?? []) {
          const slot = (acc.toolCalls[tc.index ?? 0] ??= { id: "", name: "", arguments: "" });
          if (tc.id) slot.id += tc.id;
          if (tc.function?.name) slot.name += tc.function.name;
          if (tc.function?.arguments) slot.arguments += tc.function.arguments;
        }
        if (choice.finish_reason) acc.finishReason = choice.finish_reason;
      }
    }
  } catch (err) {
    if (signal.aborted) return assemble(acc, "aborted");
    throw err;
  }
  return assemble(acc, signal.aborted ? "aborted" : acc.finishReason);
}

runAgent 换成流式调用,并处理两种特殊结局:

src/main.ts(runAgent 更新)
async function runAgent(messages: Message[], signal: AbortSignal): Promise<void> {
  while (true) {
    let printed = false;
    const result = await streamChat(messages, signal, (delta) => {
      printed = true;
      process.stdout.write(delta);
    });
    if (printed) process.stdout.write("\n");

    const assistantMsg: Message = { role: "assistant", content: result.content || null };
    if (result.toolCalls.length > 0) assistantMsg.tool_calls = result.toolCalls;
    messages.push(assistantMsg);

    if (result.finishReason === "aborted") {
      console.log("\x1b[2m[aborted — partial response kept]\x1b[0m");
      return;
    }
    if (result.toolCalls.length === 0) return;

    for (const call of result.toolCalls) {
      let toolResult: { content: string; isError: boolean };
      if (result.finishReason === "length") {
        // 输出撞上 token 上限,流式拼出来的参数可能被截断:
        // 整批不执行,报错让模型重发(pi 同款处理)
        toolResult = {
          content: `Tool call "${call.function.name}" was not executed: the response hit the output token limit, so its arguments may be truncated. Re-issue the tool call with complete arguments.`,
          isError: true,
        };
      } else {
        console.log(`\x1b[36m[${call.function.name}]\x1b[0m \x1b[2m${call.function.arguments}\x1b[0m`);
        toolResult = await executeTool(call.function.name, call.function.arguments, signal);
        console.log(`  \x1b[2m${toolResult.content.split("\n").slice(0, 5).join("\n  ")}\x1b[0m`);
      }
      messages.push({ role: "tool", tool_call_id: call.id, content: toolResult.content });
    }
  }
}

executeTooltoolBash 各自加一个可选的 signal 参数,并把它传给 execAsync({ ..., signal })——这样 Ctrl+C 连正在跑的 shell 命令一起杀。

最后是 REPL 的中断接线:agent 工作时 Ctrl+C 是"中止当前运行",空闲时才是"退出":

src/main.ts(REPL 更新)
  let running = false;
  let aborter: AbortController | null = null;
  let closed = false;
  rl.on("close", () => {
    closed = true;
  });

  rl.on("line", (line) => {
    const text = line.trim();
    if (!text || running) {
      rl.prompt();
      return;
    }
    running = true;
    aborter = new AbortController();
    messages.push({ role: "user", content: text });
    runAgent(messages, aborter.signal)
      .catch((err) => console.error(`\n\x1b[31m[error]\x1b[0m ${err.message}`))
      .finally(() => {
        running = false;
        aborter = null;
        if (!closed) rl.prompt();
      });
  });

  rl.on("SIGINT", () => {
    if (running && aborter) {
      aborter.abort(); // 中止当前运行,部分结果保留
      return;
    }
    rl.close(); // 空闲时退出
  });

验收

  • 问一个长问题("把《静夜思》逐字解释一遍"),输出应该逐字流出;中途 Ctrl+C,看到 [aborted — partial response kept],且部分回复已经进入对话历史(下一轮模型能接着看到)。
  • 让它跑一个 sleep 30 的 bash,Ctrl+C 应该立刻把命令也杀掉。

本章要点

  • SSE = 按行切的增量 JSON;delta.content 拼文本,tool_callsindex 归并拼参数字符串。
  • abort 要贯穿整条链(fetch、工具执行),部分结果保留进上下文——这是 pi 的明确设计。
  • length 结局的工具调用一律不执行,报错让模型重发。
  • pi 在流式上还走得更远:半拉 JSON 渐进解析做实时 UI 预览、统一事件流(text_delta/toolcall_delta……),见设计思想:LLM 抽象层

下一章:第 05 章:系统提示词与上下文文件