# 流式响应与前端实现

## 流式响应究竟改变了什么

传统请求等完整响应后解析一次。流式请求在响应体尚未结束时不断接收数据，因此前端更早显示结果，也必须处理拆包、解码、部分完成、取消和中途错误。本章目标是 L3：理解网络字节与业务事件之间的区别，并能构建可测试的消费端。

前置是 fetch、异步迭代和上一章状态模型。一个网络 chunk 不等于一个 token、不等于一行 JSON，也不一定包含完整的中文字符。同样一段响应可以被任意拆分或合并；程序必须按协议边界解析，而不是按每次 `reader.read()` 返回的边界解析。

## SSE、NDJSON 与 WebSocket 怎么选

SSE 使用 `text/event-stream`，事件通常以空行分隔，可包含 event、data、id 等字段。浏览器 EventSource 适合订阅式 GET 场景，但不能像 fetch 一样随意设置请求体和自定义请求头。需要 POST 请求体时，可以使用 fetch 读取 SSE 并使用正确解析器。

NDJSON 每行是一个 JSON 对象，适合自己控制前后端的教学与内部接口。WebSocket 提供双向消息通道，适合持续双向交互，但增加连接管理、重连、认证更新与背压方面的设计。仅仅逐步展示一个回答，通常不必先引入 WebSocket。

本章选择 NDJSON，明确它不是 SSE，也不是任何模型供应商的原始协议。真实项目可以在后端把供应商事件转换为稳定的应用事件，前端就不会与供应商事件格式紧密绑定。

## 定义前后端事件契约

```json
{"type":"delta","text":"你好"}
{"type":"delta","text":"，这是一段示例。"}
{"type":"done"}
```

上面三行分别是独立事件，并不组成一个合法 JSON 数组。成功必须出现 done；网络流结束但没有 done 视为不完整。错误可以在响应头发出前用 HTTP 错误码返回，也可以在响应开始后发送约定的 error 事件。头部发出之后不能再把状态码改成 500。

## 完整的离线流解析实验

保存为 `read-events.mjs`，运行 `node read-events.mjs`。使用 Node 22+ 内置 ReadableStream、Response 和 TextDecoder，不依赖模型或网络。模拟响应会故意把每个字节单独发送，从而检验中文字符跨 chunk 的情况。

```js
import assert from 'node:assert/strict';

async function consume(response, onText) {
  if (!response.ok) throw new Error(`HTTP_${response.status}`);
  if (!response.body) throw new Error('EMPTY_BODY');
  const reader = response.body.getReader();
  const decoder = new TextDecoder('utf-8', { fatal: true });
  let buffer = '';
  let completed = false;

  function parseLine(line) {
    if (line.length > 1_000_000) throw new Error('EVENT_TOO_LARGE');
    if (!line.trim()) return;
    const event = JSON.parse(line);
    if (!event || typeof event !== 'object') throw new Error('INVALID_EVENT');
    if (completed) throw new Error('EVENT_AFTER_DONE');
    if (event.type === 'delta' && typeof event.text === 'string') {
      onText(event.text);
    } else if (event.type === 'done') {
      completed = true;
    } else if (event.type === 'error') {
      throw new Error('UPSTREAM_STREAM_FAILED');
    } else {
      throw new Error('UNKNOWN_EVENT');
    }
  }

  try {
    while (true) {
      const { value, done } = await reader.read();
      if (done) break;
      // stream:true 保留尚未组成完整字符的字节。
      buffer += decoder.decode(value, { stream: true });
      let newline;
      while ((newline = buffer.indexOf('\n')) >= 0) {
        const line = buffer.slice(0, newline).replace(/\r$/, '');
        buffer = buffer.slice(newline + 1);
        parseLine(line);
      }
      // 教学上限按 JS 字符串长度计算；生产可另设字节与总响应上限。
      if (buffer.length > 1_000_000) throw new Error('EVENT_TOO_LARGE');
    }
    buffer += decoder.decode(); // 冲刷解码器剩余字节。
    if (buffer.trim()) parseLine(buffer);
    if (!completed) throw new Error('INCOMPLETE_STREAM');
  } finally {
    // 异常时主动取消消费；正常结束时 cancel 也不会重新启动流。
    await reader.cancel().catch(() => {});
    reader.releaseLock();
  }
}

const wire = [
  { type: 'delta', text: '你好' },
  { type: 'delta', text: '，前端工程师。' },
  { type: 'done' },
].map(event => JSON.stringify(event)).join('\n') + '\n';
const bytes = new TextEncoder().encode(wire);
let position = 0;
const stream = new ReadableStream({
  pull(controller) {
    if (position === bytes.length) return controller.close();
    controller.enqueue(bytes.slice(position, ++position));
  },
});
let answer = '';
await consume(new Response(stream), text => { answer += text; });
assert.equal(answer, '你好，前端工程师。');
await assert.rejects(
  () => consume(new Response('{"type":"delta","text":"未完成"}\n'), () => {}),
  /INCOMPLETE_STREAM/,
);
console.log('通过：跨字节中文解码，以及缺失 done 检测');
```

