fix: 修复流式聊天发生错误时继续处理后续事件的问题
- 增加 hasError 标记,错误发生后立即停止处理后续事件 - 在 done 事件处理前检查是否已有错误 - 在 finally 中确保 reader 被正确关闭 - 防止错误发生后,done 事件继续执行导致历史响应被错误回填feat/unify-api-and-responsive-pages
parent
c30d301fbc
commit
92ef2c38c8
141
src/api.ts
141
src/api.ts
|
|
@ -748,83 +748,94 @@ export async function regenerateMessage(
|
||||||
return await consumeSSE(resp, handlers, signal);
|
return await consumeSSE(resp, handlers, signal);
|
||||||
}
|
}
|
||||||
|
|
||||||
async function consumeSSE(resp: Response, h: StreamEvents, signal?: AbortSignal) {
|
async function consumeSSE(resp: Response, h: StreamEvents, signal?: AbortSignal) {
|
||||||
if (!resp.ok || !resp.body) {
|
if (!resp.ok || !resp.body) {
|
||||||
const txt = await resp.text().catch(() => '');
|
const txt = await resp.text().catch(() => '');
|
||||||
h.onError?.(`HTTP ${resp.status}: ${txt}`);
|
h.onError?.(`HTTP ${resp.status}: ${txt}`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
const reader = resp.body.getReader();
|
const reader = resp.body.getReader();
|
||||||
const decoder = new TextDecoder('utf-8');
|
const decoder = new TextDecoder('utf-8');
|
||||||
let buf = '';
|
let buf = '';
|
||||||
try {
|
let hasError = false; // 标记是否已发生错误,避免后续事件继续处理
|
||||||
while (true) {
|
|
||||||
const { value, done } = await reader.read();
|
try {
|
||||||
if (done) break;
|
while (true) {
|
||||||
buf += decoder.decode(value, { stream: true });
|
const { value, done } = await reader.read();
|
||||||
|
if (done || hasError) break;
|
||||||
|
|
||||||
|
buf += decoder.decode(value, { stream: true });
|
||||||
buf = buf.replace(/\r\n/g, '\n');
|
buf = buf.replace(/\r\n/g, '\n');
|
||||||
|
|
||||||
let idx;
|
let idx;
|
||||||
while ((idx = buf.indexOf('\n\n')) !== -1) {
|
while ((idx = buf.indexOf('\n\n')) !== -1 && !hasError) {
|
||||||
const raw = buf.slice(0, idx);
|
const raw = buf.slice(0, idx);
|
||||||
buf = buf.slice(idx + 2);
|
buf = buf.slice(idx + 2);
|
||||||
if (!raw.trim() || raw.startsWith(':')) continue;
|
if (!raw.trim() || raw.startsWith(':')) continue;
|
||||||
|
|
||||||
let event = 'message';
|
let event = 'message';
|
||||||
let dataStr = '';
|
let dataStr = '';
|
||||||
for (const line of raw.split('\n')) {
|
for (const line of raw.split('\n')) {
|
||||||
if (line.startsWith('event:')) event = line.slice(6).trim();
|
if (line.startsWith('event:')) event = line.slice(6).trim();
|
||||||
else if (line.startsWith('data:')) {
|
else if (line.startsWith('data:')) {
|
||||||
let part = line.slice(5);
|
let part = line.slice(5);
|
||||||
if (part.startsWith(' ')) part = part.slice(1);
|
if (part.startsWith(' ')) part = part.slice(1);
|
||||||
dataStr += (dataStr ? '\n' : '') + part;
|
dataStr += (dataStr ? '\n' : '') + part;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (!dataStr) continue;
|
if (!dataStr) continue;
|
||||||
let data: any;
|
let data: any;
|
||||||
try {
|
try {
|
||||||
data = JSON.parse(dataStr);
|
data = JSON.parse(dataStr);
|
||||||
} catch {
|
} catch {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
switch (event) {
|
switch (event) {
|
||||||
case 'meta':
|
case 'meta':
|
||||||
h.onMeta?.(data);
|
h.onMeta?.(data);
|
||||||
break;
|
break;
|
||||||
case 'retry':
|
case 'retry':
|
||||||
h.onRetry?.(data);
|
h.onRetry?.(data);
|
||||||
break;
|
break;
|
||||||
case 'reasoning_delta':
|
case 'reasoning_delta':
|
||||||
h.onReasoningDelta?.(data.content || '');
|
h.onReasoningDelta?.(data.content || '');
|
||||||
break;
|
break;
|
||||||
case 'delta':
|
case 'delta':
|
||||||
h.onDelta?.(data.content || '');
|
h.onDelta?.(data.content || '');
|
||||||
break;
|
break;
|
||||||
case 'tool_call':
|
case 'tool_call':
|
||||||
h.onToolCall?.(data);
|
h.onToolCall?.(data);
|
||||||
break;
|
break;
|
||||||
case 'tool_result':
|
case 'tool_result':
|
||||||
h.onToolResult?.(data);
|
h.onToolResult?.(data);
|
||||||
break;
|
break;
|
||||||
case 'done':
|
case 'done':
|
||||||
h.onDone?.(data);
|
// 只有在没有错误的情况下才处理 done 事件
|
||||||
break;
|
if (!hasError) {
|
||||||
case 'aborted':
|
h.onDone?.(data);
|
||||||
h.onAborted?.(data);
|
}
|
||||||
break;
|
break;
|
||||||
case 'error':
|
case 'aborted':
|
||||||
h.onError?.(data.message || 'stream error');
|
h.onAborted?.(data);
|
||||||
break;
|
break;
|
||||||
}
|
case 'error':
|
||||||
}
|
hasError = true;
|
||||||
}
|
h.onError?.(data.message || 'stream error');
|
||||||
} catch (e: any) {
|
// 发生错误后立即停止读取
|
||||||
if (signal?.aborted || e?.name === 'AbortError') {
|
break;
|
||||||
// 静默:本地已 abort
|
}
|
||||||
return;
|
}
|
||||||
}
|
}
|
||||||
h.onError?.(e?.message ?? String(e));
|
} catch (e: any) {
|
||||||
}
|
if (signal?.aborted || e?.name === 'AbortError') {
|
||||||
|
// 静默:本地已 abort
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
h.onError?.(e?.message ?? String(e));
|
||||||
|
} finally {
|
||||||
|
// 确保 reader 被释放
|
||||||
|
reader.cancel().catch(() => {});
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ============== Workflow (v1.1) ==============
|
// ============== Workflow (v1.1) ==============
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue