feat(server): add request id on cluster event (#9998)
This commit is contained in:
@@ -1,5 +1,6 @@
|
|||||||
import { INestApplication } from '@nestjs/common';
|
import { INestApplication } from '@nestjs/common';
|
||||||
import ava, { TestFn } from 'ava';
|
import ava, { TestFn } from 'ava';
|
||||||
|
import { CLS_ID, ClsServiceManager } from 'nestjs-cls';
|
||||||
import Sinon from 'sinon';
|
import Sinon from 'sinon';
|
||||||
|
|
||||||
import { EventBus } from '../../base';
|
import { EventBus } from '../../base';
|
||||||
@@ -11,6 +12,7 @@ const test = ava as TestFn<{
|
|||||||
app1: INestApplication;
|
app1: INestApplication;
|
||||||
app2: INestApplication;
|
app2: INestApplication;
|
||||||
}>;
|
}>;
|
||||||
|
|
||||||
async function createApp() {
|
async function createApp() {
|
||||||
const m = await createTestingModule(
|
const m = await createTestingModule(
|
||||||
{
|
{
|
||||||
@@ -49,14 +51,20 @@ test('should broadcast event to cluster instances', async t => {
|
|||||||
|
|
||||||
// app 2 for broadcasting
|
// app 2 for broadcasting
|
||||||
const eventbus2 = app2.get(EventBus);
|
const eventbus2 = app2.get(EventBus);
|
||||||
eventbus2.broadcast('__test__.event', { count: 0 });
|
const cls = ClsServiceManager.getClsService();
|
||||||
|
cls.run(() => {
|
||||||
|
cls.set(CLS_ID, 'test-request-id');
|
||||||
|
eventbus2.broadcast('__test__.event', { count: 0, requestId: cls.getId() });
|
||||||
|
});
|
||||||
|
|
||||||
// cause the cross instances broadcasting is asynchronization calling
|
// cause the cross instances broadcasting is asynchronization calling
|
||||||
// we should wait for the event's arriving before asserting
|
// we should wait for the event's arriving before asserting
|
||||||
await eventbus1.waitFor('__test__.event');
|
await eventbus1.waitFor('__test__.event');
|
||||||
|
|
||||||
t.true(listener.calledOnceWith({ count: 0 }));
|
t.true(listener.calledOnceWith({ count: 0, requestId: 'test-request-id' }));
|
||||||
t.true(runtimeListener.calledOnceWith({ count: 0 }));
|
t.true(
|
||||||
|
runtimeListener.calledOnceWith({ count: 0, requestId: 'test-request-id' })
|
||||||
|
);
|
||||||
|
|
||||||
off();
|
off();
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import { OnEvent } from '../../base';
|
|||||||
|
|
||||||
declare global {
|
declare global {
|
||||||
interface Events {
|
interface Events {
|
||||||
'__test__.event': { count: number };
|
'__test__.event': { count: number; requestId?: string };
|
||||||
'__test__.event2': { count: number };
|
'__test__.event2': { count: number };
|
||||||
'__test__.throw': { count: number };
|
'__test__.throw': { count: number };
|
||||||
}
|
}
|
||||||
@@ -13,8 +13,13 @@ declare global {
|
|||||||
@Injectable()
|
@Injectable()
|
||||||
export class Listeners {
|
export class Listeners {
|
||||||
@OnEvent('__test__.event')
|
@OnEvent('__test__.event')
|
||||||
onTestEvent({ count }: Events['__test__.event']) {
|
onTestEvent({ count, requestId }: Events['__test__.event']) {
|
||||||
return {
|
return requestId
|
||||||
|
? {
|
||||||
|
count: count + 1,
|
||||||
|
requestId,
|
||||||
|
}
|
||||||
|
: {
|
||||||
count: count + 1,
|
count: count + 1,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,3 +1,5 @@
|
|||||||
|
import { randomUUID } from 'node:crypto';
|
||||||
|
|
||||||
import {
|
import {
|
||||||
applyDecorators,
|
applyDecorators,
|
||||||
Injectable,
|
Injectable,
|
||||||
@@ -15,6 +17,7 @@ import {
|
|||||||
WebSocketGateway,
|
WebSocketGateway,
|
||||||
WebSocketServer,
|
WebSocketServer,
|
||||||
} from '@nestjs/websockets';
|
} from '@nestjs/websockets';
|
||||||
|
import { CLS_ID, ClsService } from 'nestjs-cls';
|
||||||
import type { Server, Socket } from 'socket.io';
|
import type { Server, Socket } from 'socket.io';
|
||||||
|
|
||||||
import { CallMetric } from '../metrics';
|
import { CallMetric } from '../metrics';
|
||||||
@@ -69,7 +72,8 @@ export class EventBus implements OnGatewayConnection, OnApplicationBootstrap {
|
|||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly emitter: EventEmitter2,
|
private readonly emitter: EventEmitter2,
|
||||||
private readonly watcher: EventEmitterReadinessWatcher
|
private readonly watcher: EventEmitterReadinessWatcher,
|
||||||
|
private readonly cls: ClsService
|
||||||
) {}
|
) {}
|
||||||
|
|
||||||
handleConnection(client: Socket) {
|
handleConnection(client: Socket) {
|
||||||
@@ -88,11 +92,15 @@ export class EventBus implements OnGatewayConnection, OnApplicationBootstrap {
|
|||||||
events.forEach(event => {
|
events.forEach(event => {
|
||||||
// Proxy all events received from server(trigger by `server.serverSideEmit`)
|
// Proxy all events received from server(trigger by `server.serverSideEmit`)
|
||||||
// to internal event system
|
// to internal event system
|
||||||
this.server?.on(event, payload => {
|
this.server?.on(event, (payload, requestId?: string) => {
|
||||||
|
this.cls.run(() => {
|
||||||
|
requestId = requestId ?? `server_event-${randomUUID()}`;
|
||||||
|
this.cls.set(CLS_ID, requestId);
|
||||||
this.logger.log(`Server Event: ${event} (Received)`);
|
this.logger.log(`Server Event: ${event} (Received)`);
|
||||||
this.emit(event, payload);
|
this.emit(event, payload);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
});
|
||||||
})
|
})
|
||||||
.catch(() => {
|
.catch(() => {
|
||||||
// startup time promise, never throw at runtime
|
// startup time promise, never throw at runtime
|
||||||
@@ -120,7 +128,7 @@ export class EventBus implements OnGatewayConnection, OnApplicationBootstrap {
|
|||||||
*/
|
*/
|
||||||
broadcast<T extends EventName>(event: T, payload: Events[T]) {
|
broadcast<T extends EventName>(event: T, payload: Events[T]) {
|
||||||
this.logger.log(`Server Event: ${event} (Send)`);
|
this.logger.log(`Server Event: ${event} (Send)`);
|
||||||
this.server?.serverSideEmit(event, payload);
|
this.server?.serverSideEmit(event, payload, this.cls.getId());
|
||||||
}
|
}
|
||||||
|
|
||||||
on<T extends EventName>(
|
on<T extends EventName>(
|
||||||
|
|||||||
Reference in New Issue
Block a user