`getReader()` 获得读取器并锁定流；`read()` 返回 Promise，结果包含 value 和 done；TextDecoder 负责字节到字符串；buffer 负责字符串到完整协议行。两种缓冲不能互相替代。`fatal:true` 让无效 UTF-8 明确失败，不静默把坏字节替换成乱码。

解析器要求成功终止事件，是因为 TCP 连接关闭并不能表达业务是否完整。这里等待网络结束后检查 completed，真实生产协议还需要空闲超时或在 done 时结束消费，防止供应商发完结果却迟迟不关闭连接。总输出大小也应限制，不能只有单行缓冲上限。

## 接到 Vue 与真实请求时怎么改

把 consume 放进独立模块并导出。组合式函数创建 AbortController，用 `fetch('/api/chat', { method:'POST', body:JSON.stringify(input), headers:{'Content-Type':'application/json'}, signal:controller.signal })` 获得 Response，再交给 consume。fetch 地址是你自己后端的契约，不是模型供应商地址。

每次请求保存唯一 attemptId，onText 更新前检查它仍是当前尝试。组件卸载时 abort，catch 中根据 signal.aborted 区分用户取消与真正错误，finally 只清理属于自己的控制器。不要在旧请求 finally 中无条件把新请求的 loading 清掉。

这一段是集成步骤说明，并非完整 Vue 组件。为了实际复用，应根据上一章的状态模型封装 useChat，再编写组件和请求模块测试。真实接口还需要登录、权限、超时、后端取消传播和安全日志，不能把流解析函数当作完整聊天服务。

## 逐 token 更新的性能问题

网络可以高频收到小片段，但 UI 不必每片段立即重新解析整个 Markdown。可以先把增量放进缓冲，在一帧或固定短时间窗口内合并更新。完成后再执行较重的语法高亮；很长的历史消息可以虚拟化，而正在生成的消息单独管理。

先测量再优化：记录首字时间、整段完成时间、更新频率和长任务。若后端按时发送但浏览器成批收到，检查反向代理缓冲、压缩和中间层。不要看到“最后一下全部出现”就断定模型没有流式返回。

## 把协议放回真实的 Vue 页面

上面的单文件实验只回答“字节能否正确拼成事件”。接下来补齐一个可以启动的前后端项目，观察 HTTP 连接、取消、业务完成和组件更新如何串联。[下载完整实验项目](downloads/streaming-lab.zip)，解压后在项目目录执行命令。下载包包含锁文件和所有源码，首次安装依赖需要网络，运行实验不需要模型密钥。

这个实验故意使用固定回答。这样，当你看到中文乱码、重复追加或取消失效时，可以排除模型输出变化。先把传输链路做对，再把服务端固定生成器替换成模型适配器；浏览器消费的事件契约保持稳定。示例只监听本机地址，没有登录和持久化，不应直接用来提供公共服务。

