Networking con Node.js
Node.js nasce come runtime per applicazioni di rete: il suo modello a event loop con I/O non bloccante è stato progettato proprio per gestire migliaia di connessioni contemporanee con un singolo thread. Oltre ai framework HTTP più noti, la libreria standard offre moduli di basso livello per TCP, UDP, DNS e TLS che permettono di costruire protocolli personalizzati, proxy e strumenti di diagnostica. In questo articolo esploreremo questi moduli con esempi completi, ponendo l'accento su framing, backpressure, timeout e chiusura ordinata delle connessioni.
I moduli di rete della libreria standard
node:net: server e client TCP, oltre ai socket di dominio Unix.node:dgram: socket UDP.node:dnsenode:dns/promises: risoluzione dei nomi, sia tramite il resolver di sistema sia tramite query DNS dirette.node:tls: connessioni cifrate costruite sopranet.fetch: client HTTP globale basato su undici, disponibile senza import.
Tutti gli esempi usano i moduli ES e richiedono una versione LTS recente di Node.js.
Un server TCP minimale
Con net.createServer() si ottiene un server la cui callback viene invocata per ogni nuova connessione. Ogni connessione è un oggetto net.Socket, ovvero uno stream duplex: si legge ascoltando l'evento data e si scrive con write().
import net from 'node:net';
const server = net.createServer((socket) => {
const peer = `${socket.remoteAddress}:${socket.remotePort}`;
console.log(`Nuova connessione da ${peer}`);
// I dati arrivano come Buffer; con setEncoding si ottengono stringhe
socket.setEncoding('utf8');
socket.on('data', (chunk) => {
socket.write(`echo: ${chunk}`);
});
socket.on('end', () => {
console.log(`${peer} ha chiuso la connessione`);
});
// Senza un listener per 'error' un ECONNRESET farebbe terminare il processo
socket.on('error', (err) => {
console.error(`Errore su ${peer}: ${err.message}`);
});
});
server.listen(9000, '0.0.0.0', () => {
console.log('Server TCP in ascolto sulla porta 9000');
});
Il listener per l'evento error è obbligatorio nella pratica: in Node.js un EventEmitter che emette error senza listener solleva un'eccezione non gestita, e un client che si disconnette bruscamente basterebbe ad abbattere l'intero server.
Il problema del framing
Il server precedente ha un difetto sottile: tratta ogni evento data come un messaggio completo. TCP però è un flusso di byte senza confini: un messaggio può arrivare spezzato in più chunk, oppure più messaggi possono arrivare in un unico chunk. Serve un livello di framing. Con il prefisso di lunghezza, ogni messaggio è preceduto da un intero a 32 bit che ne indica la dimensione.
import { Transform } from 'node:stream';
const HEADER_SIZE = 4;
const MAX_FRAME_SIZE = 16 * 1024 * 1024;
export function encodeFrame(payload) {
const body = Buffer.isBuffer(payload) ? payload : Buffer.from(payload, 'utf8');
const header = Buffer.alloc(HEADER_SIZE);
// Lunghezza in network byte order (big-endian)
header.writeUInt32BE(body.length, 0);
return Buffer.concat([header, body]);
}
export class FrameDecoder extends Transform {
#buffer = Buffer.alloc(0);
constructor() {
super({ readableObjectMode: true });
}
_transform(chunk, encoding, callback) {
this.#buffer = Buffer.concat([this.#buffer, chunk]);
// Si estraggono tutti i frame completi presenti nel buffer
while (this.#buffer.length >= HEADER_SIZE) {
const length = this.#buffer.readUInt32BE(0);
if (length > MAX_FRAME_SIZE) {
callback(new Error(`Frame troppo grande: ${length} byte`));
return;
}
if (this.#buffer.length < HEADER_SIZE + length) {
// Frame incompleto: si attende il chunk successivo
break;
}
const frame = this.#buffer.subarray(HEADER_SIZE, HEADER_SIZE + length);
this.#buffer = this.#buffer.subarray(HEADER_SIZE + length);
this.push(frame);
}
callback();
}
}
Implementare il decoder come stream Transform ha un vantaggio importante: può essere collegato al socket con pipe() o pipeline(), ereditando automaticamente la gestione della backpressure e la propagazione degli errori.
Un server JSON con framing e timeout di inattività
Usando il codec appena scritto costruiamo un server che riceve comandi JSON e risponde con oggetti JSON. Ogni connessione inattiva per troppo tempo viene chiusa con setTimeout(), evitando che client abbandonati consumino risorse.
import net from 'node:net';
import { pipeline } from 'node:stream/promises';
import { encodeFrame, FrameDecoder } from './frame-codec.js';
const IDLE_TIMEOUT_MS = 30_000;
const connections = new Set();
const handlers = {
ping: () => ({ pong: Date.now() }),
sum: ({ values }) => ({ result: values.reduce((acc, value) => acc + value, 0) }),
uptime: () => ({ seconds: Math.round(process.uptime()) }),
};
function handleMessage(socket, frame) {
let request;
try {
request = JSON.parse(frame.toString('utf8'));
} catch {
socket.write(encodeFrame(JSON.stringify({ error: 'JSON non valido' })));
return;
}
const handler = handlers[request.command];
const response = handler
? { id: request.id, ...handler(request.params ?? {}) }
: { id: request.id, error: `Comando sconosciuto: ${request.command}` };
socket.write(encodeFrame(JSON.stringify(response)));
}
const server = net.createServer(async (socket) => {
connections.add(socket);
// Chiusura automatica delle connessioni inattive
socket.setTimeout(IDLE_TIMEOUT_MS);
socket.on('timeout', () => {
socket.end(encodeFrame(JSON.stringify({ error: 'Timeout di inattività' })));
});
socket.setNoDelay(true);
const decoder = new FrameDecoder();
decoder.on('data', (frame) => handleMessage(socket, frame));
try {
await pipeline(socket, decoder);
} catch (err) {
if (err.code !== 'ECONNRESET') {
console.error(`Connessione terminata con errore: ${err.message}`);
}
socket.destroy();
} finally {
connections.delete(socket);
}
});
server.listen(9000, () => console.log('Server JSON in ascolto sulla porta 9000'));
// Chiusura ordinata: si smette di accettare connessioni e si chiudono quelle aperte
function shutdown() {
console.log('Arresto in corso...');
server.close(() => process.exit(0));
for (const socket of connections) {
socket.end();
}
// Se qualche client non chiude, si forza l'uscita dopo 5 secondi
setTimeout(() => process.exit(1), 5000).unref();
}
process.on('SIGINT', shutdown);
process.on('SIGTERM', shutdown);
server.close() impedisce nuove connessioni ma attende che quelle esistenti terminino: per questo il server tiene traccia dei socket aperti in un Set e li chiude esplicitamente. Il timer di sicurezza con unref() non impedisce al processo di terminare se tutto si chiude prima.
Il client: richieste con correlazione e timeout
Su una connessione persistente un client può inviare più richieste senza attendere le risposte precedenti. Per associare ogni risposta alla richiesta corretta si usa un identificativo, e ogni richiesta pendente è una Promise con il proprio timeout.
import net from 'node:net';
import { once } from 'node:events';
import { encodeFrame, FrameDecoder } from './frame-codec.js';
export class JsonClient {
#socket = null;
#nextId = 1;
#pending = new Map();
constructor(host, port, { connectTimeoutMs = 3000, requestTimeoutMs = 5000 } = {}) {
this.host = host;
this.port = port;
this.connectTimeoutMs = connectTimeoutMs;
this.requestTimeoutMs = requestTimeoutMs;
}
async connect() {
const socket = net.connect({ host: this.host, port: this.port });
try {
// once rifiuta la promise se viene emesso 'error' prima di 'connect' o se scade il timeout
await once(socket, 'connect', { signal: AbortSignal.timeout(this.connectTimeoutMs) });
} catch (err) {
socket.destroy();
throw err;
}
const decoder = new FrameDecoder();
socket.pipe(decoder);
decoder.on('data', (frame) => this.#resolve(frame));
socket.on('close', () => this.#rejectAll(new Error('Connessione chiusa')));
socket.on('error', (err) => this.#rejectAll(err));
this.#socket = socket;
}
request(command, params = {}) {
if (!this.#socket) {
return Promise.reject(new Error('Client non connesso'));
}
const id = this.#nextId++;
return new Promise((resolve, reject) => {
// Ogni richiesta ha un proprio timeout indipendente
const timer = setTimeout(() => {
this.#pending.delete(id);
reject(new Error(`Timeout della richiesta ${id} (${command})`));
}, this.requestTimeoutMs);
this.#pending.set(id, { resolve, reject, timer });
this.#socket.write(encodeFrame(JSON.stringify({ id, command, params })));
});
}
close() {
this.#socket?.end();
}
#resolve(frame) {
const message = JSON.parse(frame.toString('utf8'));
const entry = this.#pending.get(message.id);
if (!entry) {
return;
}
clearTimeout(entry.timer);
this.#pending.delete(message.id);
if (message.error) {
entry.reject(new Error(message.error));
} else {
entry.resolve(message);
}
}
#rejectAll(err) {
for (const { reject, timer } of this.#pending.values()) {
clearTimeout(timer);
reject(err);
}
this.#pending.clear();
}
}
// Esempio d'uso: tre richieste inviate in parallelo sulla stessa connessione
const client = new JsonClient('127.0.0.1', 9000);
await client.connect();
const [ping, sum, uptime] = await Promise.all([
client.request('ping'),
client.request('sum', { values: [1, 2, 3, 4] }),
client.request('uptime'),
]);
console.log(ping, sum, uptime);
client.close();
Backpressure: quando il destinatario è più lento
socket.write() restituisce false quando il buffer interno di scrittura supera la soglia highWaterMark. Ignorare questo segnale e continuare a scrivere significa accumulare dati in memoria senza limite. Il modo corretto è sospendere la scrittura fino all'evento drain:
import { once } from 'node:events';
async function sendMany(socket, frames) {
for (const frame of frames) {
const canContinue = socket.write(frame);
// Il buffer è pieno: si attende che venga svuotato
if (!canContinue) {
await once(socket, 'drain');
}
}
}
Quando si collega uno stream a un altro con pipeline(), questa logica è già implementata: è uno dei motivi per preferire gli stream alle scritture manuali. Un proxy TCP, ad esempio, si scrive in poche righe proprio grazie a pipeline():
import net from 'node:net';
import { pipeline } from 'node:stream/promises';
const TARGET_HOST = '127.0.0.1';
const TARGET_PORT = 5432;
const proxy = net.createServer(async (client) => {
const upstream = net.connect({ host: TARGET_HOST, port: TARGET_PORT });
try {
// Copia bidirezionale con gestione automatica di backpressure ed errori
await Promise.all([
pipeline(client, upstream),
pipeline(upstream, client),
]);
} catch (err) {
console.error(`Proxy: ${err.message}`);
} finally {
client.destroy();
upstream.destroy();
}
});
proxy.listen(15432, () => console.log(`Proxy verso ${TARGET_HOST}:${TARGET_PORT} sulla porta 15432`));
UDP con dgram
Il modulo dgram espone i socket UDP. Ogni messaggio è un datagramma indipendente, consegnato con l'indirizzo del mittente. Ecco un semplice servizio di discovery: i nodi annunciano la propria presenza in broadcast e ascoltano gli annunci degli altri.
import dgram from 'node:dgram';
import os from 'node:os';
const DISCOVERY_PORT = 41234;
const BROADCAST_ADDRESS = '255.255.255.255';
const nodeId = `${os.hostname()}-${process.pid}`;
const peers = new Map();
const socket = dgram.createSocket({ type: 'udp4', reuseAddr: true });
socket.on('message', (message, remote) => {
let announcement;
try {
announcement = JSON.parse(message.toString('utf8'));
} catch {
return;
}
// Si ignorano i propri annunci
if (announcement.nodeId === nodeId) {
return;
}
const isNew = !peers.has(announcement.nodeId);
peers.set(announcement.nodeId, { address: remote.address, lastSeen: Date.now() });
if (isNew) {
console.log(`Scoperto nodo ${announcement.nodeId} su ${remote.address}`);
}
});
socket.bind(DISCOVERY_PORT, () => {
socket.setBroadcast(true);
// Annuncio periodico della propria presenza
setInterval(() => {
const payload = Buffer.from(JSON.stringify({ nodeId, timestamp: Date.now() }));
socket.send(payload, DISCOVERY_PORT, BROADCAST_ADDRESS);
}, 2000);
// Rimozione dei nodi che non si annunciano da più di 10 secondi
setInterval(() => {
const threshold = Date.now() - 10_000;
for (const [id, peer] of peers) {
if (peer.lastSeen < threshold) {
peers.delete(id);
console.log(`Nodo ${id} non più raggiungibile`);
}
}
}, 5000);
});
Il pattern dell'annuncio periodico con scadenza compensa la natura inaffidabile di UDP: un singolo datagramma perso non ha conseguenze, perché ne arriverà un altro due secondi dopo.
Risoluzione DNS
Il modulo dns offre due famiglie di funzioni con comportamenti diversi, ed è importante conoscerne la differenza:
dns.lookup()usa il resolver del sistema operativo (getaddrinfo), rispetta/etc/hostse viene eseguito nel thread pool di libuv. È la funzione usata internamente danet.connect().dns.resolve*()interroga direttamente i server DNS tramite c-ares, in modo realmente asincrono, e permette di richiedere record di qualsiasi tipo.
import dns from 'node:dns/promises';
async function inspectDomain(domain) {
const resolver = new dns.Resolver({ timeout: 2000, tries: 2 });
// Server DNS specifici, indipendenti dalla configurazione di sistema
resolver.setServers(['1.1.1.1', '8.8.8.8']);
const settled = await Promise.allSettled([
resolver.resolve4(domain, { ttl: true }),
resolver.resolve6(domain),
resolver.resolveMx(domain),
resolver.resolveTxt(domain),
resolver.resolveNs(domain),
]);
const [a, aaaa, mx, txt, ns] = settled.map((result) =>
result.status === 'fulfilled' ? result.value : []
);
return {
a,
aaaa,
mx: mx.sort((x, y) => x.priority - y.priority),
// I record TXT sono restituiti come array di frammenti
txt: txt.map((chunks) => chunks.join('')),
ns,
};
}
async function reverse(ip) {
try {
return await dns.reverse(ip);
} catch (err) {
if (err.code === dns.NOTFOUND) {
return [];
}
throw err;
}
}
console.log(await inspectDomain('example.com'));
console.log(await reverse('1.1.1.1'));
Poiché dns.lookup() usa il thread pool di libuv, che per impostazione predefinita ha solo quattro thread, un gran numero di risoluzioni lente può bloccare anche altre operazioni come la lettura di file. In applicazioni con molte connessioni in uscita verso host diversi conviene aumentare UV_THREADPOOL_SIZE o introdurre una cache DNS.
TLS: verificare un certificato remoto
Il modulo tls crea connessioni cifrate con un'API speculare a quella di net. Dopo l'handshake, getPeerCertificate() restituisce i dettagli del certificato del server:
import tls from 'node:tls';
function checkCertificate(host, port = 443, timeoutMs = 5000) {
return new Promise((resolve, reject) => {
const socket = tls.connect({
host,
port,
// SNI: indispensabile per i server che ospitano più domini
servername: host,
minVersion: 'TLSv1.2',
timeout: timeoutMs,
});
socket.once('secureConnect', () => {
const cert = socket.getPeerCertificate();
const validTo = new Date(cert.valid_to);
resolve({
host,
authorized: socket.authorized,
protocol: socket.getProtocol(),
cipher: socket.getCipher().name,
subject: cert.subject?.CN,
issuer: cert.issuer?.O,
validTo: validTo.toISOString(),
daysLeft: Math.floor((validTo - Date.now()) / 86_400_000),
altNames: cert.subjectaltname?.split(', ').map((name) => name.replace('DNS:', '')),
});
socket.end();
});
socket.once('timeout', () => {
socket.destroy();
reject(new Error(`Timeout TLS verso ${host}`));
});
socket.once('error', reject);
});
}
const hosts = ['nodejs.org', 'github.com', 'expired.badssl.com'];
const results = await Promise.allSettled(hosts.map((host) => checkCertificate(host)));
for (const [index, result] of results.entries()) {
if (result.status === 'fulfilled') {
console.log(result.value);
} else {
console.error(`${hosts[index]}: ${result.reason.message}`);
}
}
Con le impostazioni predefinite, una connessione verso un certificato scaduto o non valido fallisce con un errore: è il comportamento corretto. L'opzione rejectUnauthorized: false va usata solo per strumenti diagnostici, mai per il traffico applicativo.
HTTP con fetch: timeout e ritentativi
Il fetch globale non ha un timeout predefinito. Il modo idiomatico per imporne uno è AbortSignal.timeout(), combinabile con un segnale esterno tramite AbortSignal.any():
const RETRYABLE_STATUS = new Set([429, 502, 503, 504]);
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
export async function fetchWithRetry(url, options = {}, { retries = 3, timeoutMs = 5000, baseDelayMs = 200 } = {}) {
let lastError;
for (let attempt = 0; attempt <= retries; attempt++) {
const signals = [AbortSignal.timeout(timeoutMs)];
if (options.signal) {
signals.push(options.signal);
}
try {
const response = await fetch(url, { ...options, signal: AbortSignal.any(signals) });
if (!RETRYABLE_STATUS.has(response.status) || attempt === retries) {
return response;
}
lastError = new Error(`HTTP ${response.status}`);
} catch (err) {
// Un annullamento esplicito del chiamante non va ritentato
if (options.signal?.aborted) {
throw err;
}
lastError = err;
}
// Backoff esponenziale con jitter per evitare ondate sincronizzate
const delay = baseDelayMs * 2 ** attempt + Math.random() * baseDelayMs;
await sleep(delay);
}
throw lastError;
}
const response = await fetchWithRetry('https://api.github.com/repos/nodejs/node', {
headers: { Accept: 'application/vnd.github+json' },
});
const repo = await response.json();
console.log(`${repo.full_name}: ${repo.stargazers_count} stelle`);
Il jitter casuale nel backoff ha uno scopo preciso: se centinaia di client falliscono nello stesso istante, senza jitter ritenterebbero tutti nello stesso momento, riproducendo il picco di carico che ha causato l'errore.
Buone pratiche
- Registrare sempre un listener per
errorsu socket e server, per evitare che un singolo client faccia terminare il processo. - Implementare il framing per qualsiasi protocollo su TCP e trattare gli eventi
datacome frammenti arbitrari del flusso. - Rispettare la backpressure, preferendo
pipeline()alle scritture manuali. - Impostare timeout ovunque: connessione, inattività e singole richieste.
- Prevedere una chiusura ordinata in risposta a
SIGTERM, fondamentale in ambienti containerizzati.
Conclusioni
La libreria standard di Node.js offre tutto il necessario per lavorare con la rete a ogni livello, dai datagrammi UDP alle richieste HTTP con ritentativi. La chiave per usarla bene è comprendere il modello sottostante: stream con backpressure, eventi asincroni e un unico thread che non deve mai bloccarsi. Con questi principi chiari, costruire protocolli personalizzati, proxy e strumenti di monitoraggio diventa un esercizio di composizione di pochi blocchi ben progettati.