Skip to content

Commit ff669c9

Browse files
authored
Merge pull request #1971 from Hexastack/fix/ws-timeout
fix(ws): web channel ws timeout
2 parents 2c96a4a + c188163 commit ff669c9

3 files changed

Lines changed: 210 additions & 19 deletions

File tree

packages/api/src/extensions/channels/web/__test__/index.spec.ts

Lines changed: 131 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,17 @@ import { WebsocketGateway } from '@/websocket/websocket.gateway';
3232

3333
import WebChannelHandler from '../index.channel';
3434

35+
const createDeferred = () => {
36+
let resolve!: () => void;
37+
let reject!: (reason?: unknown) => void;
38+
const promise = new Promise<void>((res, rej) => {
39+
resolve = res;
40+
reject = rej;
41+
});
42+
43+
return { promise, resolve, reject };
44+
};
45+
3546
describe('WebChannelHandler', () => {
3647
let module: TestingModule;
3748
let subscriberService: SubscriberService;
@@ -471,9 +482,8 @@ describe('WebChannelHandler', () => {
471482

472483
const emitMessageSpy = jest
473484
.spyOn(handler['channelEventBus'], 'emitMessage')
474-
.mockImplementation(async (event: any) => {
475-
event.setThreadId('thread-created-1');
476-
});
485+
.mockResolvedValue(undefined);
486+
let responseThreadId: string | undefined;
477487

478488
await new Promise<void>((resolve, reject) => {
479489
const timeout = setTimeout(() => {
@@ -487,16 +497,18 @@ describe('WebChannelHandler', () => {
487497
},
488498
json: (payload: any) => {
489499
clearTimeout(timeout);
490-
expect(payload.thread_id).toBe('thread-created-1');
500+
expect(typeof payload.thread_id).toBe('string');
501+
responseThreadId = payload.thread_id;
491502
resolve();
492503
},
493504
} as any as SocketResponse;
494505

495506
handler['handleEvent'](req as any, res, webSource);
496507
});
497508

509+
await Promise.resolve();
498510
expect(emitMessageSpy).toHaveBeenCalledWith(expect.anything());
499-
expect(req.session.web?.threadId).toBe('thread-created-1');
511+
expect(req.session.web?.threadId).toBe(responseThreadId);
500512
clearMock.mockRestore();
501513
emitMessageSpy.mockRestore();
502514
});
@@ -553,12 +565,13 @@ describe('WebChannelHandler', () => {
553565
handler['handleEvent'](req as any, res, webSource);
554566
});
555567

568+
await Promise.resolve();
556569
expect(req.session.web?.profile?.id).toBe(subscriber.id);
557570
expect(emitMessageSpy).toHaveBeenCalledWith(expect.anything());
558571
emitMessageSpy.mockRestore();
559572
});
560573

561-
it('broadcasts incoming user messages before async handlers finish', async () => {
574+
it('acknowledges and broadcasts incoming user messages before async handlers finish', async () => {
562575
websocketGatewayMock.broadcast.mockClear();
563576
const req = {
564577
isSocket: true,
@@ -588,13 +601,14 @@ describe('WebChannelHandler', () => {
588601
profile,
589602
};
590603

604+
const pendingChatbotWork = Promise.race<void>([]);
591605
const emitMessageSpy = jest
592606
.spyOn(handler['channelEventBus'], 'emitMessage')
593-
.mockResolvedValue(undefined);
607+
.mockImplementation(() => pendingChatbotWork);
594608

595609
await new Promise<void>((resolve, reject) => {
596610
const timeout = setTimeout(() => {
597-
reject(new Error('Timed out waiting for message response'));
611+
reject(new Error('Timed out waiting for immediate message response'));
598612
}, 2000);
599613
const res = {
600614
status: (code: number) => {
@@ -611,17 +625,121 @@ describe('WebChannelHandler', () => {
611625
handler['handleEvent'](req as any, res, webSource);
612626
});
613627

628+
await Promise.resolve();
614629
expect(websocketGatewayMock.broadcast).toHaveBeenCalled();
615630
expect(emitMessageSpy).toHaveBeenCalledWith(expect.anything());
616-
const [broadcastCallOrder] = websocketGatewayMock.broadcast.mock
617-
.invocationCallOrder as number[];
618-
const [emitMessageCallOrder] = emitMessageSpy.mock.invocationCallOrder;
631+
const broadcastCallOrder =
632+
websocketGatewayMock.broadcast.mock.invocationCallOrder.at(-1);
633+
const emitMessageCallOrder = emitMessageSpy.mock.invocationCallOrder.at(-1);
619634

620-
expect(broadcastCallOrder).toBeLessThan(emitMessageCallOrder);
635+
expect(broadcastCallOrder!).toBeLessThan(emitMessageCallOrder!);
636+
expect(req.session.web?.threadId).toEqual(expect.any(String));
621637
clearMock.mockRestore();
622638
emitMessageSpy.mockRestore();
623639
});
624640

641+
it('queues async chatbot dispatches per socket after acknowledging messages', async () => {
642+
const socket = {
643+
handshake: { address: '127.0.0.1' },
644+
data: {},
645+
};
646+
const session: Record<string, any> = {};
647+
const reqBase = {
648+
isSocket: true,
649+
query: { first_name: 'Queue', last_name: 'User' },
650+
session,
651+
headers: { 'user-agent': 'browser' },
652+
socket,
653+
user: {},
654+
} as any as SocketRequest;
655+
const profile = await handler['getOrCreateSession'](
656+
reqBase as any,
657+
webSource,
658+
);
659+
const firstDispatch = createDeferred();
660+
const secondDispatch = createDeferred();
661+
const dispatches = [firstDispatch, secondDispatch];
662+
const dispatchedTexts: string[] = [];
663+
const emitMessageSpy = jest
664+
.spyOn(handler['channelEventBus'], 'emitMessage')
665+
.mockImplementation((event: any) => {
666+
dispatchedTexts.push(event.getText());
667+
668+
return dispatches[dispatchedTexts.length - 1].promise;
669+
});
670+
const sendText = async (text: string) => {
671+
let responseBody: any;
672+
const req = {
673+
...reqBase,
674+
query: {},
675+
body: {
676+
type: 'text',
677+
data: { text },
678+
},
679+
session: {
680+
...session,
681+
web: {
682+
...session.web,
683+
profile,
684+
},
685+
},
686+
} as any as SocketRequest;
687+
688+
await new Promise<void>((resolve, reject) => {
689+
const timeout = setTimeout(() => {
690+
reject(new Error(`Timed out waiting for ${text} response`));
691+
}, 2000);
692+
const res = {
693+
status: (code: number) => {
694+
expect(code).toEqual(200);
695+
696+
return res;
697+
},
698+
json: (payload: any) => {
699+
clearTimeout(timeout);
700+
responseBody = payload;
701+
resolve();
702+
},
703+
} as any as SocketResponse;
704+
705+
handler['handleEvent'](req as any, res, webSource);
706+
});
707+
708+
Object.assign(session, req.session);
709+
710+
return responseBody;
711+
};
712+
const firstResponse = await sendText('First queued message');
713+
const secondResponse = await sendText('Second queued message');
714+
715+
expect(firstResponse.thread_id).toEqual(expect.any(String));
716+
expect(secondResponse.thread_id).toBe(firstResponse.thread_id);
717+
718+
await Promise.resolve();
719+
expect(dispatchedTexts).toEqual(['First queued message']);
720+
expect(emitMessageSpy).toHaveBeenCalledTimes(1);
721+
722+
firstDispatch.resolve();
723+
await firstDispatch.promise;
724+
await Promise.resolve();
725+
await Promise.resolve();
726+
await Promise.resolve();
727+
expect(dispatchedTexts).toEqual([
728+
'First queued message',
729+
'Second queued message',
730+
]);
731+
expect(emitMessageSpy).toHaveBeenCalledTimes(2);
732+
733+
secondDispatch.resolve();
734+
await secondDispatch.promise;
735+
await Promise.resolve();
736+
await Promise.resolve();
737+
await Promise.resolve();
738+
expect(socket.data).not.toHaveProperty('webMessageQueue');
739+
740+
emitMessageSpy.mockRestore();
741+
});
742+
625743
it('rejects chatbot sync before first user message when no thread exists', async () => {
626744
const req = {
627745
isSocket: true,
@@ -805,6 +923,7 @@ describe('WebChannelHandler', () => {
805923
handler['handleEvent'](req as any, res, webSource);
806924
});
807925

926+
await Promise.resolve();
808927
expect(createEventsSpy).toHaveBeenCalledTimes(1);
809928
expect(emitMessageSpy).toHaveBeenCalledWith(expect.anything());
810929
expect(emitStatusSpy).toHaveBeenCalledWith(expect.anything());

packages/api/src/extensions/channels/web/base-web-channel.ts

Lines changed: 44 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,10 @@ import { WebSessionService } from './services/web-session.service';
7474
import { WEB_CHANNEL_NAME } from './settings.schema';
7575
import { Web } from './types';
7676

77+
type WebSocketData = Socket['data'] & {
78+
webMessageQueue?: Promise<void>;
79+
};
80+
7781
/**
7882
* Base handler for the Socket.IO-backed "web" channel.
7983
*
@@ -161,6 +165,30 @@ export default abstract class BaseWebChannelHandler<N extends ChannelName>
161165
return normalizedSourceId.length > 0 ? normalizedSourceId : null;
162166
}
163167

168+
private enqueueMessageDispatch(
169+
req: SocketRequest,
170+
event: MessageInboundEvent,
171+
): void {
172+
const socket = req.socket as Socket & { data?: WebSocketData };
173+
const socketData = (socket.data ??= {});
174+
// Keep chatbot processing ordered per socket without making the client ack
175+
// wait for slow actions, LLM calls, or external integrations.
176+
const previous = socketData.webMessageQueue ?? Promise.resolve();
177+
const next = previous
178+
.catch(() => undefined)
179+
.then(() => this.channelEventBus.emitMessage(event))
180+
.catch((err) => {
181+
this.logger.error('Failed to process web socket message', err);
182+
});
183+
184+
socketData.webMessageQueue = next;
185+
void next.finally(() => {
186+
if (socketData.webMessageQueue === next) {
187+
delete socketData.webMessageQueue;
188+
}
189+
});
190+
}
191+
164192
@OnEvent('hook:websocket:connection', { async: true })
165193
async onWebSocketConnection(client: Socket) {
166194
try {
@@ -547,15 +575,24 @@ export default abstract class BaseWebChannelHandler<N extends ChannelName>
547575
messageEvent.setAuthorForeignId(profile.foreignId);
548576
}
549577
messageEvent.setCreatedAt(new Date());
578+
// Resolve the thread before acknowledging the socket request so the
579+
// client receives a stable thread_id, then dispatch chatbot work later.
580+
const thread = await this.sessionService.resolveThreadForIncoming(
581+
req,
582+
profile.id,
583+
{
584+
explicitThreadId: messageEvent.getThreadId(),
585+
inactivityHours: this.sessionService.resolveInactivityHours(
586+
source.settings,
587+
),
588+
sourceId: source.id,
589+
},
590+
);
591+
messageEvent.setThreadId(thread.id);
592+
messageEvent.setThreadIdOnRaw(thread.id);
550593

551594
this.broadcast(profile, StdEventType.message, messageEvent.getRaw());
552-
await this.channelEventBus.emitMessage(messageEvent);
553-
554-
const resolvedThreadId = messageEvent.getThreadId();
555-
if (resolvedThreadId) {
556-
messageEvent.setThreadIdOnRaw(resolvedThreadId);
557-
if (req.session.web) req.session.web.threadId = resolvedThreadId;
558-
}
595+
this.enqueueMessageDispatch(req, messageEvent);
559596

560597
continue;
561598
}

packages/api/src/extensions/channels/web/services/web-session.service.ts

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,10 @@ export class WebSessionService {
215215
return this.threadService.resolveThread({ subscriberId, explicitThreadId });
216216
}
217217

218+
resolveInactivityHours(settings: unknown): number {
219+
return this.threadService.resolveInactivityHours(settings);
220+
}
221+
218222
/**
219223
* Resolves a thread from the request query/body and writes the result back
220224
* to the session. Used by the subscribe/history flows.
@@ -232,4 +236,35 @@ export class WebSessionService {
232236

233237
return thread;
234238
}
239+
240+
/**
241+
* Resolves the writable thread for an incoming message and persists it on the
242+
* socket session before chatbot processing continues asynchronously.
243+
*/
244+
async resolveThreadForIncoming(
245+
req: SocketRequest,
246+
subscriberId: string,
247+
{
248+
explicitThreadId,
249+
inactivityHours,
250+
sourceId,
251+
}: {
252+
explicitThreadId?: string;
253+
inactivityHours?: number;
254+
sourceId?: string;
255+
},
256+
): Promise<Thread> {
257+
const thread = await this.threadService.resolveThreadForIncoming({
258+
subscriberId,
259+
explicitThreadId,
260+
inactivityHours,
261+
sourceId,
262+
});
263+
264+
if (req.session.web) {
265+
req.session.web.threadId = thread.id;
266+
}
267+
268+
return thread;
269+
}
235270
}

0 commit comments

Comments
 (0)