### 目录和启动顺序

```text
stream-lab/
  package.json         项目依赖和运行命令
  package-lock.json    锁定直接与间接依赖
  index.html           浏览器入口
  vite.config.js       Vue 编译与本地 API 代理
  server.mjs           Node HTTP 服务与模拟流
  src/main.js          创建 Vue 应用
  src/readEvents.js    把 Response 转换为文本增量
  src/App.vue          输入、状态、取消与展示
  test.mjs             不打开浏览器也能执行的 HTTP 验证
```

先在目录中运行 `npm ci`。终端一运行 `npm run dev:server`，看到端口 8787 的提示；终端二运行 `npm run dev:web`，打开 `http://127.0.0.1:5173`。Windows 如果不能运行 npm 的 PowerShell 脚本，可以使用 npm.cmd。两个终端分别维护两个进程，关闭其中一个会导致对应能力不可用。不要把终端正在监听误判为程序“卡住”。

前端发往 `/api/chat` 的请求先到 Vite，再由开发代理转交 Node。此时浏览器看到的是同源请求，省去教学阶段无关的跨域配置。这个代理只存在于开发服务器；执行 build 得到的静态文件不会自带 Node 后端和代理。上线时需要单独决定静态资源和 API 的路由，本地学习不需要执行部署。

### 请求与响应契约

| 输入或事件 | 约束 | 使用它的原因 |
| --- | --- | --- |
| question | 非空字符串，最多 1000 个 JS 字符串单位 | 限制输入，避免无限生成或内存占用；与码点计数不同 |
| mode | normal、error、incomplete，默认 normal | 在同一界面重现三条稳定的执行路径 |
| HTTP 400 | JSON 或问题字段无效 | 生成开始前拒绝无效请求 |
| HTTP 413 | 请求体超过 16384 字节 | 在读取请求时限制内存积累 |
| delta.text | 字符串，可为空 | 增量追加，不表示整个回答完成 |
| done | 只允许一次，位于结尾 | 显式证明应用完成，而不是仅证明连接关闭 |
| error | 流内失败事件 | HTTP 头已发送后报告业务错误 |

为什么需要两种失败通道？服务端在开始输出前仍可以返回 400；一旦已经发出 200 和部分正文，状态码就不能改成 500。此时必须使用流内 error 或断开连接，客户端把它转为失败状态。如果你的日志只统计 HTTP 状态码，就会把一部分失败的生成计入成功，这正是需要业务完成指标的原因。

### 项目配置与页面入口

以下文件全部提供完整内容。package.json 固定本次实验使用的依赖版本，锁文件包含在下载包中。你手工抄写文件时，第一次使用 npm install 生成自己的锁文件；之后使用 npm ci 复现。升级库之前先读变更说明并重新运行实验，不要随手把全部依赖替换为 latest。


#### package.json

```json package.json
{
  "name": "streaming-learning-lab",
  "private": true,
  "type": "module",
  "scripts": {
    "dev:server": "node server.mjs",
    "dev:web": "vite",
    "build": "vite build",
    "test": "node --test test.mjs"
  },
  "dependencies": { "vue": "3.5.42" },
  "devDependencies": { "vite": "8.2.2", "@vitejs/plugin-vue": "6.0.8" }
}

```


#### index.html

```html index.html
<!doctype html>
<html lang="zh-CN">
<head><meta charset="UTF-8"><meta name="viewport" content="width=device-width,initial-scale=1"><title>流式交互实验</title></head>
<body><div id="app"></div><script type="module" src="/src/main.js"></script></body>
</html>

```


#### vite.config.js

```js vite.config.js
import { defineConfig } from 'vite';
import vue from '@vitejs/plugin-vue';
export default defineConfig({
  plugins: [vue()],
  server: {
    host: '127.0.0.1', port: 5173, strictPort: true,
    // 浏览器只请求当前前端站点；开发代理转发给本地后端。
    proxy: { '/api': 'http://127.0.0.1:8787' },
  },
});

```


