返回课程首页

2

让回答逐字显示

理解 AssistantMessageEventStream 与流式事件。

源码基线:Pi v0.82.0 · commit 518855d

第一章使用 completeSimple() 等待完整回答。这种方式适合脚本,却不适合桌面聊天界面:模型可能几秒甚至更久才结束,用户会怀疑应用是不是卡住了。

这一章只改变消费响应的方式:

completeSimple() → streamSimple()

Provider、Model 和 Context 都不需要更换。

1. 最小流式版本

沿用第一章已经创建好的 modelsmodelcontext

const stream = models.streamSimple(model, context);

for await (const event of stream) {
  if (event.type === "text_delta") {
    process.stdout.write(event.delta);
  }
}

const finalMessage = await stream.result();
context.messages.push(finalMessage);

这段代码同时使用了同一个流的两种能力:

  • for await...of 按到达顺序读取事件;
  • result() 等待并取得最终的 AssistantMessage

它们不是两次模型请求,而是同一个 AssistantMessageEventStream 的两个读取角度。

2. 桌面界面应该保存什么状态

最小桌面聊天界面至少需要区分四种状态:

type ReplyState = {
  status: "idle" | "connecting" | "streaming" | "done" | "error";
  text: string;
  error?: string;
};

然后把 Pi 事件映射到界面状态:

const stream = models.streamSimple(model, context);

setReply({
  status: "connecting",
  text: "",
});

for await (const event of stream) {
  switch (event.type) {
    case "start":
      setReply((current) => ({
        ...current,
        status: "streaming",
      }));
      break;

    case "text_delta":
      setReply((current) => ({
        ...current,
        text: current.text + event.delta,
      }));
      break;

    case "done":
      setReply((current) => ({
        ...current,
        status: "done",
      }));
      break;

    case "error":
      setReply((current) => ({
        ...current,
        status: "error",
        error: event.error.errorMessage,
      }));
      break;
  }
}

这里的 setReply 使用 React 风格写法只是为了表达状态变化,不是 Pi 提供的 API。Pi 提供的是事件,桌面框架负责把事件变成界面状态。

3. Pi 的流里有哪些事件

AssistantMessageEvent 是一个 TypeScript 联合类型。事件可以分成五组:

阶段 事件 用途
整体开始 start 助手消息流开始,后面可以产生增量
普通文本 text_starttext_deltatext_end 创建、追加和结束文本块
Thinking thinking_startthinking_deltathinking_end 创建、追加和结束 Thinking 块
工具调用 toolcall_starttoolcall_deltatoolcall_end 接收工具名称和参数
整体结束 doneerror 成功结束或失败结束

源码位置:packages/ai/src/types.ts 中的 AssistantMessageEvent

统一协议只保证 start 出现在增量事件之前,不保证它表示 HTTP 连接已经建立。如果认证、Provider 定位或 lazy 模块加载失败,lazyStream() 可以直接产生 error,中间没有 start

一段纯文本响应的实际事件记录可能是:

start
text_start     contentIndex=0
text_delta     delta="你"
text_delta     delta="好"
text_end       content="你好"
done           reason="stop"

这段记录用于观察事件形状,不表示所有回答都只有一个内容块。

第二章只处理普通文本,后续章节会继续复用同一套事件协议处理 Thinking 和工具。

4. 一段文字的完整生命周期

一次只包含普通文字的回答,最容易观察到以下顺序:

sequenceDiagram
    participant Provider as Provider API
    participant Stream as AssistantMessageEventStream
    participant UI as 桌面界面

    Provider->>Stream: start
    Stream->>UI: 状态改为 streaming
    Provider->>Stream: text_start(contentIndex = 0)
    Provider->>Stream: text_delta("你")
    UI->>UI: text = "你"
    Provider->>Stream: text_delta("好")
    UI->>UI: text = "你好"
    Provider->>Stream: text_end("你好")
    Provider->>Stream: done(final AssistantMessage)
    Stream->>UI: 状态改为 done

