Skip to content

Commit 2834ca8

Browse files
committed
perf: v1 abort check
1 parent 71354ff commit 2834ca8

6 files changed

Lines changed: 206 additions & 4 deletions

File tree

packages/global/core/workflow/runtime/type.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import type { NodeInputKeyEnum } from '../constants';
1212
import { NodeOutputKeyEnum } from '../constants';
1313
import { ClassifyQuestionAgentItemSchema } from '../template/system/classifyQuestion/type';
1414
import type { NextApiResponse } from 'next';
15+
import type { IncomingMessage } from 'node:http';
1516
import type { AppSchemaType } from '../../app/type';
1617
import type { RuntimeEdgeItemType } from '../type/edge';
1718
import { ReadFileNodeResponseSchema } from '../template/system/readFiles/type';
@@ -59,6 +60,7 @@ export type WorkflowVariableStateLike = {
5960

6061
/* workflow props */
6162
export type ChatDispatchProps = {
63+
req?: IncomingMessage;
6264
res?: NextApiResponse;
6365
checkIsStopping: () => boolean;
6466
lang?: localeType;

packages/service/core/workflow/dispatch/index.ts

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ import { observeWorkflowRun, observeWorkflowStep } from '../metrics';
6363
import { withActiveSpan } from '../../../common/tracing';
6464
import { delAgentRuntimeStopSign, shouldWorkflowStop } from './workflowStatus';
6565
import { runWithContext } from '../utils/context';
66+
import { createClientAbortTracker } from './utils/clientAbort';
6667

6768
const logger = getLogger(LogCategories.MODULE.WORKFLOW.DISPATCH);
6869

@@ -165,7 +166,8 @@ export async function dispatchWorkFlow({
165166
histories,
166167
query,
167168
chatId,
168-
apiVersion
169+
apiVersion,
170+
req
169171
} = data;
170172

171173
// Check url valid
@@ -247,6 +249,14 @@ export async function dispatchWorkFlow({
247249
}
248250
}
249251

252+
const clientAbortTracker =
253+
apiVersion === 'v1'
254+
? createClientAbortTracker({
255+
req,
256+
res
257+
})
258+
: undefined;
259+
250260
const variableState = await WorkflowVariableState.create({
251261
timezone,
252262
runningAppInfo,
@@ -266,8 +276,7 @@ export async function dispatchWorkFlow({
266276
return stopping;
267277
}
268278
if (apiVersion === 'v1') {
269-
if (!res) return false;
270-
return res.closed || !!res.errored;
279+
return !!clientAbortTracker?.isClientAborted();
271280
}
272281
return false;
273282
};
@@ -315,6 +324,7 @@ export async function dispatchWorkFlow({
315324
if (checkStoppingTimer) {
316325
clearInterval(checkStoppingTimer);
317326
}
327+
clientAbortTracker?.cleanup();
318328

319329
// Close mcpClient connections
320330
Object.values(ctx.mcpClientMemory).forEach((client) => {
Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
import type { NextApiResponse } from 'next';
2+
import type { IncomingMessage } from 'node:http';
3+
4+
const getErrorCode = (error: unknown) => {
5+
if (typeof error !== 'object' || error === null || !('code' in error)) {
6+
return '';
7+
}
8+
9+
return String((error as { code?: unknown }).code);
10+
};
11+
12+
const isClientResetSocketError = (error: unknown) => {
13+
const code = getErrorCode(error);
14+
return code === 'ECONNRESET' || code === 'EPIPE';
15+
};
16+
17+
export const createClientAbortTracker = ({
18+
req,
19+
res
20+
}: {
21+
req?: IncomingMessage;
22+
res?: NextApiResponse;
23+
}) => {
24+
let clientAborted = false;
25+
let requestAborted = !!req?.aborted;
26+
let responseCompleted = !!(res?.writableEnded || res?.writableFinished);
27+
let responseError = !!res?.errored;
28+
let serverSocketError = false;
29+
30+
/**
31+
* v1 工作流只在客户端主动断开当前响应时停止。
32+
*
33+
* `socket.destroyed`、`res.destroyed`、`writableAborted` 这类快照过宽,不能单独证明用户取消。
34+
* `req.aborted` 是首选信号;但 Next API 运行时下 fetch abort 可能只稳定落到未 finish 的
35+
* `res.close`,因此把 close 作为 fallback。服务端 response/socket error 会屏蔽该 fallback,
36+
* 但客户端 reset 类 socket error 仍属于断开信号,不能抢先屏蔽后续 `req.aborted`/`res.close`。
37+
*/
38+
const responseFinished = () =>
39+
responseCompleted || !!(res?.writableEnded || res?.writableFinished);
40+
const responseErrored = () => responseError || serverSocketError || !!res?.errored;
41+
const canAcceptRequestAbort = () => !responseFinished() && !responseErrored();
42+
const isRequestAbortedSnapshot = () => requestAborted || !!req?.aborted;
43+
const markResponseCompleted = () => {
44+
responseCompleted = true;
45+
};
46+
const markResponseError = (error: unknown) => {
47+
if (!isClientResetSocketError(error)) {
48+
responseError = true;
49+
}
50+
};
51+
const markSocketError = (error: unknown) => {
52+
if (!isClientResetSocketError(error)) {
53+
serverSocketError = true;
54+
}
55+
};
56+
const markClientAborted = () => {
57+
requestAborted = true;
58+
clientAborted = true;
59+
};
60+
const markResponseClosed = () => {
61+
if (canAcceptRequestAbort()) {
62+
clientAborted = true;
63+
}
64+
};
65+
66+
req?.on('aborted', markClientAborted);
67+
req?.socket?.on('error', markSocketError);
68+
res?.on('finish', markResponseCompleted);
69+
res?.on('error', markResponseError);
70+
res?.on('close', markResponseClosed);
71+
72+
return {
73+
isClientAborted: () => clientAborted || isRequestAbortedSnapshot(),
74+
cleanup: () => {
75+
req?.off('aborted', markClientAborted);
76+
req?.socket?.off('error', markSocketError);
77+
res?.off('finish', markResponseCompleted);
78+
res?.off('error', markResponseError);
79+
res?.off('close', markResponseClosed);
80+
}
81+
};
82+
};
Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,107 @@
1+
import { EventEmitter } from 'node:events';
2+
import { describe, expect, test } from 'vitest';
3+
import { createClientAbortTracker } from '@fastgpt/service/core/workflow/dispatch/utils/clientAbort';
4+
5+
const createMockReq = () =>
6+
Object.assign(new EventEmitter(), {
7+
aborted: false,
8+
socket: new EventEmitter()
9+
});
10+
11+
const createMockRes = () =>
12+
Object.assign(new EventEmitter(), {
13+
writableEnded: false,
14+
writableFinished: false,
15+
errored: false
16+
});
17+
18+
describe('createClientAbortTracker', () => {
19+
test('detects client request abort before response finish', () => {
20+
const req = createMockReq();
21+
const res = createMockRes();
22+
const tracker = createClientAbortTracker({ req, res: res as any });
23+
24+
req.emit('aborted');
25+
26+
expect(tracker.isClientAborted()).toBe(true);
27+
tracker.cleanup();
28+
});
29+
30+
test('does not treat close after normal finish as client abort', () => {
31+
const req = createMockReq();
32+
const res = createMockRes();
33+
const tracker = createClientAbortTracker({ req, res: res as any });
34+
35+
res.emit('finish');
36+
res.writableEnded = true;
37+
res.emit('close');
38+
39+
expect(tracker.isClientAborted()).toBe(false);
40+
tracker.cleanup();
41+
});
42+
43+
test('does not treat response close after server error as client abort', () => {
44+
const req = createMockReq();
45+
const res = createMockRes();
46+
const tracker = createClientAbortTracker({ req, res: res as any });
47+
48+
res.emit('error');
49+
res.errored = true;
50+
res.emit('close');
51+
52+
expect(tracker.isClientAborted()).toBe(false);
53+
tracker.cleanup();
54+
});
55+
56+
test('keeps client reset socket errors eligible for close fallback', () => {
57+
const req = createMockReq();
58+
const res = createMockRes();
59+
const tracker = createClientAbortTracker({ req, res: res as any });
60+
const resetError = Object.assign(new Error('socket reset'), { code: 'ECONNRESET' });
61+
62+
req.socket.emit('error', resetError);
63+
res.emit('close');
64+
65+
expect(tracker.isClientAborted()).toBe(true);
66+
tracker.cleanup();
67+
});
68+
69+
test('keeps client reset response errors eligible for close fallback', () => {
70+
const req = createMockReq();
71+
const res = createMockRes();
72+
const tracker = createClientAbortTracker({ req, res: res as any });
73+
const resetError = Object.assign(new Error('write reset'), { code: 'EPIPE' });
74+
75+
res.emit('error', resetError);
76+
res.emit('close');
77+
78+
expect(tracker.isClientAborted()).toBe(true);
79+
tracker.cleanup();
80+
});
81+
82+
test('keeps accepted client abort sticky after later response error', () => {
83+
const req = createMockReq();
84+
const res = createMockRes();
85+
const tracker = createClientAbortTracker({ req, res: res as any });
86+
87+
res.emit('close');
88+
res.emit('error');
89+
res.errored = true;
90+
91+
expect(tracker.isClientAborted()).toBe(true);
92+
tracker.cleanup();
93+
});
94+
95+
test('keeps explicit request abort even after response error', () => {
96+
const req = createMockReq();
97+
const res = createMockRes();
98+
const tracker = createClientAbortTracker({ req, res: res as any });
99+
100+
res.emit('error');
101+
res.errored = true;
102+
req.emit('aborted');
103+
104+
expect(tracker.isClientAborted()).toBe(true);
105+
tracker.cleanup();
106+
});
107+
});

pro

Submodule pro updated from 7b80643 to 780360a

projects/app/src/pages/api/v1/chat/completions.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -276,6 +276,7 @@ async function handler(req: NextApiRequest, res: NextApiResponse) {
276276
if (app.version === 'v2') {
277277
return dispatchWorkFlow({
278278
apiVersion: 'v1',
279+
req,
279280
res,
280281
lang: getLocale(req),
281282
requestOrigin: req.headers.origin,

0 commit comments

Comments
 (0)