#### src/main.js

```js src/main.js
import { createApp } from 'vue';
import App from './App.vue';
createApp(App).mount('#app');

```


`createApp(App).mount('#app')` 把根组件挂载到 HTML 中的同名节点。vite.config.js 的 proxy 只负责转发请求，不能代替服务端的参数校验。strictPort 让被占用的端口直接报错，避免工具自动改端口后你仍然访问旧页面。这几项配置都服务于可复现运行，与模型能力没有关系。

### 服务端：先校验，再逐步写入


#### server.mjs

```js server.mjs
import { createServer } from 'node:http';
import { once } from 'node:events';
import { setTimeout as delay } from 'node:timers/promises';
import { pathToFileURL } from 'node:url';

function readJson(request, limit = 16_384) {
  return new Promise((resolve, reject) => {
    let size = 0;
    let rejected = false;
    const chunks = [];
    request.on('data', chunk => {
      size += chunk.length;
      if (size > limit) {
        chunks.length = 0;
        if (!rejected) reject(Object.assign(new Error('INPUT_TOO_LARGE'), { status: 413 }));
        rejected = true;
        return; // 继续消费请求，但不再把内容保存在内存。
      }
      if (!rejected) chunks.push(chunk);
    });
    request.on('end', () => {
      if (rejected) return;
      try { resolve(JSON.parse(Buffer.concat(chunks).toString('utf8'))); }
      catch { reject(Object.assign(new Error('INVALID_JSON'), { status: 400 })); }
    });
    request.on('error', reject);
  });
}

export function createDemoServer({ intervalMs = 90 } = {}) {
  return createServer(async (request, response) => {
    if (request.method !== 'POST' || request.url !== '/api/chat') {
      response.writeHead(404, { 'Content-Type': 'application/json' });
      return response.end(JSON.stringify({ error: 'NOT_FOUND' }));
    }
    const controller = new AbortController();
    // 响应连接关闭时，停止上游生成与等待；不要把请求体读取完当作取消。
    response.on('close', () => controller.abort());
    try {
      const input = await readJson(request);
      if (!input || typeof input.question !== 'string' || !input.question.trim()
          || input.question.length > 1000) {
        throw Object.assign(new Error('INVALID_QUESTION'), { status: 400 });
      }
      const mode = input.mode ?? 'normal';
      if (!['normal', 'error', 'incomplete'].includes(mode)) {
        throw Object.assign(new Error('INVALID_MODE'), { status: 400 });
      }
      response.writeHead(200, {
        'Content-Type': 'application/x-ndjson; charset=utf-8',
        'Cache-Control': 'no-store',
        'X-Content-Type-Options': 'nosniff',
      });
      response.flushHeaders();
      async function send(event) {
        controller.signal.throwIfAborted();
        // write=false 表示应等消费者消化缓冲，不能继续无限追加。
        if (!response.write(JSON.stringify(event) + '\n')) {
          await once(response, 'drain', { signal: controller.signal });
        }
      }
      const answer = `这是离线模拟回答。你问的是：${input.question.trim()}。本实验只验证流式交互。`;
      let count = 0;
      for (const char of answer) {
        await delay(intervalMs, undefined, { signal: controller.signal });
        if (++count === 9 && mode === 'error') {
          await send({ type: 'error', code: 'DEMO_FAILURE' });
          return response.end();
        }
        if (count === 9 && mode === 'incomplete') return response.end();
        await send({ type: 'delta', text: char });
      }
      await send({ type: 'done' });
      response.end();
    } catch (error) {
      if (controller.signal.aborted) return;
      if (!response.headersSent) {
        response.writeHead(error.status ?? 500, { 'Content-Type': 'application/json' });
        response.end(JSON.stringify({ error: error.status ? error.message : 'SERVER_ERROR' }));
      } else {
        response.end(JSON.stringify({ type: 'error', code: 'STREAM_FAILED' }) + '\n');
      }
    }
  });
}

if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) {
  createDemoServer().listen(8787, '127.0.0.1', () => {
    console.log('离线模拟后端：http://127.0.0.1:8787');
  });
}

```