不过,不能把这个简单顺序当成所有响应的固定规则。

仓库文档明确说明:不同内容块的事件不保证连续。文本、Thinking 和 ToolCall 可能交错出现。消费者必须通过 contentIndex 判断一个事件属于 AssistantMessage.content 中的哪个内容块。

因此,只有一个纯文本气泡时,直接累加 event.delta 足够;当界面同时展示 Thinking、文字和工具时,应按 contentIndex 管理多个内容块。

5. contentIndex 为什么重要

Pi 的 AssistantMessage.content 是内容块数组。OpenAI Responses 适配器收到一个新的输出项时,会先创建对应内容块,并把它加入这个数组:

const block: TextContent = {
  type: "text",
  text: "",
};

output.content.push(block);

const slot = {
  type: "text",
  block,
  contentIndex: output.content.length - 1,
};

stream.push({
  type: "text_start",
  contentIndex: slot.contentIndex,
  partial: output,
});

源码位置:packages/ai/src/api/openai-responses-shared.tsprocessResponsesStream()createSlot

后续增量到达时,适配器修改同一个文本块,并携带相同的 contentIndex

slot.block.text += event.delta;

stream.push({
  type: "text_delta",
  contentIndex: slot.contentIndex,
  delta: event.delta,
  partial: output,
});

源码位置:packages/ai/src/api/openai-responses-shared.ts 中的 response.output_text.delta 分支

这说明 contentIndex 不是流事件的序号,而是 AssistantMessage.content 的数组索引。它把增量事件关联到最终内容数组中的位置,通常也决定内容块的显示顺序。

6. 不要把 partial 当成不可变快照

每个增量事件还携带 partial,它表示当前正在构造的 AssistantMessage

在 OpenAI Responses 的实现中,适配器创建一个 output 对象,然后在流式处理期间持续修改其中的 contentusage。事件中的 partial 指向这个正在变化的对象。

因此桌面 UI 有两种稳妥做法:

  1. 对简单文本,直接使用 event.delta 累加自己的界面字符串。
  2. 需要读取完整 partial 时,把需要的字段转换成自己的 UI 数据对象,不依赖 Provider 对象引用变化触发刷新。

例如不要只写:

setMessage(event.partial);

对于依赖引用比较的 UI 框架,同一个对象被原地更新时,这种写法可能无法表达出每次增量变化。本章只有纯文本时,可以只取出 UI 真正需要的值:

const textBlocks = event.partial.content.flatMap(
  (block, contentIndex) =>
    block.type === "text"
      ? [{ contentIndex, text: block.text }]
      : [],
);

setTextBlocks(textBlocks);

只做一层对象展开不等于通用深复制,因为 ToolCall.argumentsusage.cost 等字段还有嵌套对象。这里讨论的是 OpenAI Responses 适配器中可以从源码确认的对象更新方式。应用层最好把 Provider 事件归一化成自己的 UI 数据结构,而不是依赖某个 Provider 的对象实现细节。

7. EventStream 内部怎样保存事件

AssistantMessageEventStream 继承自通用的 EventStream<T, R>

它内部维护两个重要数组:

private queue: T[] = [];
private waiting: ((value: IteratorResult<T>) => void)[] = [];

源码位置:packages/ai/src/utils/event-stream.ts 中的 EventStream

可以把它理解成一个很小的“生产者—消费者”通道:

flowchart LR
    Producer["Provider<br/>push(event)"] --> Decision{"此时有人等待吗?"}
    Decision -->|"有"| Waiting["直接交给 waiting 中的消费者"]
    Decision -->|"没有"| Queue["暂存在 queue"]
    Queue --> Iterator["for await...of"]
    Waiting --> Iterator

push() 的行为是:

  • 如果有正在等待下一条事件的消费者,直接把事件交给它;
  • 否则把事件加入 queue
  • 如果事件是最终事件,同时解析 result() 对应的 Promise。

