Пошук уроків, статей та іншого контенту
Реалізуєте власний Server і ClientProxy для підключення NestJS до нестандартного протоколу обміну повідомленнями.
NestJS підтримує стандартні транспортери: TCP, Redis, NATS, MQTT, Kafka та інші. Якщо потрібний протокол відсутній у списку, можна реалізувати власний транспортер.
Власний транспортер складається з двох частин:
Server — приймає повідомлення, знаходить обробник NestJS і повертає результат;
ClientProxy — серіалізує виклики клієнта, надсилає їх серверу та перетворює відповіді на результат send() або emit().
Логічно транспортер має виконувати такий цикл:
клієнт формує пакет із pattern, data та ідентифікатором;
транспорт передає пакет серверу;
сервер знаходить обробник, зареєстрований через @MessagePattern() або @EventPattern();
сервер викликає обробник;
результат повертається клієнту;
клієнт зіставляє відповідь із початковим запитом.
NestJS не вимагає конкретного фізичного протоколу. Це може бути TCP, WebSocket, Unix socket, власна черга або інший механізм обміну.
Для прикладу використаємо простий протокол на основі TCP.
Кожне повідомлення буде JSON-об’єктом, завершеним символом нового рядка:
{"id":"1","type":"request","pattern":"math.add","data":{"a":2,"b":3}}Сервер повертатиме:
{"id":"1","ok":true,"data":5}Для помилки:
{"id":"1","ok":false,"error":"Некоректні аргументи"}Поле id потрібне для зіставлення відповіді з конкретним запитом. Одне TCP-з’єднання може одночасно містити багато запитів, тому покладатися лише на порядок відповідей не можна.
Пакет події не очікує відповіді:
{"id":"2","type":"event","pattern":"notifications.created","data":{"id":"n-1"}}Власний сервер має успадковуватися від Server і реалізовувати CustomTransportStrategy.
Server уже містить внутрішній реєстр обробників NestJS. Після запуску мікросервісу декоратори @MessagePattern() та @EventPattern() реєструють у ньому свої обробники.
Мінімальний контракт стратегії:
listen(callback) — запустити приймання повідомлень;
close() — звільнити ресурси.
У сервері потрібно самостійно:
розібрати вхідний протокол;
знайти обробник за pattern;
викликати його;
дочекатися результату Promise або Observable;
сформувати відповідь.
Клієнт має успадковуватися від ClientProxy.
ClientProxy використовується NestJS для методів:
client.send(pattern, data) — запит із відповіддю;
client.emit(pattern, data) — подія без обов’язкової відповіді.
У власному клієнті потрібно реалізувати:
connect() — встановлення з’єднання;
close() — завершення з’єднання;
publish() — відправлення запиту;
dispatchEvent() — відправлення події.
Метод publish() отримує callback. Його потрібно викликати після отримання відповіді від сервера.
Нижче наведено робочий приклад транспортера на TCP з JSON-повідомленнями, розділеними символом нового рядка.
Встановлення залежностей:
npm install @nestjs/common @nestjs/core @nestjs/microservices rxjs
npm install -D ts-node typescript @types/nodeФайл main.ts:
import * as net from 'node:net';
import {
Controller,
Injectable,
Module,
} from '@nestjs/common';
import { NestFactory } from '@nestjs/core';
import {
ClientProxy,
CustomTransportStrategy,
MessagePattern,
ReadPacket,
Server,
Transport,
WritePacket,
} from '@nestjs/microservices';
import { isObservable, lastValueFrom } from 'rxjs';
type TransportRequest = {
id: string;
type: 'request' | 'event';
pattern: string;
data: unknown;
};
type TransportResponse = {
id: string;
ok: boolean;
data?: unknown;
error?: string;
};
type LineTransportOptions = {
host?: string;
port: number;
};
function serialize(value: unknown): string {
return JSON.stringify(value === undefined ? null : value);
}
async function resolveHandlerResult(value: unknown): Promise<unknown> {
if (isObservable(value)) {
return lastValueFrom(value);
}
return await value;
}
class LineServer
extends Server
implements CustomTransportStrategy
{
private readonly host: string;
private readonly port: number;
private tcpServer?: net.Server;
constructor(options: LineTransportOptions) {
super();
this.host = options.host ?? '127.0.0.1';
this.port = options.port;
}
listen(callback: () => void): void {
this.tcpServer = net.createServer((socket) => {
let buffer = '';
socket.setEncoding('utf8');
socket.on('data', (chunk: string) => {
buffer += chunk;
let newlineIndex = buffer.indexOf('\n');
while (newlineIndex !== -1) {
const line = buffer.slice(0, newlineIndex);
buffer = buffer.slice(newlineIndex + 1);
if (line.trim().length > 0) {
void this.handleLine(socket, line);
}
newlineIndex = buffer.indexOf('\n');
}
});
});
this.tcpServer.listen(this.port, this.host, callback);
}
close(): void {
this.tcpServer?.close();
}
private async handleLine(
socket: net.Socket,
line: string,
): Promise<void> {
let packet: TransportRequest;
try {
packet = JSON.parse(line) as TransportRequest;
} catch {
socket.write(
`${JSON.stringify({
ok: false,
error: 'Некоректний JSON-пакет',
})}\n`,
);
return;
}
const handler = this.getHandlerByPattern(packet.pattern);
if (!handler) {
if (packet.type === 'request') {
const response: TransportResponse = {
id: packet.id,
ok: false,
error: `Обробник для "${packet.pattern}" не знайдено`,
};
socket.write(`${JSON.stringify(response)}\n`);
}
return;
}
try {
const result = await resolveHandlerResult(
handler(packet.data, socket),
);
if (packet.type === 'request') {
const response: TransportResponse = {
id: packet.id,
ok: true,
data: result,
};
socket.write(`${JSON.stringify(response)}\n`);
}
} catch (error) {
if (packet.type !== 'request') {
return;
}
const message =
error instanceof Error
? error.message
: 'Невідома помилка обробника';
const response: TransportResponse = {
id: packet.id,
ok: false,
error: message,
};
socket.write(`${JSON.stringify(response)}\n`);
}
}
}
class LineClient extends ClientProxy {
private readonly host: string;
private readonly port: number;
private socket?: net.Socket;
private buffer = '';
private sequence = 0;
private readonly pending = new Map<
string,
(packet: WritePacket) => void
>();
constructor(options: LineTransportOptions) {
super();
this.host = options.host ?? '127.0.0.1';
this.port = options.port;
}
connect(): Promise<net.Socket> {
if (this.socket && !this.socket.destroyed) {
return Promise.resolve(this.socket);
}
return new Promise((resolve, reject) => {
const socket = net.createConnection({
host: this.host,
port: this.port,
});
socket.setEncoding('utf8');
const onConnect = () => {
socket.off('error', onInitialError);
resolve(socket);
};
const onInitialError = (error: Error) => {
socket.off('connect', onConnect);
reject(error);
};
socket.once('connect', onConnect);
socket.once('error', onInitialError);
socket.on('data', (chunk: string) => {
this.handleData(chunk);
});
socket.on('error', (error: Error) => {
for (const callback of this.pending.values()) {
callback({ err: error.message });
}
this.pending.clear();
});
socket.on('close', () => {
this.socket = undefined;
});
this.socket = socket;
});
}
close(): void {
this.socket?.end();
this.socket = undefined;
}
protected publish(
packet: ReadPacket,
callback: (packet: WritePacket) => void,
): () => void {
const id = String(++this.sequence);
void this.connect()
.then((socket) => {
this.pending.set(id, callback);
const request: TransportRequest = {
id,
type: 'request',
pattern: String(packet.pattern),
data: packet.data,
};
socket.write(`${JSON.stringify(request)}\n`);
})
.catch((error: Error) => {
callback({ err: error.message });
});
return () => {
this.pending.delete(id);
};
}
protected async dispatchEvent(
packet: ReadPacket,
): Promise<void> {
const socket = await this.connect();
const event: TransportRequest = {
id: String(++this.sequence),
type: 'event',
pattern: String(packet.pattern),
data: packet.data,
};
socket.write(`${JSON.stringify(event)}\n`);
}
private handleData(chunk: string): void {
this.buffer += chunk;
let newlineIndex = this.buffer.indexOf('\n');
while (newlineIndex !== -1) {
const line = this.buffer.slice(0, newlineIndex);
this.buffer = this.buffer.slice(newlineIndex + 1);
if (line.trim().length > 0) {
this.handleResponse(line);
}
newlineIndex = this.buffer.indexOf('\n');
}
}
private handleResponse(line: string): void {
let response: TransportResponse;
try {
response = JSON.parse(line) as TransportResponse;
} catch {
return;
}
const callback = this.pending.get(response.id);
if (!callback) {
return;
}
this.pending.delete(response.id);
if (response.ok) {
callback({
response: response.data,
isDisposed: true,
});
} else {
callback({
err: response.error ?? 'Помилка віддаленого обробника',
});
}
}
}
@Controller()
class MathController {
@MessagePattern('math.add')
add(data: { a: number; b: number }): number {
if (
typeof data?.a !== 'number' ||
typeof data?.b !== 'number'
) {
throw new Error('Поля "a" та "b" мають бути числами');
}
return data.a + data.b;
}
}
@Module({
controllers: [MathController],
})
class AppModule {}
async function startServer(): Promise<void> {
const app = await NestFactory.createMicroservice(AppModule, {
strategy: new LineServer({
host: '127.0.0.1',
port: 4000,
}),
});
await app.listen();
console.log('Власний транспортер слухає 127.0.0.1:4000');
}
async function startClient(): Promise<void> {
const client = new LineClient({
host: '127.0.0.1',
port: 4000,
});
const result = await client
.send<number>('math.add', { a: 2, b: 3 })
.toPromise();
console.log('Результат:', result);
client.close();
}
const mode = process.argv[2];
if (mode === 'server') {
void startServer();
} else if (mode === 'client') {
void startClient();
} else {
console.log('Використання:');
console.log(' npx ts-node main.ts server');
console.log(' npx ts-node main.ts client');
}Запустіть сервер в одному процесі:
npx ts-node main.ts serverВ іншому процесі запустіть клієнта:
npx ts-node main.ts clientРезультат клієнта:
Результат: 5Після створення мікросервісу NestJS реєструє метод add() під шаблоном math.add.
У сервері виконується:
const handler = this.getHandlerByPattern(packet.pattern);Якщо клієнт передав math.add, сервер отримає функцію методу контролера і викличе її:
handler(packet.data, socket)Другий аргумент можна використати як контекст транспорту: TCP-сокет, ідентифікатор з’єднання або власний об’єкт метаданих.
Обробник NestJS може повертати:
звичайне значення;
Promise;
Observable.
Тому сервер не повинен просто серіалізувати результат виклику. Спочатку потрібно дочекатися асинхронного результату:
async function resolveHandlerResult(value: unknown): Promise<unknown> {
if (isObservable(value)) {
return lastValueFrom(value);
}
return await value;
}Для Observable важливо враховувати його завершення. Якщо Observable не завершується, lastValueFrom() чекатиме нескінченно. Для потокових відповідей потрібен інший протокол, який передає кілька повідомлень, а не одне фінальне значення.
У попередньому прикладі клієнт створюється напряму. У застосунку NestJS його можна зареєструвати як provider.
import { Module } from '@nestjs/common';
@Module({
providers: [
{
provide: 'LINE_TRANSPORT',
useFactory: () =>
new LineClient({
host: '127.0.0.1',
port: 4000,
}),
},
],
exports: ['LINE_TRANSPORT'],
})
export class TransportModule {}Після цього клієнт можна інжектувати в сервіс:
import { Inject, Injectable } from '@nestjs/common';
import { ClientProxy } from '@nestjs/microservices';
import { firstValueFrom } from 'rxjs';
@Injectable()
export class CalculatorService {
constructor(
@Inject('LINE_TRANSPORT')
private readonly client: ClientProxy,
) {}
async add(a: number, b: number): Promise<number> {
return firstValueFrom(
this.client.send<number>('math.add', { a, b }),
);
}
}send() повертає Observable, тому для отримання одного результату зручно використовувати firstValueFrom().
У реальному застосунку provider також повинен закривати клієнт під час завершення модуля. Для цього клас-клієнт можна доповнити методом життєвого циклу або явно викликати close() у відповідному сервісі.
Запит використовує модель request-response:
const result$ = client.send('math.add', { a: 2, b: 3 });Для нього клієнт:
створює id;
зберігає callback у pending;
надсилає пакет;
чекає пакет із таким самим id;
передає результат у callback.
Подія не має обов’язкової відповіді:
client.emit('notifications.created', {
id: 'notification-1',
});У власному протоколі поле type допомагає серверу відрізнити ці випадки. Для події сервер викликає обробник, але не відправляє результат клієнту.
Окремо слід визначити поведінку подій у разі помилки. У прикладі помилка події ігнорується на рівні відповіді, оскільки клієнт не очікує callback. У production-транспортері можна додати:
підтвердження доставки;
повторну доставку;
окремий канал помилок;
dead-letter чергу.
TCP не зберігає межі повідомлень. Один виклик data може містити:
лише частину JSON;
кілька JSON-повідомлень;
JSON-повідомлення разом із частиною наступного.
Тому код використовує буфер:
let buffer = '';
socket.on('data', (chunk: string) => {
buffer += chunk;
let newlineIndex = buffer.indexOf('\n');
while (newlineIndex !== -1) {
const line = buffer.slice(0, newlineIndex);
buffer = buffer.slice(newlineIndex + 1);
// Обробка одного повного повідомлення
newlineIndex = buffer.indexOf('\n');
}
});Розділення повідомлень символом нового рядка підходить для простого протоколу, але має обмеження. Наприклад, потрібно гарантувати, що серіалізований JSON не містить необробленого символу нового рядка поза рядковими значеннями. JSON.stringify() коректно екранує такі символи.
Для складніших протоколів можна використовувати:
префікс із довжиною повідомлення;
бінарний заголовок;
окремий кадр протоколу;
готовий фреймінг, якщо його визначає зовнішня система.
Власний транспортер відповідає за перетворення помилок у формат протоколу.
На сервері помилка обробника перетворюється на:
{
id: '1',
ok: false,
error: 'Текст помилки'
}На клієнті вона передається в callback як err:
callback({
err: response.error ?? 'Помилка віддаленого обробника',
});Після цього NestJS завершує Observable, повернений send(), помилкою.
Не варто передавати клієнту безпосередньо весь об’єкт помилки. Він може містити:
циклічні посилання;
конфіденційні дані;
внутрішній stack trace;
значення, які неможливо серіалізувати.
Краще визначити стабільний формат помилки, наприклад:
type TransportError = {
code: string;
message: string;
details?: unknown;
};Метод connect() має бути ідемпотентним: повторний виклик не повинен створювати нове з’єднання, якщо чинне вже існує.
У прикладі перевіряється стан сокета:
if (this.socket && !this.socket.destroyed) {
return Promise.resolve(this.socket);
}У складнішій реалізації бажано також зберігати Promise поточного підключення. Інакше кілька паралельних викликів send() можуть одночасно почати кілька спроб підключення.
Потрібно також визначити поведінку під час розриву з’єднання:
завершити всі очікувані запити помилкою;
видалити їх із pending;
дозволити наступному запиту виконати повторне підключення;
не викликати callback для одного запиту більше одного разу.
Не можна вважати, що одна подія data містить одне повне повідомлення. Завжди потрібен буфер і механізм фреймінгу.
Якщо клієнт зіставляє відповіді лише за порядком, паралельні запити можуть отримати чужі результати. Кожен запит повинен мати унікальний id.
Callback із publish() потрібно викликати і при помилці, і при розриві з’єднання. Інакше Observable може залишитися активним назавжди.
Обробник NestJS не обов’язково повертає Promise. Якщо Observable передати в JSON.stringify() без обробки, клієнт не отримає очікуване значення.
Клієнт може передати помилковий pattern, а сервер може не мати відповідного обробника. Такий випадок повинен завершувати запит структурованою помилкою.
Не кожне значення можна передати через JSON. Не підтримуються, зокрема, функції, BigInt, циклічні структури та деякі спеціальні об’єкти. Контракт транспорту повинен визначати допустимі дані.
Сервер повинен закривати TCP-сервер, а клієнт — сокет і незавершені запити. Інакше процес може не завершуватися або залишати очікувані операції в пам’яті.
Приклад призначений для локального обміну повідомленнями. Для мережевого використання потрібно окремо визначити автентифікацію, шифрування, обмеження розміру пакета та захист від некоректних даних.
Власний транспортер NestJS реалізується через Server і ClientProxy.
Сервер повинен реалізувати listen() та close(), а також розбір протоколу й виклик зареєстрованих обробників.
Клієнт повинен реалізувати connect(), close(), publish() і dispatchEvent().
send() використовує request-response, а emit() — подієву модель.
Для паралельних запитів потрібні унікальні ідентифікатори та реєстр очікуваних відповідей.
TCP потребує власного фреймінгу повідомлень.
Власний транспортер має явно визначати формат помилок, поведінку під час розриву з’єднання та правила серіалізації даних.