`readJson(request, limit)` 的返回值是 Promise，成功得到解析后的普通对象，失败得到带状态码的 Error。读取大小按照 Buffer 的字节数累计，超过上限立即拒绝，并清空已保存的片段。后续到达的数据继续被消费但不再储存。这是教学实现；面对真实公网慢速请求，还要配置网关、请求时限和连接限制。

`createDemoServer({ intervalMs })` 返回未监听端口的 Server。把创建和监听拆开，测试就能用端口 0 让操作系统选择空闲端口，避免“机器上已经有服务占用 8787”导致测试失败。生产入口与测试使用的是同一个 handler，因此成功、失败和参数约束不会各实现一遍。

send 把一个事件序列化成一行，并检查 write 的返回值。false 表示 Node 的可写缓冲已经达到背压条件，应等待 drain；它不表示这次数据已经丢失，更不能重新发送同一行。若忽略返回值，在很慢的消费者下持续生成大量内容，进程会积累内存。背压是生产者配合消费者速度的机制，不是让网络瞬间变快。

这里监听 response 的 close：客户端关闭响应连接时，让控制器停止等待与后续写入。不要把“请求体已经读完”当作“用户取消回答”，POST 请求体可能早已发送结束，但响应还要生成几十秒。真实模型适配器还应接收相同的 signal；只取消本地延时却没有取消上游调用，费用和资源占用可能继续发生。

### 客户端：解析器是独立模块


#### src/readEvents.js

```js src/readEvents.js

export async function consume(response, onText) {
  if (!response.ok) throw new Error(`HTTP_${response.status}`);
  if (!response.body) throw new Error('EMPTY_BODY');
  const reader = response.body.getReader();
  const decoder = new TextDecoder('utf-8', { fatal: true });
  let buffer = '';
  let completed = false;

  function parseLine(line) {
    if (line.length > 1_000_000) throw new Error('EVENT_TOO_LARGE');
    if (!line.trim()) return;
    const event = JSON.parse(line);
    if (!event || typeof event !== 'object') throw new Error('INVALID_EVENT');
    if (completed) throw new Error('EVENT_AFTER_DONE');
    if (event.type === 'delta' && typeof event.text === 'string') {
      onText(event.text);
    } else if (event.type === 'done') {
      completed = true;
    } else if (event.type === 'error') {
      throw new Error('UPSTREAM_STREAM_FAILED');
    } else {
      throw new Error('UNKNOWN_EVENT');
    }
  }

  try {
    while (true) {
      const { value, done } = await reader.read();
      if (done) break;
      // stream:true 保留尚未组成完整字符的字节。
      buffer += decoder.decode(value, { stream: true });
      let newline;
      while ((newline = buffer.indexOf('\n')) >= 0) {
        const line = buffer.slice(0, newline).replace(/\r$/, '');
        buffer = buffer.slice(newline + 1);
        parseLine(line);
      }
      // 教学上限按 JS 字符串长度计算；生产可另设字节与总响应上限。
      if (buffer.length > 1_000_000) throw new Error('EVENT_TOO_LARGE');
    }
    buffer += decoder.decode(); // 冲刷解码器剩余字节。
    if (buffer.trim()) parseLine(buffer);
    if (!completed) throw new Error('INCOMPLETE_STREAM');
  } finally {
    // 异常时主动取消消费；正常结束时 cancel 也不会重新启动流。
    await reader.cancel().catch(() => {});
    reader.releaseLock();
  }
}

```


解析器的输入是 Response 和同步回调 onText；它直到读完整个响应并确认 done 后才兑现 Promise。每个 delta 调用一次回调；任何格式错误、未知事件或缺少完成标记都拒绝 Promise。回调不要偷偷返回一个未等待的异步写入任务，否则 consume 完成不代表那些任务也完成。需要异步消费时，应把回调契约改为可等待并在解析过程中顺序 await。

