mirror of
https://github.com/n8n-io/n8n.git
synced 2025-03-05 20:50:17 -08:00
fix(core): Fix parsing of push messages (no-changelog) (#13136)
This commit is contained in:
parent
298a7b0038
commit
67b951ee07
|
@ -7,6 +7,8 @@ export type * from './user';
|
||||||
export type * from './api-keys';
|
export type * from './api-keys';
|
||||||
|
|
||||||
export type { Collaborator } from './push/collaboration';
|
export type { Collaborator } from './push/collaboration';
|
||||||
|
export type { HeartbeatMessage } from './push/heartbeat';
|
||||||
|
export { createHeartbeatMessage, heartbeatMessageSchema } from './push/heartbeat';
|
||||||
export type { SendWorkerStatusMessage } from './push/worker';
|
export type { SendWorkerStatusMessage } from './push/worker';
|
||||||
|
|
||||||
export type { BannerName } from './schemas/bannerName.schema';
|
export type { BannerName } from './schemas/bannerName.schema';
|
||||||
|
|
11
packages/@n8n/api-types/src/push/heartbeat.ts
Normal file
11
packages/@n8n/api-types/src/push/heartbeat.ts
Normal file
|
@ -0,0 +1,11 @@
|
||||||
|
import { z } from 'zod';
|
||||||
|
|
||||||
|
export const heartbeatMessageSchema = z.object({
|
||||||
|
type: z.literal('heartbeat'),
|
||||||
|
});
|
||||||
|
|
||||||
|
export type HeartbeatMessage = z.infer<typeof heartbeatMessageSchema>;
|
||||||
|
|
||||||
|
export const createHeartbeatMessage = (): HeartbeatMessage => ({
|
||||||
|
type: 'heartbeat',
|
||||||
|
});
|
|
@ -1,4 +1,4 @@
|
||||||
import type { PushMessage } from '@n8n/api-types';
|
import { createHeartbeatMessage, type PushMessage } from '@n8n/api-types';
|
||||||
import { Container } from '@n8n/di';
|
import { Container } from '@n8n/di';
|
||||||
import { EventEmitter } from 'events';
|
import { EventEmitter } from 'events';
|
||||||
import { Logger } from 'n8n-core';
|
import { Logger } from 'n8n-core';
|
||||||
|
@ -107,7 +107,8 @@ describe('WebSocketPush', () => {
|
||||||
expect(mockWebSocket2.send).toHaveBeenCalledWith(expectedMsg);
|
expect(mockWebSocket2.send).toHaveBeenCalledWith(expectedMsg);
|
||||||
});
|
});
|
||||||
|
|
||||||
it('emits message event when connection receives data', () => {
|
it('emits message event when connection receives data', async () => {
|
||||||
|
jest.useRealTimers();
|
||||||
const mockOnMessageReceived = jest.fn();
|
const mockOnMessageReceived = jest.fn();
|
||||||
webSocketPush.on('message', mockOnMessageReceived);
|
webSocketPush.on('message', mockOnMessageReceived);
|
||||||
webSocketPush.add(pushRef1, userId, mockWebSocket1);
|
webSocketPush.add(pushRef1, userId, mockWebSocket1);
|
||||||
|
@ -118,10 +119,30 @@ describe('WebSocketPush', () => {
|
||||||
|
|
||||||
mockWebSocket1.emit('message', buffer);
|
mockWebSocket1.emit('message', buffer);
|
||||||
|
|
||||||
|
// Flush the event loop
|
||||||
|
await new Promise(process.nextTick);
|
||||||
|
|
||||||
expect(mockOnMessageReceived).toHaveBeenCalledWith({
|
expect(mockOnMessageReceived).toHaveBeenCalledWith({
|
||||||
msg: data,
|
msg: data,
|
||||||
pushRef: pushRef1,
|
pushRef: pushRef1,
|
||||||
userId,
|
userId,
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("emits doesn' emit message for client heartbeat", async () => {
|
||||||
|
const mockOnMessageReceived = jest.fn();
|
||||||
|
webSocketPush.on('message', mockOnMessageReceived);
|
||||||
|
webSocketPush.add(pushRef1, userId, mockWebSocket1);
|
||||||
|
webSocketPush.add(pushRef2, userId, mockWebSocket2);
|
||||||
|
|
||||||
|
const data = createHeartbeatMessage();
|
||||||
|
const buffer = Buffer.from(JSON.stringify(data));
|
||||||
|
|
||||||
|
mockWebSocket1.emit('message', buffer);
|
||||||
|
|
||||||
|
// Flush the event loop
|
||||||
|
await new Promise(process.nextTick);
|
||||||
|
|
||||||
|
expect(mockOnMessageReceived).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
|
@ -1,3 +1,4 @@
|
||||||
|
import { heartbeatMessageSchema } from '@n8n/api-types';
|
||||||
import { Service } from '@n8n/di';
|
import { Service } from '@n8n/di';
|
||||||
import { ApplicationError } from 'n8n-workflow';
|
import { ApplicationError } from 'n8n-workflow';
|
||||||
import type WebSocket from 'ws';
|
import type WebSocket from 'ws';
|
||||||
|
@ -18,11 +19,19 @@ export class WebSocketPush extends AbstractPush<WebSocket> {
|
||||||
|
|
||||||
super.add(pushRef, userId, connection);
|
super.add(pushRef, userId, connection);
|
||||||
|
|
||||||
const onMessage = (data: WebSocket.RawData) => {
|
const onMessage = async (data: WebSocket.RawData) => {
|
||||||
try {
|
try {
|
||||||
const buffer = Array.isArray(data) ? Buffer.concat(data) : Buffer.from(data);
|
const buffer = Array.isArray(data) ? Buffer.concat(data) : Buffer.from(data);
|
||||||
|
const msg: unknown = JSON.parse(buffer.toString('utf8'));
|
||||||
|
|
||||||
this.onMessageReceived(pushRef, JSON.parse(buffer.toString('utf8')));
|
// Client sends application level heartbeat messages to react
|
||||||
|
// to connection issues. This is in addition to the protocol
|
||||||
|
// level ping/pong mechanism used by the server.
|
||||||
|
if (await this.isClientHeartbeat(msg)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.onMessageReceived(pushRef, msg);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
this.errorReporter.error(
|
this.errorReporter.error(
|
||||||
new ApplicationError('Error parsing push message', {
|
new ApplicationError('Error parsing push message', {
|
||||||
|
@ -67,4 +76,10 @@ export class WebSocketPush extends AbstractPush<WebSocket> {
|
||||||
connection.isAlive = false;
|
connection.isAlive = false;
|
||||||
connection.ping();
|
connection.ping();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async isClientHeartbeat(msg: unknown) {
|
||||||
|
const result = await heartbeatMessageSchema.safeParseAsync(msg);
|
||||||
|
|
||||||
|
return result.success;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -1,7 +1,7 @@
|
||||||
import { useHeartbeat } from '@/push-connection/useHeartbeat';
|
import { useHeartbeat } from '@/push-connection/useHeartbeat';
|
||||||
import { useReconnectTimer } from '@/push-connection/useReconnectTimer';
|
import { useReconnectTimer } from '@/push-connection/useReconnectTimer';
|
||||||
import { ref } from 'vue';
|
import { ref } from 'vue';
|
||||||
|
import { createHeartbeatMessage } from '@n8n/api-types';
|
||||||
export type UseWebSocketClientOptions<T> = {
|
export type UseWebSocketClientOptions<T> = {
|
||||||
url: string;
|
url: string;
|
||||||
onMessage: (data: T) => void;
|
onMessage: (data: T) => void;
|
||||||
|
@ -31,7 +31,7 @@ export const useWebSocketClient = <T>(options: UseWebSocketClientOptions<T>) =>
|
||||||
const { startHeartbeat, stopHeartbeat } = useHeartbeat({
|
const { startHeartbeat, stopHeartbeat } = useHeartbeat({
|
||||||
interval: 30_000,
|
interval: 30_000,
|
||||||
onHeartbeat: () => {
|
onHeartbeat: () => {
|
||||||
socket.value?.send(JSON.stringify({ type: 'heartbeat' }));
|
socket.value?.send(JSON.stringify(createHeartbeatMessage()));
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
Loading…
Reference in a new issue