Skip to content

Commit

Permalink
fix(core): Fix parsing of push messages (no-changelog) (#13136)
Browse files Browse the repository at this point in the history
  • Loading branch information
tomi authored Feb 10, 2025
1 parent 298a7b0 commit 67b951e
Show file tree
Hide file tree
Showing 5 changed files with 55 additions and 6 deletions.
2 changes: 2 additions & 0 deletions packages/@n8n/api-types/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ export type * from './user';
export type * from './api-keys';

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 { BannerName } from './schemas/bannerName.schema';
Expand Down
11 changes: 11 additions & 0 deletions packages/@n8n/api-types/src/push/heartbeat.ts
Original file line number Diff line number Diff line change
@@ -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',
});
25 changes: 23 additions & 2 deletions packages/cli/src/push/__tests__/websocket.push.test.ts
Original file line number Diff line number Diff line change
@@ -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 { EventEmitter } from 'events';
import { Logger } from 'n8n-core';
Expand Down Expand Up @@ -107,7 +107,8 @@ describe('WebSocketPush', () => {
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();
webSocketPush.on('message', mockOnMessageReceived);
webSocketPush.add(pushRef1, userId, mockWebSocket1);
Expand All @@ -118,10 +119,30 @@ describe('WebSocketPush', () => {

mockWebSocket1.emit('message', buffer);

// Flush the event loop
await new Promise(process.nextTick);

expect(mockOnMessageReceived).toHaveBeenCalledWith({
msg: data,
pushRef: pushRef1,
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();
});
});
19 changes: 17 additions & 2 deletions packages/cli/src/push/websocket.push.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { heartbeatMessageSchema } from '@n8n/api-types';
import { Service } from '@n8n/di';
import { ApplicationError } from 'n8n-workflow';
import type WebSocket from 'ws';
Expand All @@ -18,11 +19,19 @@ export class WebSocketPush extends AbstractPush<WebSocket> {

super.add(pushRef, userId, connection);

const onMessage = (data: WebSocket.RawData) => {
const onMessage = async (data: WebSocket.RawData) => {
try {
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) {
this.errorReporter.error(
new ApplicationError('Error parsing push message', {
Expand Down Expand Up @@ -67,4 +76,10 @@ export class WebSocketPush extends AbstractPush<WebSocket> {
connection.isAlive = false;
connection.ping();
}

private async isClientHeartbeat(msg: unknown) {
const result = await heartbeatMessageSchema.safeParseAsync(msg);

return result.success;
}
}
4 changes: 2 additions & 2 deletions packages/editor-ui/src/push-connection/useWebSocketClient.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { useHeartbeat } from '@/push-connection/useHeartbeat';
import { useReconnectTimer } from '@/push-connection/useReconnectTimer';
import { ref } from 'vue';

import { createHeartbeatMessage } from '@n8n/api-types';
export type UseWebSocketClientOptions<T> = {
url: string;
onMessage: (data: T) => void;
Expand Down Expand Up @@ -31,7 +31,7 @@ export const useWebSocketClient = <T>(options: UseWebSocketClientOptions<T>) =>
const { startHeartbeat, stopHeartbeat } = useHeartbeat({
interval: 30_000,
onHeartbeat: () => {
socket.value?.send(JSON.stringify({ type: 'heartbeat' }));
socket.value?.send(JSON.stringify(createHeartbeatMessage()));
},
});

Expand Down

0 comments on commit 67b951e

Please # to comment.