解析分三层：TextDecoder 保留跨 chunk 的 UTF-8 字节；buffer 保留跨读取的半行；parseLine 根据应用协议判断一行的类型。不要把三层合成 `JSON.parse(decoder.decode(value))`，它碰巧在某次网络分块下成功，也会在代理、带宽和浏览器变化后失败。

本实现同时限制完整单行和未完成的缓冲长度，避免只限制残留 buffer 却允许巨大的完整事件绕过约束。它没有限制整个回答的累计长度，页面演示因服务端固定短回答而有界；接入真实模型前应增加总输出、持续时间和空闲时间限制。单行上限、总量上限和超时限制解决的是三个不同的问题。

### Vue 组件：谁拥有当前请求


#### src/App.vue

```vue src/App.vue
<script setup>
import { computed, onUnmounted, ref } from 'vue';
import { consume } from './readEvents.js';

const question = ref('流式响应和普通响应有什么区别？');
const mode = ref('normal');
const answer = ref('');
const phase = ref('idle');
const error = ref('');
let attemptId = 0;
let activeController = null;
const busy = computed(() => ['connecting', 'streaming'].includes(phase.value));
const labels = {
  idle: '等待输入', connecting: '正在连接', streaming: '正在生成',
  completed: '回答完成', cancelled: '已停止，当前内容可能不完整', failed: '生成失败',
};

async function submit() {
  if (!question.value.trim()) return;
  activeController?.abort();
  const mine = new AbortController();
  const myAttempt = ++attemptId;
  activeController = mine;
  answer.value = '';
  error.value = '';
  phase.value = 'connecting';
  try {
    const response = await fetch('/api/chat', {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify({ question: question.value, mode: mode.value }),
      signal: mine.signal,
    });
    await consume(response, text => {
      if (myAttempt !== attemptId || mine.signal.aborted) return; // 忽略旧请求与取消后的增量。
      phase.value = 'streaming';
      answer.value += text;
    });
    mine.signal.throwIfAborted();
    if (myAttempt === attemptId) phase.value = 'completed';
  } catch (cause) {
    if (myAttempt !== attemptId) return;
    phase.value = mine.signal.aborted ? 'cancelled' : 'failed';
    if (!mine.signal.aborted) error.value = cause.message;
  } finally {
    // 旧请求清理时不能移除新请求正在使用的控制器。
    if (activeController === mine) activeController = null;
  }
}

function stop() { activeController?.abort(); }
function onKeydown(event) {
  // 中文输入法确认候选词时，不误触发发送。
  if (event.key === 'Enter' && !event.shiftKey && !event.isComposing) {
    event.preventDefault();
    if (!busy.value) submit();
  }
}
onUnmounted(() => { attemptId++; activeController?.abort(); });
</script>

<template>
  <main>
    <h1>流式交互实验</h1>
    <p>本地固定回答，无模型调用或费用。</p>
    <label for="question">问题</label>
    <textarea id="question" v-model="question" rows="4" @keydown="onKeydown"></textarea>
    <label for="mode">实验场景</label>
    <select id="mode" v-model="mode" :disabled="busy">
      <option value="normal">正常完成</option>
      <option value="error">服务端中途报错</option>
      <option value="incomplete">缺少完成事件</option>
    </select>
    <div class="actions">
      <button :disabled="busy || !question.trim()" @click="submit">发送</button>
      <button :disabled="!busy" @click="stop">停止生成</button>
    </div>
    <p role="status" aria-live="polite">{{ labels[phase] }}</p>
    <p v-if="error" role="alert">{{ error }}</p>
    <!-- 使用插值按纯文本显示，模型输出不会被当成 HTML 执行。 -->
    <pre aria-label="回答内容">{{ answer || '回答将在这里逐步显示。' }}</pre>
  </main>
</template>

<style scoped>
main { max-width: 720px; margin: 40px auto; padding: 20px; font: 16px/1.8 system-ui; }
label { display: block; margin-top: 16px; }
textarea, select { width: 100%; font: inherit; box-sizing: border-box; padding: 8px; }
.actions { display: flex; gap: 12px; margin-top: 16px; }
button { font: inherit; padding: 6px 16px; }
pre { white-space: pre-wrap; overflow-wrap: anywhere; font: inherit; background: #f3f5f4; color: #213b2e; padding: 20px; }
</style>

```


