SSE Client
SSE Client 是前端消费 Server-Sent Events 的模块,负责建立连接、解析事件、处理重连和错误恢复。它是 AI 应用流式交互的前端基础。
#type / concept
#status / evergreen
#tech / frontend
#tech / network
[!info] related notes
- 所属 MOC: AI Agent Application MOC
- 协议: SSE (Server-Sent Events)
- 相关: Event Reducer, Abort Generation
- 实践: 前端 SSE 消费
SSE Client
一句话定义
SSE Client 是前端消费 Server-Sent Events 的模块。它不只是 new EventSource(),还需要处理 POST 请求、自定义 Header、错误恢复、断线重连和事件解析。
它解决什么问题
浏览器原生的 EventSource API 有几个限制:
- 只支持 GET: AI 场景通常用 POST 发送请求体
- 不支持自定义 Header: 无法传递 Auth Token
- 自动重连时丢失上下文: 重连后不知道从哪里继续
- 不支持请求体: 无法发送对话历史等数据
所以 AI 应用通常不用原生 EventSource,而是用 fetch + ReadableStream 手动实现 SSE 客户端。
核心原理
两种实现方式
| 方式 | 优点 | 缺点 |
|---|---|---|
EventSource | 自动重连、简单 | 只支持 GET、无自定义 Header |
fetch + Stream | 支持 POST/Headers/Body | 需手动处理重连和解析 |
fetch 实现 SSE 客户端
async function* sseClient(url: string, options: RequestInit) {
const response = await fetch(url, {
...options,
headers: {
...options.headers,
'Accept': 'text/event-stream',
},
});
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const events = parseSSEBuffer(buffer);
buffer = events.remaining;
for (const event of events.parsed) {
yield event;
}
}
}
function parseSSEBuffer(buffer: string) {
const lines = buffer.split('\n');
const parsed = [];
let remaining = '';
let currentEvent: any = {};
for (const line of lines) {
if (line === '') {
if (currentEvent.event || currentEvent.data) {
parsed.push(currentEvent);
currentEvent = {};
}
} else if (line.startsWith('event: ')) {
currentEvent.event = line.slice(7);
} else if (line.startsWith('data: ')) {
currentEvent.data = (currentEvent.data || '') + line.slice(6);
} else if (line.startsWith('id: ')) {
currentEvent.id = line.slice(4);
} else {
remaining += line + '\n';
}
}
return { parsed, remaining };
}
与 AbortController 集成
class SSEClient {
private controller: AbortController | null = null;
async connect(url: string, body: any, onEvent: (event: SSEEvent) => void) {
this.controller = new AbortController();
const response = await fetch(url, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(body),
signal: this.controller.signal,
});
const reader = response.body!.getReader();
// ... 读取流
}
disconnect() {
this.controller?.abort();
}
}
典型工程实现
React Hook
function useSSE() {
const [events, setEvents] = useState<SSEEvent[]>([]);
const [status, setStatus] = useState<'idle' | 'connecting' | 'streaming' | 'done' | 'error'>('idle');
const clientRef = useRef<SSEClient | null>(null);
const connect = async (url: string, body: any) => {
setStatus('connecting');
clientRef.current = new SSEClient();
await clientRef.current.connect(url, body, (event) => {
setEvents(prev => [...prev, event]);
setStatus('streaming');
if (event.event === 'done') {
setStatus('done');
}
});
};
const disconnect = () => {
clientRef.current?.disconnect();
setStatus('idle');
};
return { events, status, connect, disconnect };
}
常见坑
- 不做 buffer 处理: 一次
read()可能包含不完整的 SSE 事件 - 不做 AbortController: 组件卸载时连接没有断开
- 不处理网络错误: 断网时没有重连机制
- 不解析 event 类型: 所有事件都当
message处理 - 内存泄漏: 事件累积在数组中不清理