1
0

sse.ts 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161
  1. /**
  2. * SSE 流消费的统一实现。
  3. *
  4. */
  5. /** 本模块需要读取的最小事件形状;各业务自己的事件类型(SseEvent / SSERawEvent)都满足它 */
  6. export interface SseEventLike {
  7. type?: string;
  8. session_id?: string;
  9. thread_id?: string;
  10. }
  11. /** 消费完一条 SSE 流的结果 */
  12. export interface SseConsumeResult {
  13. /** 终止事件携带的会话 id,未携带时为空串 */
  14. sessionId: string;
  15. /** 终止事件携带的线程 id,未携带时为空串 */
  16. threadId: string;
  17. /** 解析失败被丢弃的 data 行数,0 表示整条流都解析正常 */
  18. parseErrors: number;
  19. }
  20. /** 一次解析尝试的结果 */
  21. type FlushOutcome = 'terminated' | 'dispatched' | 'incomplete' | 'invalid' | 'empty';
  22. /** 默认的终止事件类型:收到其中之一即停止读取并释放连接 */
  23. const DEFAULT_TERMINATE_ON = ['done', 'error', 'interrupt'];
  24. /**
  25. * 消费一个 SSE 响应体,逐事件回调,直到收到终止事件或流结束。
  26. * @param response - fetch 返回的响应(须带 body)
  27. * @param onEvent - 每个事件的回调
  28. * @param options.terminateOn - 终止事件类型,默认 done / error / interrupt
  29. */
  30. export async function consumeSse<E extends SseEventLike>(
  31. response: Response,
  32. onEvent: (event: E) => void,
  33. options: { terminateOn?: readonly string[] } = {}
  34. ): Promise<SseConsumeResult> {
  35. const terminateOn = options.terminateOn ?? DEFAULT_TERMINATE_ON;
  36. const body = response.body;
  37. if (!body) {
  38. throw new Error('SSE 响应缺少可读流');
  39. }
  40. const reader = body.getReader();
  41. const decoder = new TextDecoder();
  42. const result: SseConsumeResult = { sessionId: '', threadId: '', parseErrors: 0 };
  43. let buffer = '';
  44. /** 当前事件已收集、尚未解析的 data 行 */
  45. let pending: string[] = [];
  46. /** 丢弃坏数据时记账:解析失败不再无声跳过 */
  47. const dropBadData = (raw: string) => {
  48. result.parseErrors += 1;
  49. console.warn('SSE数据解析失败:', raw);
  50. };
  51. /**
  52. * 尝试把当前累积的 data 行解析成一个事件并派发。
  53. * 解析失败时保留 pending 不清理——是"多行事件的前半段"还是"坏数据"由调用方判断。
  54. */
  55. const flushPending = (): FlushOutcome => {
  56. if (pending.length === 0) return 'empty';
  57. const raw = pending.join('\n');
  58. // 规范规定 data 为空的事件不派发;服务端也有用 `data:` 空行做心跳的写法,故不算解析错误
  59. if (raw.trim() === '') {
  60. pending = [];
  61. return 'empty';
  62. }
  63. let event: E;
  64. try {
  65. event = JSON.parse(raw) as E;
  66. } catch {
  67. return 'incomplete';
  68. }
  69. // 非对象负载(`data: null` / `data: 123`)没有事件语义,交给回调只会让它崩在取字段上。
  70. // 按坏数据丢弃并继续读流——与重构前"回调抛错被 try 吞掉后继续"的最终表现一致。
  71. if (typeof event !== 'object' || !event) {
  72. dropBadData(raw);
  73. pending = [];
  74. return 'invalid';
  75. }
  76. pending = [];
  77. onEvent(event);
  78. if (event.type && terminateOn.includes(event.type)) {
  79. result.sessionId = event.session_id || '';
  80. result.threadId = event.thread_id || '';
  81. return 'terminated';
  82. }
  83. return 'dispatched';
  84. };
  85. /** 处理一行;返回 true 表示已收到终止事件 */
  86. const handleLine = (rawLine: string): boolean => {
  87. // 兼容 CRLF / CR 行尾
  88. const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine;
  89. // 空行是规范定义的事件边界:到这里还解析不出来,才能判定为坏数据并丢弃
  90. if (line === '') {
  91. const outcome = flushPending();
  92. if (outcome === 'incomplete') {
  93. dropBadData(pending.join('\n'));
  94. pending = [];
  95. }
  96. return outcome === 'terminated';
  97. }
  98. const colon = line.indexOf(':');
  99. // 没有冒号的整行都是字段名、冒号在行首则是注释行(心跳),两者都不消费;
  100. // 字段名与值之间允许 0 个空格,所以这里不能用 startsWith('data: ') 判断
  101. if (colon <= 0 || line.slice(0, colon) !== 'data') return false;
  102. const value = line.slice(colon + 1);
  103. pending.push(value.startsWith(' ') ? value.slice(1) : value);
  104. // 每收到一行就尝试解析累积的 data 行:
  105. // - 能解析 → 立即派发(兼容事件之间不发送空行的服务端,与原实现行为一致)
  106. // - 非对象负载 → 已记坏数据并清理干净,结束
  107. // - 解析不出来且已积累多行 → 首行是坏数据,丢掉它再试,避免一行坏 JSON 连坐后续事件
  108. // - 只剩一行 → 保留,它可能是多行事件的前半段
  109. for (;;) {
  110. const outcome = flushPending();
  111. if (outcome === 'terminated') return true;
  112. if (outcome === 'dispatched' || outcome === 'invalid') return false;
  113. if (pending.length <= 1) return false;
  114. dropBadData(pending.shift() as string);
  115. }
  116. };
  117. try {
  118. for (;;) {
  119. const { done, value } = await reader.read();
  120. if (done) break;
  121. buffer += decoder.decode(value, { stream: true });
  122. const lines = buffer.split('\n');
  123. buffer = lines.pop() || '';
  124. for (const line of lines) {
  125. if (handleLine(line)) return result;
  126. }
  127. }
  128. return result;
  129. } finally {
  130. try {
  131. await reader.cancel();
  132. } catch {
  133. // 流已自然结束,或已被调用方的 AbortController 中断,取消失败无需处理
  134. }
  135. }
  136. }