注意这里存在三种不同的值。响应式 phase 决定用户看见的状态；attemptId 决定某个回调还有没有更新界面的资格；AbortController 决定底层请求是否应停止。后两者不必用于模板，所以可以是普通局部变量。把控制器存成全局单例，会让多个聊天窗口互相取消，组合式函数实例或组件实例应拥有自己的请求资源。

submit 捕获 myAttempt 和 mine。这是该次异步过程的固定身份。每个增量先比较身份，再检查取消状态；finally 也只有在仍拥有控制器时才能清理。设想请求 A 的 finally 晚于请求 B 开始执行，如果直接把 activeController 设为 null，B 的停止按钮就失效。这类错误不靠增加 loading 判断解决，而靠明确资源所有权解决。

消费完成后还检查取消信号，避免用户刚点停止、done 又已经在本地缓冲时，界面反而跳成“完成”。组件卸载时递增 attemptId 并中止请求，既释放资源，也使已经排队的回调失效。保留部分文本是界面选择，completed、cancelled、failed 的标识才说明文本是否完整，不能仅凭“气泡里有内容”判断成功。

### 实际验证与故障实验


#### test.mjs

```js test.mjs
import test from 'node:test';
import assert from 'node:assert/strict';
import { once } from 'node:events';
import { createDemoServer } from './server.mjs';
import { consume } from './src/readEvents.js';

async function withServer(run) {
  const server = createDemoServer({ intervalMs: 0 });
  server.listen(0, '127.0.0.1');
  await once(server, 'listening');
  const url = `http://127.0.0.1:${server.address().port}/api/chat`;
  try { await run(url); }
  finally { server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); }
}
const post = (url, body, signal) => fetch(url, {
  method: 'POST', headers: { 'Content-Type': 'application/json' },
  body: JSON.stringify(body), signal,
});

test('真实HTTP连接逐步返回完整文本', async () => withServer(async url => {
  let text = '';
  await consume(await post(url, { question: '你好' }), delta => { text += delta; });
  assert.ok(text.includes('你问的是：你好'));
}));
test('错误事件不能被当成成功', async () => withServer(async url => {
  await assert.rejects(() => post(url, { question: '你好', mode: 'error' })
    .then(response => consume(response, () => {})), /UPSTREAM_STREAM_FAILED/);
}));
test('没有done的连接关闭必须报告不完整', async () => withServer(async url => {
  await assert.rejects(() => post(url, { question: '你好', mode: 'incomplete' })
    .then(response => consume(response, () => {})), /INCOMPLETE_STREAM/);
}));
test('空问题与超大请求被拒绝', async () => withServer(async url => {
  assert.equal((await post(url, { question: ' ' })).status, 400);
  assert.equal((await post(url, { question: 'x'.repeat(20_000) })).status, 413);
}));

test('收到首个增量后取消，请求以 AbortError 结束', async () => withServer(async url => {
  const controller = new AbortController();
  let received = 0;
  await assert.rejects(async () => {
    const response = await post(url, { question: '取消实验' }, controller.signal);
    await consume(response, () => { received++; controller.abort(); });
  }, error => error.name === 'AbortError');
  assert.ok(received >= 1);
}));

test('完整超长事件也受限制，不能靠换行绕过上限', async () => {
  const response = new Response(JSON.stringify({ type: 'delta', text: 'a'.repeat(1_000_001) }) + '\n');
  await assert.rejects(() => consume(response, () => {}), /EVENT_TOO_LARGE/);
});

