Kaynağa Gözat

[Mod 0000]AI助手sse事件逻辑修改

wangkeyi 1 gün önce
ebeveyn
işleme
eccff637b5

+ 3 - 42
src/views/ventAI/dataPicker/api.ts

@@ -1,3 +1,4 @@
+import { consumeSse } from '../sse';
 import type { SSERawEvent } from './types';
 
 enum Api {
@@ -47,48 +48,8 @@ async function ssePost(
     throw new Error(`HTTP错误: ${response.status}`);
   }
 
-  const reader = response.body!.getReader();
-  const decoder = new TextDecoder();
-  let buffer = '';
-  let threadId = '';
-  let sessionId = '';
-  let isDone = false;
-
-  // eslint-disable-next-line no-constant-condition
-  while (true) {
-    const { done, value } = await reader.read();
-    if (done || isDone) break;
-
-    buffer += decoder.decode(value, { stream: true });
-    const lines = buffer.split('\n');
-    buffer = lines.pop() || '';
-
-    for (const line of lines) {
-      if (line.trim() === '') continue;
-
-      if (line.startsWith('data: ')) {
-        const dataStr = line.replace(/^data: /, '');
-
-        try {
-          const event: SSERawEvent = JSON.parse(dataStr);
-          onEvent(event);
-
-          if (event.type === 'done') {
-            threadId = event.thread_id || '';
-            sessionId = event.session_id || '';
-            isDone = true;
-            break;
-          }
-          if (event.type === 'error') {
-            isDone = true;
-            break;
-          }
-        } catch (e) {
-          console.warn('SSE数据解析失败:', dataStr);
-        }
-      }
-    }
-  }
+  // 该接口只把 done / error 当终止事件(与统一对话不同,不含 interrupt)
+  const { sessionId, threadId } = await consumeSse<SSERawEvent>(response, onEvent, { terminateOn: ['done', 'error'] });
 
   return { thread_id: threadId, session_id: sessionId };
 }

+ 19 - 132
src/views/ventAI/manageAssistent/api.ts

@@ -1,5 +1,6 @@
 import { defHttp } from '/@/utils/http/axios';
 import { getToken } from '/@/utils/auth';
