Na nova forma, a API faz apenas o que é barato e previsível: aplica limites de tamanho enquanto lê o corpo, grava em disco temporário em vez de memória, calcula o hash, guarda o arquivo em um prefixo de quarentena no bucket com uma chave gerada pelo servidor, registra a importação, enfileira um job e responde 202 com o endereço onde o status pode ser consultado. O custo dessa requisição é de entrada e saída, cresce de forma linear com o tamanho do arquivo e a memória fica constante, porque nada é montado no heap.
Antes: tudo dentro da requisição
navegador --POST 38 MB--> API: lê, descompacta, valida e grava 410 mil linhas --> 200 depois de 3 min
(o mesmo processo atende todas as outras rotas)
Depois: receber e processar são etapas diferentes
navegador --POST--> API: limite de tamanho, hash, objeto em quarentena --> 202 + Location
|
+--> fila (jobId = id da importação)
|
v
worker (concorrência 2 por instância)
1. baixa e confere o tamanho
2. inspeciona: tipo real, zip, texto
3. processo filho: heap de 384 MB, 5 min, leitura em streaming
4. grava em lotes, em uma transação
|
v
status: concluida | rejeitada | falhou
navegador --GET /importacoes/:id (2 s, depois 10 s)--> status e contagem de linhasCREATE TABLE importacoes (
id uuid PRIMARY KEY,
tenant_id uuid NOT NULL,
sha256 text NOT NULL,
chave text NOT NULL, -- chave do objeto no bucket, gerada pelo servidor
extensao text NOT NULL, -- '.csv' ou '.xlsx', já validada na entrada
tamanho bigint NOT NULL,
status text NOT NULL, -- na_fila | processando | concluida | rejeitada | falhou
motivo text,
linhas_ok integer,
linhas_erro integer,
criada_em timestamptz NOT NULL DEFAULT now(),
UNIQUE (tenant_id, sha256) -- o mesmo arquivo do mesmo cliente vira a mesma importação
);
import express from 'express';
import multer from 'multer';
import IORedis from 'ioredis';
import pg from 'pg';
import { Queue } from 'bullmq';
import { S3Client } from '@aws-sdk/client-s3';
import { Upload } from '@aws-sdk/lib-storage';
import { createHash, randomUUID } from 'node:crypto';
import { createReadStream } from 'node:fs';
import { unlink } from 'node:fs/promises';
import path from 'node:path';
import { autenticar } from './autenticacao.js'; // preenche req.tenantId
const MiB = 1024 * 1024;
const EXTENSOES_ACEITAS = new Set(['.csv', '.xlsx']);
const BUCKET = process.env.IMPORTACOES_BUCKET;
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL });
const s3 = new S3Client({});
const conexao = new IORedis(process.env.REDIS_URL, { maxRetriesPerRequest: null });
const fila = new Queue('importacoes', { connection: conexao });
// Disco temporário em vez de memória, e limites que o multer aplica enquanto lê o corpo
const upload = multer({
dest: '/tmp/recebidos',
limits: { fileSize: 50 * MiB, files: 1, fields: 4, parts: 5 },
});
async function sha256DoArquivo(caminho) {
const hash = createHash('sha256');
for await (const pedaco of createReadStream(caminho)) hash.update(pedaco);
return hash.digest('hex');
}
const buscarExistente = (tenantId, sha256) =>
pool.query('SELECT id, status FROM importacoes WHERE tenant_id = $1 AND sha256 = $2', [
tenantId,
sha256,
]);
async function receberImportacao(req, res) {
const arquivo = req.file;
if (!arquivo) return res.status(400).json({ erro: 'arquivo_ausente' });
try {
const extensao = path.extname(arquivo.originalname).toLowerCase();
if (!EXTENSOES_ACEITAS.has(extensao)) {
return res.status(415).json({ erro: 'tipo_nao_aceito' });
}
// Duplo clique, retentativa do navegador ou reenvio no dia seguinte: mesma importação
const sha256 = await sha256DoArquivo(arquivo.path);
const existente = await buscarExistente(req.tenantId, sha256);
if (existente.rowCount) return res.status(200).json(existente.rows[0]);
// A chave é do servidor; o nome original nunca vira caminho
const id = randomUUID();
const chave = `quarentena/${req.tenantId}/${id}${extensao}`;
await new Upload({
client: s3,
params: { Bucket: BUCKET, Key: chave, Body: createReadStream(arquivo.path) },
}).done();
const inserida = await pool.query(
`INSERT INTO importacoes (id, tenant_id, sha256, chave, extensao, tamanho, status)
VALUES ($1, $2, $3, $4, $5, $6, 'na_fila')
ON CONFLICT (tenant_id, sha256) DO NOTHING
RETURNING id, status`,
[id, req.tenantId, sha256, chave, extensao, arquivo.size],
);
if (!inserida.rowCount) {
// Corrida entre dois envios do mesmo arquivo: o objeto duplicado expira pela
// regra de ciclo de vida do prefixo de quarentena
const { rows } = await buscarExistente(req.tenantId, sha256);
return res.status(200).json(rows[0]);
}
// jobId igual ao id da importação: reenfileirar a mesma importação não cria job novo
await fila.add(
'importar',
{ importacaoId: id },
{ jobId: id, attempts: 3, backoff: { type: 'exponential', delay: 30_000 } },
);
res.status(202).location(`/importacoes/${id}`).json({ id, status: 'na_fila' });
} finally {
await unlink(arquivo.path).catch(() => {});
}
}
const app = express();
app.post('/importacoes', autenticar, upload.single('arquivo'), receberImportacao);
app.get('/importacoes/:id', autenticar, async (req, res) => {
if (!/^[0-9a-f-]{36}$/.test(req.params.id)) return res.status(404).end();
const { rows } = await pool.query(
`SELECT id, status, motivo, linhas_ok, linhas_erro
FROM importacoes WHERE id = $1 AND tenant_id = $2`,
[req.params.id, req.tenantId],
);
if (!rows.length) return res.status(404).end();
res.json(rows[0]);
});
// Limites do multer viram respostas claras, não 500
app.use((erro, req, res, next) => {
if (erro instanceof multer.MulterError) {
const status = erro.code === 'LIMIT_FILE_SIZE' ? 413 : 400;
return res.status(status).json({ erro: erro.code });
}
next(erro);
});
app.listen(3000);
- O limite de tamanho precisa existir em todas as camadas. O multer interrompe a leitura no primeiro byte acima de 50 MB e devolve 413, mas o proxy ou o balanceador deve ter um limite parecido, para que um corpo de 5 GB não chegue nem a ocupar um processo da aplicação.
- A chave do objeto é gerada pelo servidor. O nome original serve apenas para exibição e nunca vira caminho de arquivo, o que elimina de uma vez sobrescrita de objetos e travessia de diretório com nomes como ../../config.
- O hash torna o envio idempotente. Duplo clique, retentativa do navegador e reenvio do mesmo arquivo no dia seguinte devolvem a mesma importação, e a restrição única no banco resolve a corrida entre dois envios simultâneos.
- A ordem é objeto, linha, job. Se o enfileiramento falhar depois do INSERT, a importação fica na_fila sem job; um processo de reconciliação reenfileira importações nesse estado há mais de cinco minutos, e o jobId igual ao id da importação torna o reenvio inofensivo.
- A extensão é só a primeira triagem, e é barata. Ela não prova nada sobre o conteúdo, que é verificado no worker antes de qualquer parser abrir o arquivo.
O 202 muda o contrato com a interface, e essa é a parte que costuma gerar resistência. A tela passa a mostrar a importação como recebida e consulta o status no endereço do cabeçalho Location, com intervalo crescente: a cada 2 segundos no início e a cada 10 depois do primeiro minuto. O ganho é que a resposta deixa de depender do tamanho do arquivo. Um envio de 50 MB responde em segundos, o usuário pode fechar a aba, e o resultado continua disponível na lista de importações, com a contagem de linhas aceitas e recusadas.