```


本教材编写时在 Node 22.22.0 实际执行了 npm test，六项检查通过，并执行 npm run build 通过 Vue 编译。测试使用真实本机 HTTP 连接验证正常完成、服务端报错、不完整结束、输入约束与取消，另检查完整超长事件限制。这个结论不包含浏览器键盘、屏幕阅读器或真实模型供应商的行为；这些需要分别验收。

| 操作 | 应观察到什么 | 发生偏差时先查哪里 |
| --- | --- | --- |
| 正常模式发送 | 部分内容逐步出现，最终回答完成 | Network 是否持续接收字节，done 是否存在 |
| error 模式发送 | 保留前八段文字，状态失败 | error 是否被当作普通文本吞掉 |
| incomplete 模式发送 | 保留部分文字，INCOMPLETE_STREAM | 是否错误地把 EOF 当成功 |
| 生成过程中停止 | 内容停止增加，显示不完整 | signal 是否传入 fetch，旧增量是否被忽略 |
| 停止后再次发送 | 新回答从空内容开始 | 旧 finally 是否清除了新控制器 |
| 输入中文时确认候选词 | 不应意外发送 | isComposing 与浏览器输入法事件 |
| 关闭后端再发送 | 界面显示失败并能再次操作 | 开发代理错误、按钮是否一直锁住 |

如果所有字最后一起显示，按三段排查：服务端是否每次调用 write；直连 Node 是否逐步收到数据；经过代理后是否仍然逐步到达。再看前端是否一直等待 response.text 才更新。只有在确认数据已经逐步到达浏览器后，讨论渲染批次才有意义。不要用“加一个 setTimeout”掩盖未知的缓冲层。

### 在这个实验上继续练到什么程度

第一步，增加 requestId 并显示于错误详情，不能把原始错误堆栈展示给普通用户。第二步，增加最大输出长度和空闲超时，并确保 done 到达时会清除定时器。第三步，把 answer 改成消息数组，每条消息仍拥有稳定 ID 和独立尝试，路由切换不能误取消其他窗口。第四步，再接入后端模型适配层，将供应商事件转换为本章的三类事件。

达到 L2 的标准是能照着契约实现这套链路并定位输入、网络、解析错误。你的前端主轴目标是 L3：能解释旧请求为什么会污染新请求，设计连续取消与重试的状态语义，用可控实验重现竞态，并把长文本、安全渲染与无障碍融入相同流程。单纯看到打字效果只完成了其中很小的一部分。


## 练习

把测试的分块方式改成随机 1 到 7 字节；加入 CRLF、多个事件同一 chunk、空行、损坏 JSON、缺少 done 和 done 后额外事件。增加总字符上限，超限后取消读流。

<details>
<summary>参考答案</summary>

分块测试应保持原始 wire 内容不变，只改变网络切分，所有合法分块都必须得到同一答案。统计所有 delta 长度而不只是 buffer 长度，以限制总回答大小。损坏 JSON 与协议违例进入 failed，已接收文本保留为不完整。生产中应记录错误类别，不记录可能包含敏感信息的整个坏事件。

</details>

## 验收与自测

- 任意字节边界拆分不会破坏中文或事件解析。
- HTTP 错误、空响应、协议错误和中途关闭能被区分。
- 用户取消和组件卸载不会留下活动请求。
- 快速连续请求不会产生状态串线。

**问：每次 read 得到一个 JSON 吗？** 答：没有保证，必须自行按协议聚合。

**问：SSE 是用换行分隔 JSON 吗？** 答：SSE 有自己的字段与空行事件边界，不能直接套 NDJSON 解析器。

**问：首字出现快就代表总耗时低吗？** 答：不是，需要同时衡量首字和完成时间。

## 官方参考

- [MDN：ReadableStream.getReader](https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream/getReader)
- [MDN：TextDecoder.decode](https://developer.mozilla.org/en-US/docs/Web/API/TextDecoder/decode)
- [MDN：使用服务器发送事件](https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events/Using_server-sent_events)
- [Vue：性能优化](https://cn.vuejs.org/guide/best-practices/performance)