异步迭代器的行为正好相反:

  • queue 有事件时,立即取出最早的一条;
  • 流已经结束时,退出循环;
  • 两者都不是时,进入 waiting 等待下一次 push()

这就是 for await...of 可以边生成边读取的原因。

8. EventStream 不是广播订阅器

EventStream 内部只有一份共享的 queuewaitingpush() 使用 waiting.shift(),异步迭代器读取事件时也使用 queue.shift()

这意味着它是单消费队列,不是广播:

错误方式:
同一个 stream → 消费者 A
             → 消费者 B

A 和 B 会分走事件,而不是各自收到完整事件副本。

桌面应用应该只启动一个 for await...of 消费循环,再把归约后的应用状态分发给多个界面组件。这里的“归约”是指:把“旧状态 + 新事件”计算成“新状态”。

9. result() 是怎样得到最终消息的

AssistantMessageEventStream 在构造时告诉父类:

event.type === "done" || event.type === "error"

就是完成条件。

如果收到 done,最终结果取 event.message;如果收到 error,最终结果取 event.error

源码位置:packages/ai/src/utils/event-stream.ts 中的 AssistantMessageEventStream

所以:

const finalMessage = await stream.result();

无论成功还是失败,类型上都得到一个 AssistantMessage。失败信息通过它的 stopReasonerrorMessage 表达,而不是要求每个 Provider 用不同的异常形状。

失败的最终消息是否追加到下一轮 Context,由上层会话策略决定。仓库文档允许把 aborted 消息加入 Context 继续对话,但应用不应在不了解失败类型时无条件保存所有错误消息。

错误的详细处理留到第 3 章。

10. 为什么调用后立刻就能拿到 Stream

模型调用前通常要做异步工作:

  • 解析认证;
  • 延迟加载 Provider 的 API 模块;
  • 创建网络请求。

models.streamSimple() 本身会立即返回 AssistantMessageEventStream。这是 lazyStream() 的作用。

export function lazyStream(model, setup) {
  const outer = new AssistantMessageEventStream();

  setup()
    .then((inner) => forwardStream(outer, inner))
    .catch((error) => {
      // 把初始化失败转换成 error 事件
    });

  return outer;
}

源码位置:packages/ai/src/api/lazy.ts 中的 lazyStream

它先创建外层流,再在后台进行异步初始化。内层 Provider 流产生事件后,forwardStream() 把事件逐一推到外层流。

flowchart LR
    App["桌面应用"] --> Outer["外层 Stream<br/>立即返回"]
    Setup["认证与模块加载"] --> Inner["Provider 内层 Stream"]
    Inner -->|"forwardStream"| Outer
    Outer --> UI["for await 消费"]

这样桌面应用不需要等待认证和模块加载完成,便可以先进入“正在连接”的状态。

11. OpenAI Provider 怎样接到具体实现

OpenAI Provider 的工厂代码很短:

export function openaiProvider() {
  return createProvider({
    id: "openai",
    name: "OpenAI",
    baseUrl: "https://api.openai.com/v1",
    auth: {
      apiKey: envApiKeyAuth(
        "OpenAI API key",
        ["OPENAI_API_KEY"],
      ),
    },
    models: Object.values(OPENAI_MODELS),
    api: openAIResponsesApi(),
  });
}

源码位置:packages/ai/src/providers/openai.ts

openAIResponsesApi() 又通过 lazyApi() 延迟加载真正的 openai-responses.ts。也就是说,注册 Provider 时不会立即载入完整 SDK 实现;第一次调用会触发动态 import(),后续模块求值由 JavaScript 运行时的模块缓存复用。

从桌面应用看到的完整路径是:

models.streamSimple()
  → 找到 OpenAI Provider
  → 应用认证
  → Provider.streamSimple()
  → 延迟加载 openai-responses.ts
  → 创建 Provider 内部 Stream
  → 把 Provider 事件转发到应用拿到的 Stream

