Cada falha reagenda a entrega com recuo exponencial: trinta segundos na primeira, dobrando a cada tentativa, até um teto de seis horas. O jitter usa metade do intervalo fixa e metade aleatória, para que entregas que falharam juntas não voltem juntas, e o valor de Retry-After é respeitado quando o destino o envia. Há dois limites de desistência, vinte tentativas ou setenta e duas horas de idade, e o que vier primeiro manda a entrega para o estado morta, de onde ela só sai por reenvio explícito.
| Tentativa | Espera antes da próxima | Observação |
|---|
| 1 | 15 a 30 segundos | Absorve reinícios e falhas momentâneas |
| 3 | 1 a 2 minutos | Cobre um deploy do lado do cliente |
| 6 | 8 a 16 minutos | O destino provavelmente já está em pausa |
| 10 | 2 h 08 a 4 h 16 | Falha que já é um incidente do cliente |
| 11 a 20 | 3 a 6 horas | Teto; a vigésima acontece entre 31 e 63 horas depois do evento |
O worker abaixo junta as peças: reserva entregas até o limite do processo, envia com timeout de cinco segundos para a chamada inteira, assina o corpo com HMAC incluindo o timestamp, classifica a resposta e registra o resultado. Ele usa fetch nativo do Node 18 ou superior e o driver pg. A cada falha, além de reagendar a entrega, ele incrementa as falhas seguidas do destino, o que alimenta a pausa descrita na próxima seção.
import crypto from 'node:crypto';
import pg from 'pg';
import { RESERVA_SQL } from './reserva-sql.js';
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL, max: 10 });
const TIMEOUT_MS = 5_000;
const MAX_EM_VOO = 64; // entregas simultaneas deste processo
const MAX_TENTATIVAS = 20;
const IDADE_MAXIMA_MS = 72 * 3_600_000;
const BASE_MS = 30_000;
const TETO_MS = 6 * 3_600_000;
const FALHAS_PARA_PAUSAR = 20;
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
// Recuo exponencial com jitter: metade fixa, metade aleatoria. Nunca antes
// do que o cliente pediu em Retry-After, limitado ao teto.
export function proximoAtrasoMs(tentativas, retryAfterMs = 0) {
const exp = Math.min(TETO_MS, BASE_MS * 2 ** (tentativas - 1));
const atraso = exp / 2 + Math.random() * (exp / 2);
return Math.max(Math.round(atraso), Math.min(retryAfterMs, TETO_MS));
}
export function parseRetryAfter(valor) {
if (!valor) return 0;
const segundos = Number(valor);
if (Number.isFinite(segundos)) return Math.max(0, segundos * 1000);
const data = Date.parse(valor);
return Number.isNaN(data) ? 0 : Math.max(0, data - Date.now());
}
async function enviar(e) {
const corpo = JSON.stringify({ id: e.evento_id, tipo: e.tipo, criado_em: e.criado_em, dados: e.payload });
const ts = String(Math.floor(Date.now() / 1000));
const assinatura = crypto.createHmac('sha256', e.segredo).update(ts + '.' + corpo).digest('hex');
try {
const res = await fetch(e.url, {
method: 'POST',
redirect: 'manual',
signal: AbortSignal.timeout(TIMEOUT_MS),
headers: {
'content-type': 'application/json',
'webhook-id': e.evento_id,
'webhook-timestamp': ts,
'webhook-signature': 'v1=' + assinatura,
},
body: corpo,
});
await res.body?.cancel(); // o corpo da resposta nao interessa
return { status: res.status, retryAfterMs: parseRetryAfter(res.headers.get('retry-after')) };
} catch (err) {
const erro = err.name === 'TimeoutError' ? 'timeout' : String(err.cause?.code || err.message);
return { status: 0, erro };
}
}
async function concluir(e, r) {
if (r.status >= 200 && r.status < 300) {
await pool.query(
"UPDATE webhook_entregas SET status = 'entregue', entregue_em = now(), ultimo_erro = NULL WHERE id = $1",
[e.id],
);
await pool.query(
'UPDATE webhook_endpoints SET falhas_seguidas = 0, primeira_falha_em = NULL, pausado_ate = NULL WHERE id = $1',
[e.endpoint_id],
);
return;
}
const erro = r.erro || 'HTTP ' + r.status;
if (r.status === 410) {
// O cliente disse explicitamente que o endpoint nao existe mais
await pool.query("UPDATE webhook_endpoints SET status = 'desativado' WHERE id = $1", [e.endpoint_id]);
await pool.query("UPDATE webhook_entregas SET status = 'morta', ultimo_erro = $2 WHERE id = $1", [e.id, erro]);
return;
}
const idadeMs = Date.now() - new Date(e.criado_em).getTime();
const desistir = r.status === 413 || e.tentativas >= MAX_TENTATIVAS || idadeMs >= IDADE_MAXIMA_MS;
if (desistir) {
await pool.query("UPDATE webhook_entregas SET status = 'morta', ultimo_erro = $2 WHERE id = $1", [e.id, erro]);
} else {
await pool.query(
"UPDATE webhook_entregas SET status = 'pendente', proxima_em = now() + $2::float8 * interval '1 millisecond', " +
'ultimo_erro = $3 WHERE id = $1',
[e.id, proximoAtrasoMs(e.tentativas, r.retryAfterMs), erro],
);
}
// Falhas seguidas pausam o destino inteiro; cada rodada de sondagem que
// falha estende a pausa, ate uma hora
await pool.query(
'UPDATE webhook_endpoints SET falhas_seguidas = falhas_seguidas + 1, ' +
'primeira_falha_em = coalesce(primeira_falha_em, now()), ' +
'pausado_ate = CASE WHEN falhas_seguidas + 1 >= $2 ' +
"THEN now() + least(interval '1 hour', (falhas_seguidas + 2 - $2) * interval '1 minute') " +
'ELSE pausado_ate END WHERE id = $1',
[e.endpoint_id, FALHAS_PARA_PAUSAR],
);
}
async function reservar(limite) {
const client = await pool.connect();
try {
await client.query('BEGIN');
// Serializa a reserva entre processos: sem isso, dois workers contam as
// mesmas vagas e passam do limite de concorrencia do destino
await client.query('SELECT pg_advisory_xact_lock(4201)');
// Devolve a fila o que ficou preso por um worker que morreu no meio do envio
await client.query(
"UPDATE webhook_entregas SET status = 'pendente', proxima_em = now() " +
"WHERE status = 'enviando' AND reservado_em < now() - interval '1 minute'",
);
const { rows } = await client.query(RESERVA_SQL, [limite]);
await client.query('COMMIT');
return rows;
} catch (err) {
await client.query('ROLLBACK').catch(() => {});
throw err;
} finally {
client.release();
}
}
export async function rodar() {
const emVoo = new Set();
for (;;) {
const vagas = MAX_EM_VOO - emVoo.size;
let entregas = [];
if (vagas > 0) {
try {
entregas = await reservar(vagas);
} catch (err) {
console.error('falha ao reservar entregas', err);
}
}
for (const e of entregas) {
const p = enviar(e)
.then((r) => concluir(e, r))
.catch((err) => console.error('falha ao registrar entrega', e.id, err))
.finally(() => emVoo.delete(p));
emVoo.add(p);
}
if (entregas.length === 0) await Promise.race([...emVoo, sleep(500)]);
}
}
- O timeout de cinco segundos é aplicado com AbortSignal.timeout, que cobre conexão, TLS, envio e espera da resposta.
- Uma entrega presa em enviando por mais de um minuto, porque o processo morreu no meio do envio, volta para a fila na próxima reserva; o cliente pode receber o evento duas vezes, e por isso o identificador do evento vai no cabeçalho.
- O teto de seis horas vale também para Retry-After: um destino que pede para esperar trinta dias não bloqueia a entrega além do teto.
- A retentativa nunca é imediata, nem na primeira falha, porque a primeira falha de um destino sobrecarregado é exatamente o momento em que repetir mais piora tudo.