+import { consumeSse } from '../sse';
 import type { SseEvent } from './components/chatModal/types';
 
 enum Api {
@@ -190,12 +191,13 @@ export const uploadSkill = (file: File) => {
  * @param params.session_id - 会话唯一标识ID,不传参时服务端自动生成全新会话ID
  * @param params.file - 上传的PDF文件,可选
  * @param onChunk - 流式数据回调函数
- * @returns Promise<{ session_id: string }>
+ * @returns Promise<{ session_id: string; parse_errors: number }>
+ *          parse_errors 大于 0 表示这条流里有事件被丢弃(正文可能不完整)
  */
 export const unifiedStream = async (
   params: { message: string; session_id?: string; file?: File; mode?: string; signal?: AbortSignal },
   onChunk: (data: SseEvent) => void
-): Promise<{ session_id: string }> => {
+): Promise<{ session_id: string; parse_errors: number }> => {
   try {
     const formData = new FormData();
     formData.append('message', params.message);
@@ -220,49 +222,10 @@ export const unifiedStream = async (
       await throwStreamHttpError(response);
     }
 
-    const reader = response.body!.getReader();
-    const decoder = new TextDecoder();
-    let buffer = '';
-    let sessionId = '';
-    let isDone = false;
-
-    // eslint-disable-next-line no-constant-condition
-    while (true) {
-      const { done, value } = await reader.read();
-      if (done || isDone) break;
-
-      buffer += decoder.decode(value, { stream: true });
-      const lines = buffer.split('\n');
-      buffer = lines.pop() || '';
-
-      for (const line of lines) {
-        if (line.trim() === '') continue;
-
-        if (line.startsWith('data: ')) {
-          const dataStr = line.replace(/^data: /, '');
-
-          try {
-            const data = JSON.parse(dataStr) as SseEvent;
-
-            onChunk(data);
-
-            if (data.type === 'done' || data.type === 'error' || data.type === 'interrupt') {
-              if (data.session_id) {
-                sessionId = data.session_id;
-              }
-              isDone = true;
-              break;
-            }
-          } catch (e) {
-            console.warn('SSE数据解析失败:', dataStr);
-          }
-        }
-      }
-
-      if (isDone) break;
-    }
+    const { sessionId, parseErrors } = await consumeSse<SseEvent>(response, onChunk);
 
-    return { session_id: sessionId };
+    // parseErrors > 0 表示这条流里有事件被丢弃,正文可能不完整
+    return { session_id: sessionId, parse_errors: parseErrors };
   } catch (error) {
     if ((error as Error)?.name !== 'AbortError') {
       console.error('统一对话流式请求失败:', error);
@@ -277,12 +240,13 @@ export const unifiedStream = async (
  * @param params.session_id - 会话唯一标识ID,可选
  * @param params.message - 附加消息/审查指令,可选
  * @param onChunk - 流式数据回调函数
- * @returns Promise<{ session_id: string }>
+ * @returns Promise<{ session_id: string; parse_errors: number }>
+ *          parse_errors 大于 0 表示这条流里有事件被丢弃(正文可能不完整)
  */
 export const reviewPdfStream = async (
   params: { file: File; session_id?: string; message?: string; mode?: string; signal?: AbortSignal },
   onChunk: (data: SseEvent) => void
-): Promise<{ session_id: string }> => {
+): Promise<{ session_id: string; parse_errors: number }> => {
   try {
     const formData = new FormData();
     formData.append('file', params.file);
@@ -307,49 +271,10 @@ export const reviewPdfStream = async (
       await throwStreamHttpError(response);
     }
 
-    const reader = response.body!.getReader();
-    const decoder = new TextDecoder();
-    let buffer = '';
-    let sessionId = '';
-    let isDone = false;
-
-    // eslint-disable-next-line no-constant-condition
-    while (true) {
-      const { done, value } = await reader.read();
-      if (done || isDone) break;
-
-      buffer += decoder.decode(value, { stream: true });
-      const lines = buffer.split('\n');
-      buffer = lines.pop() || '';
-
-      for (const line of lines) {
-        if (line.trim() === '') continue;
-
-        if (line.startsWith('data: ')) {
-          const dataStr = line.replace(/^data: /, '');
-
-          try {
-            const data = JSON.parse(dataStr) as SseEvent;
-
-            onChunk(data);
-
-            if (data.type === 'done' || data.type === 'error' || data.type === 'interrupt') {
-              if (data.session_id) {
-                sessionId = data.session_id;
-              }
-              isDone = true;
-              break;
-            }
-          } catch (e) {
-            console.warn('SSE数据解析失败:', dataStr);
-          }
-        }
-      }
+    const { sessionId, parseErrors } = await consumeSse<SseEvent>(response, onChunk);
 
-      if (isDone) break;
-    }
-
-    return { session_id: sessionId };
+    // parseErrors > 0 表示这条流里有事件被丢弃,正文可能不完整
+    return { session_id: sessionId, parse_errors: parseErrors };
   } catch (error) {
     if ((error as Error)?.name !== 'AbortError') {
       console.error('PDF审查流式请求失败:', error);
@@ -365,12 +290,13 @@ export const reviewPdfStream = async (
  * @param params.status - ask_user响应状态:"answered"(默认,消费answers)/ "cancelled"(取消)/ "error"(出错)
  * @param params.answers - ask_user答案列表,与中断事件ask_user.questions一一对应(status=answered时必填)
  * @param onChunk - 流式数据回调函数
- * @returns Promise<{ session_id: string }>
+ * @returns Promise<{ session_id: string; parse_errors: number }>
+ *          parse_errors 大于 0 表示这条流里有事件被丢弃(正文可能不完整)
  */
 export const chatResumeStream = async (
   params: { session_id: string; thread_id?: string; status?: string; answers?: string[]; signal?: AbortSignal },
   onChunk: (data: SseEvent) => void
-): Promise<{ session_id: string }> => {
+): Promise<{ session_id: string; parse_errors: number }> => {
   try {
     const body: Record<string, any> = {
       session_id: params.session_id,
@@ -397,49 +323,10 @@ export const chatResumeStream = async (
       await throwStreamHttpError(response);
     }
 
-    const reader = response.body!.getReader();
-    const decoder = new TextDecoder();
-    let buffer = '';
-    let sessionId = '';
-    let isDone = false;
-
-    // eslint-disable-next-line no-constant-condition
-    while (true) {
-      const { done, value } = await reader.read();
-      if (done || isDone) break;
-
-      buffer += decoder.decode(value, { stream: true });
-      const lines = buffer.split('\n');
-      buffer = lines.pop() || '';
-
-      for (const line of lines) {
-        if (line.trim() === '') continue;
-
-        if (line.startsWith('data: ')) {
-          const dataStr = line.replace(/^data: /, '');
-
-          try {
-            const data = JSON.parse(dataStr) as SseEvent;
-
-            onChunk(data);
-
-            if (data.type === 'done' || data.type === 'error' || data.type === 'interrupt') {
-              if (data.session_id) {
-                sessionId = data.session_id;
-              }
-              isDone = true;
-              break;
-            }
-          } catch (e) {
-            console.warn('SSE数据解析失败:', dataStr);
-          }
-        }
-      }
-
-      if (isDone) break;
-    }
+    const { sessionId, parseErrors } = await consumeSse<SseEvent>(response, onChunk);
 
-    return { session_id: sessionId };
+    // parseErrors > 0 表示这条流里有事件被丢弃,正文可能不完整
+    return { session_id: sessionId, parse_errors: parseErrors };
   } catch (error) {
     if ((error as Error)?.name !== 'AbortError') {
       console.error('恢复会话流式请求失败:', error);

+ 161 - 0
src/views/ventAI/sse.ts

@@ -0,0 +1,161 @@
+/**
+ * SSE 流消费的统一实现。
+ *
+ */
+
+/** 本模块需要读取的最小事件形状;各业务自己的事件类型(SseEvent / SSERawEvent)都满足它 */
+export interface SseEventLike {
+  type?: string;
+  session_id?: string;
+  thread_id?: string;
+}
+
+/** 消费完一条 SSE 流的结果 */
+export interface SseConsumeResult {
+  /** 终止事件携带的会话 id,未携带时为空串 */
+  sessionId: string;
+  /** 终止事件携带的线程 id,未携带时为空串 */
+  threadId: string;
+  /** 解析失败被丢弃的 data 行数,0 表示整条流都解析正常 */
+  parseErrors: number;
+}
+
+/** 一次解析尝试的结果 */
+type FlushOutcome = 'terminated' | 'dispatched' | 'incomplete' | 'invalid' | 'empty';
+
+/** 默认的终止事件类型:收到其中之一即停止读取并释放连接 */
+const DEFAULT_TERMINATE_ON = ['done', 'error', 'interrupt'];
+
+/**
+ * 消费一个 SSE 响应体,逐事件回调,直到收到终止事件或流结束。
+ * @param response - fetch 返回的响应(须带 body)
+ * @param onEvent - 每个事件的回调
+ * @param options.terminateOn - 终止事件类型,默认 done / error / interrupt
+ */
+export async function consumeSse<E extends SseEventLike>(
+  response: Response,
+  onEvent: (event: E) => void,
+  options: { terminateOn?: readonly string[] } = {}
+): Promise<SseConsumeResult> {
+  const terminateOn = options.terminateOn ?? DEFAULT_TERMINATE_ON;
+
+  const body = response.body;
+  if (!body) {
+    throw new Error('SSE 响应缺少可读流');
+  }
+
+  const reader = body.getReader();
+  const decoder = new TextDecoder();
+
+  const result: SseConsumeResult = { sessionId: '', threadId: '', parseErrors: 0 };
+  let buffer = '';
+  /** 当前事件已收集、尚未解析的 data 行 */
+  let pending: string[] = [];
+
+  /** 丢弃坏数据时记账:解析失败不再无声跳过 */
+  const dropBadData = (raw: string) => {
+    result.parseErrors += 1;
+    console.warn('SSE数据解析失败:', raw);
+  };
+
+  /**
+   * 尝试把当前累积的 data 行解析成一个事件并派发。
+   * 解析失败时保留 pending 不清理——是"多行事件的前半段"还是"坏数据"由调用方判断。
+   */
+  const flushPending = (): FlushOutcome => {
+    if (pending.length === 0) return 'empty';
+
+    const raw = pending.join('\n');
+    // 规范规定 data 为空的事件不派发;服务端也有用 `data:` 空行做心跳的写法,故不算解析错误
+    if (raw.trim() === '') {
+      pending = [];
+      return 'empty';
+    }
+
+    let event: E;
+    try {
+      event = JSON.parse(raw) as E;
+    } catch {
+      return 'incomplete';
+    }
+
+    // 非对象负载(`data: null` / `data: 123`)没有事件语义,交给回调只会让它崩在取字段上。
+    // 按坏数据丢弃并继续读流——与重构前"回调抛错被 try 吞掉后继续"的最终表现一致。
+    if (typeof event !== 'object' || !event) {
+      dropBadData(raw);
+      pending = [];
+      return 'invalid';
+    }
+
+    pending = [];
+    onEvent(event);
+
+    if (event.type && terminateOn.includes(event.type)) {
+      result.sessionId = event.session_id || '';
+      result.threadId = event.thread_id || '';
+      return 'terminated';
+    }
+    return 'dispatched';
+  };
+
+  /** 处理一行;返回 true 表示已收到终止事件 */
+  const handleLine = (rawLine: string): boolean => {
+    // 兼容 CRLF / CR 行尾
+    const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine;
+
+    // 空行是规范定义的事件边界:到这里还解析不出来,才能判定为坏数据并丢弃
+    if (line === '') {
+      const outcome = flushPending();
+      if (outcome === 'incomplete') {
+        dropBadData(pending.join('\n'));
+        pending = [];
+      }
+      return outcome === 'terminated';
+    }
+
+    const colon = line.indexOf(':');
+    // 没有冒号的整行都是字段名、冒号在行首则是注释行(心跳),两者都不消费;
+    // 字段名与值之间允许 0 个空格,所以这里不能用 startsWith('data: ') 判断
+    if (colon <= 0 || line.slice(0, colon) !== 'data') return false;
+
+    const value = line.slice(colon + 1);
+    pending.push(value.startsWith(' ') ? value.slice(1) : value);
+
+    // 每收到一行就尝试解析累积的 data 行:
+    // - 能解析 → 立即派发(兼容事件之间不发送空行的服务端,与原实现行为一致)
+    // - 非对象负载 → 已记坏数据并清理干净,结束
+    // - 解析不出来且已积累多行 → 首行是坏数据,丢掉它再试,避免一行坏 JSON 连坐后续事件
+    // - 只剩一行 → 保留,它可能是多行事件的前半段
+    for (;;) {
+      const outcome = flushPending();
+      if (outcome === 'terminated') return true;
+      if (outcome === 'dispatched' || outcome === 'invalid') return false;
+      if (pending.length <= 1) return false;
+
+      dropBadData(pending.shift() as string);
+    }
+  };
+
+  try {
+    for (;;) {
+      const { done, value } = await reader.read();
+      if (done) break;
+
+      buffer += decoder.decode(value, { stream: true });
+      const lines = buffer.split('\n');
+      buffer = lines.pop() || '';
+
+      for (const line of lines) {
+        if (handleLine(line)) return result;
+      }
+    }
+
+    return result;
+  } finally {
+    try {
+      await reader.cancel();
+    } catch {
+      // 流已自然结束,或已被调用方的 AbortController 中断,取消失败无需处理
+    }
+  }
+}