12. 为什么 Pi 使用事件,而不是只返回字符串

如果 API 只返回一个最终字符串,桌面应用会失去四类信息:

  1. 增量:不知道新到达的是哪一段文字。
  2. 异构内容:无法在同一回答中区分 Text、Thinking 和 ToolCall。
  3. 生命周期:不知道内容块何时创建、更新和结束。
  4. 终止状态:无法统一表达成功、长度限制、工具请求、失败和取消。

事件协议把“回答是什么”和“回答怎样产生”同时交给应用。最终 AssistantMessage 仍然保留完整结果,所以应用不需要在流式体验和最终结构化数据之间二选一。

这是本章最重要的设计复盘:

字符串适合表示结果;事件适合表示一个仍在发生的过程。

13. 一个更完整的界面事件归约器

把事件处理从组件中拿出来,可以减少界面代码与 Pi 类型的耦合:

type ReplyState = {
  status: "connecting" | "streaming" | "done" | "error";
  textBlocks: Record<number, string>;
  error?: string;
};

function reduceReply(
  state: ReplyState,
  event: AssistantMessageEvent,
): ReplyState {
  switch (event.type) {
    case "start":
      return { ...state, status: "streaming" };

    case "text_start":
      return {
        ...state,
        textBlocks: {
          ...state.textBlocks,
          [event.contentIndex]: "",
        },
      };

    case "text_delta":
      return {
        ...state,
        textBlocks: {
          ...state.textBlocks,
          [event.contentIndex]:
            (state.textBlocks[event.contentIndex] ?? "") + event.delta,
        },
      };

    case "done":
      return { ...state, status: "done" };

    case "error":
      return {
        ...state,
        status: "error",
        error: event.error.errorMessage,
      };

    default:
      return state;
  }
}

“归约器”是一个根据旧状态和新事件计算新状态的函数。这个归约器是课程中的桌面应用示例,不是 Pi 仓库内置函数。它的输入严格使用 Pi 的事件结构,并通过 contentIndex 管理多个文本块。

14. 本章小结

  • streamSimple() 返回 AssistantMessageEventStream
  • 同一个 Stream 既能异步迭代事件,也能通过 result() 得到最终消息。
  • 普通文字通过 text_starttext_deltatext_end 表达。
  • contentIndex 把增量事件关联到最终内容数组中的块。
  • 不同内容块的事件可能交错,不能假设事件总是连续出现。
  • EventStream 使用 queuewaiting 连接事件生产者与消费者。
  • EventStream 是单消费队列,不会向多个迭代器广播完整副本。
  • lazyStream() 让应用立即得到外层 Stream,再异步完成认证和模块加载。
  • 桌面应用应把 Pi 事件转换为自己的不可变 UI 状态。
  • 事件同时表达增量、异构内容、生命周期和终止状态。

下一章会在同一个事件循环中加入 Thinking、Token、停止原因和错误状态。

15. 自测

  1. for await...ofstream.result() 是否会发起两次请求?
  2. 为什么不能假设两个 text_delta 之间不会出现工具事件?
  3. contentIndex 关联的是哪个数据结构?
  4. 当消费者还没有调用下一次迭代时,新事件保存在哪里?
  5. lazyStream() 为什么要同时维护外层流和内层流?
  6. 为什么不应该让两个 for await...of 同时读取一个 EventStream?

本章源码依据

  • packages/ai/src/types.ts
  • packages/ai/src/utils/event-stream.ts
  • packages/ai/src/api/lazy.ts
  • packages/ai/src/models.ts
  • packages/ai/src/providers/openai.ts
  • packages/ai/src/api/openai-responses.lazy.ts
  • packages/ai/src/api/openai-responses.ts
  • packages/ai/src/api/openai-responses-shared.ts
  • packages/ai/README.md 的 Quick Start 和 Complete Event Reference