Filas de Tarefas e Processamento Assíncrono

[133] Filas de Tarefas e Processamento Assíncrono

A fila separa aceitar o pedido de executá-lo: a API confirma em milissegundos e o trabalho pesado corre em outro processo. O artigo monta isso com BullMQ e Redis — retentativa com backoff, concorrência, progresso reportado ao cliente, jobs agendados por cron e um painel para enxergar o que está acontecendo.
Javascript

47 min de leitura

Módulo 9 — Tópicos Avançados

Introdução

Imagine que um usuário faz upload de uma planilha com 10.000 produtos para importar em massa. Processar isso sincronamente — na mesma requisição HTTP — significaria deixar o cliente esperando por minutos enquanto o servidor trabalha. Se a conexão cair no meio, todo o trabalho se perde. Se dois usuários fizerem isso ao mesmo tempo, o servidor engasga.

Ou imagine enviar um email de boas-vindas ao cadastrar. Conectar ao servidor SMTP, aguardar a resposta e só então responder ao usuário adiciona latência desnecessária. Se o SMTP estiver lento ou indisponível, o cadastro falha — mesmo que isso não devesse impedir o registro.

O problema central em ambos os casos é o mesmo: estamos misturando a requisição do usuário com trabalho que pode — e deve — acontecer de forma independente e assíncrona. Filas de tarefas (job queues) resolvem isso desacoplando a aceitação de um trabalho da sua execução.

O padrão produtor-consumidor

A ideia central de uma fila de tarefas é simples. Um produtor aceita o trabalho e o coloca em uma fila. Um ou mais consumidores (workers) processam as tarefas da fila independentemente, em paralelo, no seu próprio ritmo. O produtor pode responder imediatamente ao cliente — "seu trabalho foi enfileirado" — sem esperar o resultado.

Sem filas:
  Cliente → POST /importar → [processa 10.000 itens... espera 3 min] → 200 OK

Com filas:
  Cliente → POST /importar → [enfileira job] → 202 Accepted (imediato)
                                    ↓
                              [Worker processa em background]
                                    ↓
                              [Notifica o cliente via WebSocket quando terminar]

As vantagens vão além da velocidade de resposta. Filas permitem controle de concorrência — você decide quantos workers processam em paralelo, evitando sobrecarregar serviços externos. Permitem retry automático com backoff exponencial quando um job falha. Permitem priorizar tarefas urgentes. E dão visibilidade sobre o estado de cada job — pendente, processando, concluído, falhou.

BullMQ — filas de tarefas com Redis

O BullMQ é a biblioteca mais robusta e popular para filas de tarefas em Node.js. Usa o Redis como backend — o Redis é extremamente rápido para esta finalidade, com suporte nativo a listas, pub/sub e atomicidade.

npm install bullmq ioredis

A arquitetura do BullMQ tem três componentes principais. A Queue é o ponto de entrada — onde os jobs são adicionados. O Worker é o processador — consome jobs da fila e executa a lógica. O QueueEvents é o observador — permite reagir a eventos como job concluído ou falhado.

Configuração base

Antes de criar filas específicas, vamos centralizar a configuração da conexão Redis que todas as filas vão compartilhar.

// src/infrastructure/queue/redisConfig.js
// A configuração Redis é compartilhada entre todas as filas
// O BullMQ cria conexões próprias — não compartilhamos com o Socket.IO
// para evitar interferência entre os sistemas

const redisConfig = {
  host: process.env.REDIS_HOST || 'localhost',
  port: Number(process.env.REDIS_PORT) || 6379,
  password: process.env.REDIS_PASSWORD || undefined,

  // Configurações importantes para produção
  maxRetriesPerRequest: null,  // BullMQ requer este valor específico
  enableReadyCheck: false,     // evita erro ao conectar durante startup

  // Reconnect automático com backoff exponencial
  retryStrategy: (tentativas) => {
    if (tentativas > 10) return null; // desiste após 10 tentativas
    return Math.min(tentativas * 200, 5000); // espera até 5s entre tentativas
  },
};

module.exports = { redisConfig };
// src/infrastructure/queue/filas.js
// Define e exporta todas as filas da aplicação
// Centralizar aqui evita nomes duplicados e facilita listar todas as filas

const { Queue } = require('bullmq');
const { redisConfig } = require('./redisConfig');

// Cada fila tem um nome único — o nome é a chave no Redis
// Agrupe jobs relacionados na mesma fila para simplificar o código dos workers

const filaEmail = new Queue('email', {
  connection: redisConfig,
  defaultJobOptions: {
    // Tenta até 3 vezes antes de mover para a fila de falhados
    attempts: 3,

    // Backoff exponencial: 1s, 2s, 4s entre tentativas
    // Evita bombardear um serviço que está temporariamente indisponível
    backoff: {
      type: 'exponential',
      delay: 1000,
    },

    // Remove jobs bem-sucedidos após 24h — libera memória no Redis
    removeOnComplete: { age: 86400, count: 100 },

    // Mantém jobs falhados por 7 dias para análise
    removeOnFail: { age: 604800 },
  },
});

const filaImportacao = new Queue('importacao', {
  connection: redisConfig,
  defaultJobOptions: {
    // Importações têm mais tentativas — são operações longas e mais suscetíveis
    attempts: 2,
    backoff: { type: 'exponential', delay: 5000 },
    removeOnComplete: { age: 3600 },  // 1 hora — importações são únicas
    removeOnFail: { age: 604800 },
  },
});

const filaNotificacao = new Queue('notificacao', {
  connection: redisConfig,
  defaultJobOptions: {
    attempts: 5,  // notificações têm mais tentativas — importante entregá-las
    backoff: { type: 'exponential', delay: 2000 },
    removeOnComplete: { age: 3600 },
    removeOnFail: { age: 604800 },
  },
});

const filaRelatorio = new Queue('relatorio', {
  connection: redisConfig,
  defaultJobOptions: {
    attempts: 2,
    backoff: { type: 'fixed', delay: 10000 },
    removeOnComplete: { age: 86400 },
    removeOnFail: { age: 604800 },
  },
});

module.exports = { filaEmail, filaImportacao, filaNotificacao, filaRelatorio };

Workers — processando os jobs

Cada worker é um processo ou módulo separado que fica ouvindo a fila e processando jobs conforme eles chegam. Em produção, workers podem rodar em processos completamente separados do servidor HTTP — isso isola falhas e permite escalar independentemente.

// src/workers/emailWorker.js
// Worker responsável por processar todos os jobs de email
// Este arquivo pode rodar como processo separado: node src/workers/emailWorker.js

const { Worker } = require('bullmq');
const nodemailer = require('nodemailer');
const { redisConfig } = require('../infrastructure/queue/redisConfig');

// Configura o transportador de email uma vez — reutilizado por todos os jobs
const transportador = nodemailer.createTransport({
  host: process.env.SMTP_HOST,
  port: Number(process.env.SMTP_PORT) || 587,
  secure: false,
  auth: {
    user: process.env.SMTP_USER,
    pass: process.env.SMTP_PASS,
  },
});

// Mapa de templates — cada tipo de email tem sua própria função de renderização
// Separar templates do worker facilita manutenção e testes
const templates = {
  boasVindas: (dados) => ({
    subject: `Bem-vindo(a), ${dados.nome}!`,
    html: `
      <h1>Olá, ${dados.nome}!</h1>
      <p>Sua conta foi criada com sucesso.</p>
      <p>Acesse: <a href="${process.env.FRONTEND_URL}/login">Entrar</a></p>
    `,
  }),

  tarefaCriada: (dados) => ({
    subject: `Nova tarefa: ${dados.tarefa.titulo}`,
    html: `
      <h2>Tarefa criada com sucesso</h2>
      <p><strong>${dados.tarefa.titulo}</strong></p>
      ${dados.tarefa.prazo ? `<p>Prazo: ${new Date(dados.tarefa.prazo).toLocaleDateString('pt-BR')}</p>` : ''}
      <p><a href="${process.env.FRONTEND_URL}/tarefas">Ver minhas tarefas</a></p>
    `,
  }),

  redefinirSenha: (dados) => ({
    subject: 'Redefinição de senha',
    html: `
      <h2>Redefinição de senha</h2>
      <p>Clique no link abaixo para redefinir sua senha:</p>
      <p><a href="${process.env.FRONTEND_URL}/nova-senha?token=${dados.token}">
        Redefinir senha
      </a></p>
      <p>O link expira em 1 hora.</p>
      <p>Se não foi você, ignore este email.</p>
    `,
  }),
};

// O worker escuta a fila 'email' e processa um job por vez (concurrency: 1)
// Para emails, processamento serial evita sobrecarregar o SMTP
// Aumente a concorrência para outros tipos de jobs se necessário
const worker = new Worker(
  'email',
  async (job) => {
    const { tipo, para, dados } = job.data;

    // Registra o progresso — visível no painel do BullMQ
    await job.updateProgress(10);

    // Valida que temos um template para este tipo de email
    const gerarTemplate = templates[tipo];
    if (!gerarTemplate) {
      // Lançar erro aqui faz o BullMQ tentar novamente (até o limite de attempts)
      throw new Error(`Template de email desconhecido: ${tipo}`);
    }

    const template = gerarTemplate(dados);
    await job.updateProgress(30);

    // Envia o email — pode lançar erro se o SMTP estiver indisponível
    // BullMQ vai retentar automaticamente com o backoff configurado
    const resultado = await transportador.sendMail({
      from: `"${process.env.EMAIL_NOME || 'Gestão App'}" <${process.env.EMAIL_FROM}>`,
      to: para,
      subject: template.subject,
      html: template.html,
    });

    await job.updateProgress(100);

    // O retorno do processador é salvo no job — útil para debugging
    return {
      messageId: resultado.messageId,
      aceito: resultado.accepted,
    };
  },
  {
    connection: redisConfig,
    concurrency: 1,       // processa um email por vez para não sobrecarregar SMTP

    // Limita a taxa de processamento: no máximo 10 emails por segundo
    // Útil para respeitar limites de API de provedores de email
    limiter: {
      max: 10,
      duration: 1000,
    },
  }
);

// Eventos do worker para logging e monitoramento
worker.on('completed', (job) => {
  console.log(`[EmailWorker] ✅ Job ${job.id} (${job.data.tipo}) concluído`);
});

worker.on('failed', (job, erro) => {
  console.error(
    `[EmailWorker] ❌ Job ${job?.id} (${job?.data?.tipo}) falhou: ${erro.message}`
  );
  // Em produção, você enviaria este erro para o Sentry
});

worker.on('error', (erro) => {
  console.error('[EmailWorker] Erro no worker:', erro.message);
});

console.log('[EmailWorker] 🚀 Aguardando jobs na fila "email"...');

module.exports = worker; // exporta para testes
// src/workers/importacaoWorker.js
// Worker para importação de dados em massa — jobs de longa duração

const { Worker } = require('bullmq');
const { parse } = require('csv-parse/sync');
const { redisConfig } = require('../infrastructure/queue/redisConfig');
const Produto = require('../infrastructure/database/models/ProdutoModel');
const { notificarUsuario } = require('../socketHandlers/notificacoes');

// 'io' é injetado na criação — permite notificar o usuário via WebSocket
function criarImportacaoWorker(io) {
  const worker = new Worker(
    'importacao',
    async (job) => {
      const { usuarioId, conteudoCSV, nomeArquivo } = job.data;
      const resultados = { criados: 0, atualizados: 0, erros: [] };

      await job.updateProgress(5);

      // Parseia o CSV recebido como string
      let registros;
      try {
        registros = parse(conteudoCSV, {
          columns: true,        // usa a primeira linha como cabeçalho
          skip_empty_lines: true,
          trim: true,
        });
      } catch (erro) {
        throw new Error(`CSV inválido: ${erro.message}`);
      }

      const total = registros.length;
      console.log(`[ImportacaoWorker] Processando ${total} registros de ${nomeArquivo}`);

      // Processa em lotes de 100 para não sobrecarregar o banco
      // e para poder reportar progresso gradualmente
      const TAMANHO_LOTE = 100;

      for (let i = 0; i < registros.length; i += TAMANHO_LOTE) {
        const lote = registros.slice(i, i + TAMANHO_LOTE);

        // Processa cada registro do lote com tratamento de erro individual
        // Um registro inválido não deve abortar toda a importação
        for (const registro of lote) {
          try {
            // Upsert: cria ou atualiza o produto pelo código
            await Produto.findOneAndUpdate(
              { codigo: registro.codigo },
              {
                nome: registro.nome,
                preco: parseFloat(registro.preco),
                estoque: parseInt(registro.estoque, 10),
                categoria: registro.categoria,
              },
              { upsert: true, new: true, runValidators: true }
            );

            resultados.criados++;
          } catch (erro) {
            // Registra o erro mas continua o processamento
            resultados.erros.push({
              linha: i + lote.indexOf(registro) + 2, // +2: 1-indexed + header
              codigo: registro.codigo,
              erro: erro.message,
            });
          }
        }

        // Atualiza o progresso proporcionalmente ao lote concluído
        const progresso = Math.round(((i + lote.length) / total) * 90) + 5;
        await job.updateProgress(progresso);

        // Notifica o usuário do progresso via WebSocket
        if (io) {
          notificarUsuario(io, usuarioId, {
            tipo: 'importacao:progresso',
            mensagem: `Importando... ${progresso}%`,
            jobId: job.id,
          });
        }
      }

      await job.updateProgress(100);

      // Notifica o usuário que a importação terminou
      if (io) {
        notificarUsuario(io, usuarioId, {
          tipo: 'importacao:concluida',
          mensagem: `Importação concluída: ${resultados.criados} produtos, ${resultados.erros.length} erros`,
          jobId: job.id,
          resultados,
        });
      }

      return resultados;
    },
    {
      connection: redisConfig,
      // Importações são pesadas — processa uma por vez para não travar o banco
      concurrency: 1,
    }
  );

  worker.on('failed', (job, erro) => {
    console.error(`[ImportacaoWorker] ❌ Job ${job?.id} falhou:`, erro.message);

    if (io && job?.data?.usuarioId) {
      notificarUsuario(io, job.data.usuarioId, {
        tipo: 'importacao:falhou',
        mensagem: `Importação falhou: ${erro.message}`,
        jobId: job.id,
      });
    }
  });

  return worker;
}

module.exports = { criarImportacaoWorker };

Produtores — enfileirando jobs a partir da API

Agora que os workers existem, precisamos enfileirar jobs a partir dos controllers e dos eventos de domínio.

// src/infrastructure/queue/produtores.js
// Funções de alto nível para enfileirar cada tipo de job
// Encapsular a lógica aqui evita que os controllers conheçam o BullMQ diretamente

const { filaEmail, filaImportacao, filaNotificacao } = require('./filas');

// Enfileira um email para envio assíncrono
// Retorna imediatamente — o envio acontece em background
async function enfileirarEmail(tipo, para, dados, opcoes = {}) {
  const job = await filaEmail.add(
    tipo,    // nome do job — aparece no painel de monitoramento
    { tipo, para, dados },
    {
      // Jobs de email de redefinição de senha têm prioridade maior
      priority: opcoes.prioridade || 0,

      // Atraso opcional — útil para emails de "lembrete" agendados
      delay: opcoes.atrasoMs || 0,

      // ID único para evitar duplicatas — útil para emails de boas-vindas
      jobId: opcoes.jobId,
    }
  );

  console.log(`[Fila] Email enfileirado: ${tipo} → ${para} (job ${job.id})`);
  return job;
}

// Enfileira uma importação de CSV
async function enfileirarImportacao(usuarioId, conteudoCSV, nomeArquivo) {
  const job = await filaImportacao.add(
    'importar-csv',
    { usuarioId, conteudoCSV, nomeArquivo },
    {
      // Importações podem demorar muito — aumenta o timeout
      // O default do BullMQ é sem timeout — ajuste conforme necessário
    }
  );

  console.log(`[Fila] Importação enfileirada: ${nomeArquivo} (job ${job.id})`);
  return job;
}

// Enfileira geração de relatório — pode ser agendado para horário de baixo tráfego
async function enfileirarRelatorio(usuarioId, tipo, filtros, agendar) {
  const opcoes = {};

  if (agendar) {
    // Agenda para um horário específico — útil para relatórios diários
    const agora = new Date();
    const horarioAlvo = new Date(agendar);
    const atrasoMs = horarioAlvo.getTime() - agora.getTime();

    if (atrasoMs > 0) {
      opcoes.delay = atrasoMs;
    }
  }

  return filaRelatorio.add('gerar-relatorio', { usuarioId, tipo, filtros }, opcoes);
}

module.exports = { enfileirarEmail, enfileirarImportacao, enfileirarRelatorio };
// Integrando nos eventos de domínio — container.js
// Quando uma tarefa é criada, enfileira o email automaticamente
// sem que o caso de uso precise saber sobre filas ou email

const { enfileirarEmail } = require('./infrastructure/queue/produtores');
const EVENTOS = require('./domain/events/EventosDominio');

eventEmitter.on(EVENTOS.TAREFA_CRIADA, async ({ tarefa, usuario }) => {
  // Invalida cache (síncrono)
  cache.del(`stats:${usuario._id}`);

  // Enfileira email (assíncrono — não bloqueia a resposta HTTP)
  await enfileirarEmail(
    'tarefaCriada',
    usuario.email,
    { tarefa: tarefa.toJSON(), nome: usuario.nome },
    // ID único para evitar email duplicado se o evento for emitido duas vezes
    { jobId: `tarefa-criada:${tarefa.id}` }
  );
});
// src/interfaces/http/routes/importacao.js
// Rota para upload de CSV — responde imediatamente após enfileirar

const express = require('express');
const multer = require('multer');  // npm install multer
const { autenticar } = require('../middlewares/auth');
const { enfileirarImportacao } = require('../../../infrastructure/queue/produtores');
const { filaImportacao } = require('../../../infrastructure/queue/filas');

const router = express.Router();

// Multer armazena o arquivo em memória — para arquivos grandes, use disco
const upload = multer({
  storage: multer.memoryStorage(),
  limits: { fileSize: 10 * 1024 * 1024 }, // 10MB máximo
  fileFilter: (req, file, cb) => {
    if (file.mimetype !== 'text/csv' && !file.originalname.endsWith('.csv')) {
      return cb(new Error('Apenas arquivos CSV são aceitos.'));
    }
    cb(null, true);
  },
});

router.post('/produtos', autenticar, upload.single('arquivo'), async (req, res, next) => {
  try {
    if (!req.file) {
      return res.status(400).json({ erro: 'Arquivo CSV não enviado.' });
    }

    const conteudoCSV = req.file.buffer.toString('utf-8');
    const nomeArquivo = req.file.originalname;

    // Enfileira e responde IMEDIATAMENTE — não espera o processamento
    const job = await enfileirarImportacao(
      req.usuario._id.toString(),
      conteudoCSV,
      nomeArquivo
    );

    // 202 Accepted — trabalho aceito mas ainda não concluído
    res.status(202).json({
      mensagem: 'Importação enfileirada com sucesso.',
      jobId: job.id,
      // O cliente pode usar este endpoint para verificar o progresso
      statusUrl: `/importacao/status/${job.id}`,
    });
  } catch (erro) {
    next(erro);
  }
});

// Endpoint para verificar o status de um job de importação
router.get('/status/:jobId', autenticar, async (req, res, next) => {
  try {
    const job = await filaImportacao.getJob(req.params.jobId);

    if (!job) {
      return res.status(404).json({ erro: 'Job não encontrado.' });
    }

    // Verifica se o job pertence ao usuário autenticado
    if (job.data.usuarioId !== req.usuario._id.toString()) {
      return res.status(403).json({ erro: 'Sem acesso a este job.' });
    }

    const estado = await job.getState();
    const progresso = job.progress;

    res.json({
      jobId: job.id,
      estado,          // 'waiting', 'active', 'completed', 'failed', 'delayed'
      progresso,       // 0-100
      resultado: estado === 'completed' ? job.returnvalue : null,
      erro: estado === 'failed' ? job.failedReason : null,
      criadoEm: new Date(job.timestamp).toISOString(),
    });
  } catch (erro) {
    next(erro);
  }
});

module.exports = router;

Jobs repetidos — agendamento com cron

O BullMQ suporta jobs que se repetem automaticamente em um horário definido, similar ao cron do Unix. Isso é perfeito para relatórios diários, limpeza de dados antigos e alertas periódicos.

// src/infrastructure/queue/agendamentos.js
// Jobs recorrentes — configurados uma vez no startup da aplicação

const { filaEmail, filaRelatorio, filaNotificacao } = require('./filas');

async function configurarAgendamentos() {
  // Remove agendamentos antigos antes de recriar
  // Evita duplicatas se a aplicação reiniciar
  await filaRelatorio.obliterate({ force: true }).catch(() => {});
  await filaNotificacao.obliterate({ force: true }).catch(() => {});

  // Relatório diário de tarefas — todo dia às 8h (horário de Brasília)
  await filaRelatorio.add(
    'relatorio-diario',
    { tipo: 'tarefas-do-dia' },
    {
      repeat: {
        pattern: '0 8 * * *',      // cron: minuto hora dia mês diaSemana
        tz: 'America/Sao_Paulo',   // fuso horário
      },
      jobId: 'relatorio-diario',   // ID fixo para evitar duplicatas
    }
  );

  // Verificação de tarefas atrasadas — toda hora
  await filaNotificacao.add(
    'verificar-tarefas-atrasadas',
    { tipo: 'verificacao-atraso' },
    {
      repeat: { pattern: '0 * * * *' },
      jobId: 'verificar-atraso',
    }
  );

  // Limpeza de tokens expirados — todo domingo à meia-noite
  await filaEmail.add(
    'limpeza-tokens',
    { tipo: 'limpeza-interna' },
    {
      repeat: { pattern: '0 0 * * 0' },
      jobId: 'limpeza-tokens',
    }
  );

  console.log('[Agendamentos] ✅ Jobs recorrentes configurados.');
}

module.exports = { configurarAgendamentos };
// src/workers/notificacaoWorker.js
// Processa jobs de notificação — incluindo o job recorrente de tarefas atrasadas

const { Worker } = require('bullmq');
const { redisConfig } = require('../infrastructure/queue/redisConfig');
const Tarefa = require('../infrastructure/database/models/TarefaModel');
const Usuario = require('../infrastructure/database/models/UsuarioModel');
const { enfileirarEmail } = require('../infrastructure/queue/produtores');

function criarNotificacaoWorker(io) {
  return new Worker(
    'notificacao',
    async (job) => {
      const { tipo } = job.data;

      if (tipo === 'verificacao-atraso') {
        // Busca todas as tarefas pendentes com prazo passado
        const atrasadas = await Tarefa.find({
          status: { $in: ['pendente', 'em_progresso'] },
          prazo: { $lt: new Date() },
          // Evita notificar múltiplas vezes — usa um campo de controle
          notificacaoAtrasoEnviada: { $ne: true },
        })
          .populate('usuario', 'email nome')
          .lean();

        console.log(`[NotificacaoWorker] ${atrasadas.length} tarefas atrasadas encontradas.`);

        for (const tarefa of atrasadas) {
          if (!tarefa.usuario?.email) continue;

          // Enfileira um email para cada tarefa atrasada
          await enfileirarEmail(
            'tarefaAtrasada',
            tarefa.usuario.email,
            { tarefa, nome: tarefa.usuario.nome },
            { jobId: `atraso:${tarefa._id}` }
          );

          // Marca como notificado para não enviar novamente
          await Tarefa.findByIdAndUpdate(tarefa._id, {
            notificacaoAtrasoEnviada: true,
          });
        }

        return { tarefasNotificadas: atrasadas.length };
      }
    },
    { connection: redisConfig, concurrency: 2 }
  );
}

module.exports = { criarNotificacaoWorker };

Monitoramento — Bull Board

O Bull Board é uma interface web que permite visualizar e gerenciar todas as filas em tempo real — ver jobs pendentes, ativos, concluídos e falhados, inspecionar dados e reintentar jobs manualmente.

// npm install @bull-board/express @bull-board/api

const { createBullBoard } = require('@bull-board/api');
const { BullMQAdapter } = require('@bull-board/api/bullMQAdapter');
const { ExpressAdapter } = require('@bull-board/express');
const { filaEmail, filaImportacao, filaNotificacao, filaRelatorio } = require('./filas');

function configurarBullBoard(app) {
  const adaptadorExpress = new ExpressAdapter();

  // Rota base do painel — proteja com autenticação em produção!
  adaptadorExpress.setBasePath('/admin/filas');

  createBullBoard({
    queues: [
      new BullMQAdapter(filaEmail),
      new BullMQAdapter(filaImportacao),
      new BullMQAdapter(filaNotificacao),
      new BullMQAdapter(filaRelatorio),
    ],
    serverAdapter: adaptadorExpress,
  });

  // Middleware simples de autenticação para o painel
  // Em produção, use autenticação mais robusta
  app.use(
    '/admin/filas',
    (req, res, next) => {
      const senhaAdmin = req.headers['x-admin-key'];
      if (senhaAdmin !== process.env.ADMIN_KEY) {
        return res.status(401).json({ erro: 'Não autorizado.' });
      }
      next();
    },
    adaptadorExpress.getRouter()
  );

  console.log('[BullBoard] 📊 Painel de filas disponível em /admin/filas');
}

module.exports = { configurarBullBoard };

Front-end — acompanhando o progresso de jobs

// src/hooks/useProgressoImportacao.js
// Hook para acompanhar o progresso de um job de importação
// Combina polling REST com atualizações via WebSocket
import { useState, useEffect, useCallback } from 'react';
import { useSocket } from './useSocket';

export function useProgressoImportacao(jobId) {
  const [estado, setEstado] = useState(null);
  const [progresso, setProgresso] = useState(0);
  const [resultado, setResultado] = useState(null);
  const [erro, setErro] = useState(null);
  const { ouvir } = useSocket();

  // Polling de fallback — caso o WebSocket não esteja disponível
  const verificarStatus = useCallback(async () => {
    if (!jobId) return;
    try {
      const res = await fetch(`/importacao/status/${jobId}`, {
        headers: { Authorization: `Bearer ${token}` },
      });
      const dados = await res.json();

      setEstado(dados.estado);
      setProgresso(dados.progresso || 0);

      if (dados.estado === 'completed') setResultado(dados.resultado);
      if (dados.estado === 'failed') setErro(dados.erro);
    } catch { /* ignora erros de rede durante polling */ }
  }, [jobId]);

  useEffect(() => {
    if (!jobId) return;

    // Ouve eventos WebSocket de progresso — atualização em tempo real
    const removerProgresso = ouvir('notificacao:nova', (dados) => {
      const notif = dados.notificacao;
      if (notif.jobId !== jobId) return;

      if (notif.tipo === 'importacao:progresso') {
        // Extrai porcentagem da mensagem "Importando... 45%"
        const match = notif.mensagem.match(/(\d+)%/);
        if (match) setProgresso(Number(match[1]));
        setEstado('active');
      }

      if (notif.tipo === 'importacao:concluida') {
        setEstado('completed');
        setProgresso(100);
        setResultado(notif.resultados);
      }

      if (notif.tipo === 'importacao:falhou') {
        setEstado('failed');
        setErro(notif.mensagem);
      }
    });

    // Polling a cada 3 segundos como fallback
    verificarStatus();
    const intervalo = setInterval(verificarStatus, 3000);

    return () => {
      removerProgresso();
      clearInterval(intervalo);
    };
  }, [jobId, ouvir, verificarStatus]);

  const concluido = estado === 'completed';
  const falhou = estado === 'failed';
  const processando = estado === 'active' || estado === 'waiting';

  return { estado, progresso, resultado, erro, concluido, falhou, processando };
}
// src/components/ImportacaoProdutos.jsx
// Componente completo de upload e acompanhamento de importação
import { useState } from 'react';
import { useProgressoImportacao } from '../hooks/useProgressoImportacao';

export function ImportacaoProdutos() {
  const [jobId, setJobId] = useState(null);
  const [enviando, setEnviando] = useState(false);
  const { progresso, resultado, erro, concluido, falhou, processando } =
    useProgressoImportacao(jobId);

  async function handleUpload(e) {
    const arquivo = e.target.files[0];
    if (!arquivo) return;

    setEnviando(true);

    const formData = new FormData();
    formData.append('arquivo', arquivo);

    try {
      const res = await fetch('/importacao/produtos', {
        method: 'POST',
        headers: { Authorization: `Bearer ${token}` },
        body: formData,
        // Não definir Content-Type — o browser define com o boundary correto para multipart
      });

      const dados = await res.json();

      if (res.ok) {
        // Salva o jobId para acompanhar o progresso
        setJobId(dados.jobId);
      } else {
        console.error('Erro ao enfileirar importação:', dados.erro);
      }
    } finally {
      setEnviando(false);
    }
  }

  return (
    <div className="importacao">
      <h2>Importar Produtos</h2>
      <p className="instrucao">
        Faça upload de um arquivo CSV com colunas: codigo, nome, preco, estoque, categoria
      </p>

      {/* Área de upload — só mostra se não há importação em andamento */}
      {!jobId && (
        <label className="upload-area">
          <input
            type="file"
            accept=".csv"
            onChange={handleUpload}
            disabled={enviando}
            className="upload-input"
          />
          <div className="upload-conteudo">
            {enviando ? (
              <p>Enviando arquivo...</p>
            ) : (
              <>
                <span className="upload-icone">📁</span>
                <p>Clique ou arraste o CSV aqui</p>
                <small>Máximo: 10MB</small>
              </>
            )}
          </div>
        </label>
      )}

      {/* Barra de progresso — exibida durante o processamento */}
      {jobId && (
        <div className="progresso-container">
          <div className="progresso-header">
            <p>
              {processando && 'Importando produtos...'}
              {concluido && '✅ Importação concluída!'}
              {falhou && '❌ Importação falhou'}
            </p>
            <span>{progresso}%</span>
          </div>

          <div className="barra-progresso">
            <div
              className={`barra-preenchimento ${concluido ? 'completa' : ''}`}
              style={{ width: `${progresso}%` }}
            />
          </div>

          {/* Resultado detalhado quando concluído */}
          {concluido && resultado && (
            <div className="resultado-importacao">
              <p>✅ {resultado.criados} produtos importados</p>
              {resultado.erros?.length > 0 && (
                <details className="erros-detalhes">
                  <summary>⚠️ {resultado.erros.length} erros</summary>
                  <ul>
                    {resultado.erros.map((e, i) => (
                      <li key={i}>
                        Linha {e.linha} (código: {e.codigo}): {e.erro}
                      </li>
                    ))}
                  </ul>
                </details>
              )}
              <button onClick={() => setJobId(null)} className="btn-primario">
                Nova importação
              </button>
            </div>
          )}

          {falhou && (
            <div className="erro-importacao">
              <p>{erro}</p>
              <button onClick={() => setJobId(null)} className="btn-secundario">
                Tentar novamente
              </button>
            </div>
          )}
        </div>
      )}
    </div>
  );
}

Tarefa para você

Adicione filas de tarefas na aplicação do Módulo 6:

# 1. Configure o Redis localmente para desenvolvimento
#    docker run -d -p 6379:6379 redis:7-alpine

# 2. Crie a estrutura de filas:
#    src/infrastructure/queue/redisConfig.js
#    src/infrastructure/queue/filas.js (filaEmail, filaNotificacao)
#    src/infrastructure/queue/produtores.js

# 3. Implemente o emailWorker.js com dois templates:
#    - boasVindas: enviado quando um usuário se cadastra
#    - tarefaCriada: enviado quando uma tarefa é criada
#    Teste com Ethereal (https://ethereal.email) — captura emails sem enviar

# 4. Integre com o EventEmitter de domínio:
#    - TAREFA_CRIADA → enfileirarEmail('tarefaCriada', ...)
#    - Use jobId único para evitar emails duplicados

# 5. Crie a rota POST /importacao/produtos e o worker de importação
#    Teste com um CSV de 100 produtos fictícios
#    Observe o progresso no console do worker

# 6. Configure o Bull Board em /admin/filas
#    Proteja com a variável de ambiente ADMIN_KEY
#    Acesse o painel e observe os jobs sendo processados

# 7. Adicione o job recorrente de verificação de tarefas atrasadas
#    Configure para rodar a cada 1 minuto em desenvolvimento:
#    pattern: '* * * * *' (toda vez que o minuto muda)

# 8. Implemente o componente ImportacaoProdutos no front-end
#    Acompanhe o progresso via polling ou WebSocket
Ver solução — as filas inteiras — email, importação de CSV, painel protegido e job de minuto
// FILAS DE TAREFAS com BullMQ.
//
//   1 Redis                → docker run -d -p 6379:6379 redis:7-alpine
//   2 estrutura de filas   → src/infrastructure/queue/*.js
//   3 emailWorker          → src/workers/emailWorker.js (boasVindas, tarefaCriada)
//   4 EventEmitter + jobId → src/infrastructure/queue/produtores.js
//   5 importação de CSV    → src/rotas/importacao.js + src/workers/importacaoWorker.js
//   6 Bull Board           → src/admin/bullBoard.js (protegido por ADMIN_KEY)
//   7 job recorrente       → src/jobs/recorrentes.js ('* * * * *')
//   8 front-end            → src/componentes/ImportacaoProdutos.jsx
//
// 23 testes verdes contra um Redis DE VERDADE, e um deles envia um email real
// pelo SMTP do Ethereal. O item 1 é a única parte que não rodou como está
// escrita: não havia docker na máquina, e o Redis subiu pelo
// redis-memory-server (mesmo servidor, compilado do fonte).

// ---- src/infrastructure/queue/redisConfig.js
// 1 — A CONEXÃO. Um objeto de opções, não um cliente: o BullMQ abre as
// conexões de que precisa (o Worker usa uma exclusiva, bloqueante).
function conexaoRedis(env = process.env) {
  return {
    host: env.REDIS_HOST || "127.0.0.1",
    port: Number(env.REDIS_PORT || 6379),
    password: env.REDIS_PASSWORD || undefined,
    // Historicamente obrigatório para o Worker: o comando BZPOPMIN fica
    // bloqueado esperando job, e o ioredis o daria por perdido. O BullMQ 6
    // aceita a conexão sem isto, mas manter é o que evita "Reached the max
    // retries per request limit" em rede instável.
    maxRetriesPerRequest: null,
  };
}

module.exports = { conexaoRedis };

// ---- src/infrastructure/queue/filas.js
// 2 — AS FILAS. Uma por assunto, não uma para tudo: assim dá para dar
// concorrência, prioridade e política de retentativa diferentes a cada uma.
const { Queue } = require("bullmq");
const { conexaoRedis } = require("./redisConfig");

const PADRAO_DO_JOB = {
  attempts: 3,
  backoff: { type: "exponential", delay: 1000 },
  // Sem teto, a lista de concluídos cresce para sempre e o Redis vira o gargalo.
  // ATENÇÃO: isto interage com o jobId — ver produtores.js.
  removeOnComplete: { age: 3600, count: 1000 },
  removeOnFail: { age: 24 * 3600 },
};

function criarFilas(conexao = conexaoRedis()) {
  const opcoes = { connection: conexao, defaultJobOptions: PADRAO_DO_JOB };
  return {
    filaEmail: new Queue("email", opcoes),
    filaNotificacao: new Queue("notificacao", opcoes),
    filaImportacao: new Queue("importacao", {
      connection: conexao,
      defaultJobOptions: { ...PADRAO_DO_JOB, attempts: 1 }, // importar duas vezes duplicaria produtos
    }),
  };
}

module.exports = { criarFilas, PADRAO_DO_JOB };

// ---- src/infrastructure/queue/produtores.js
// 4 — QUEM ENFILEIRA. O produtor não sabe processar nada; só descreve o
// trabalho e some. É o que permite responder ao HTTP em milissegundos.
const TAREFA_CRIADA = "tarefa:criada";
const USUARIO_CADASTRADO = "usuario:cadastrado";

function criarProdutores({ filaEmail, filaNotificacao }) {
  // jobId determinístico: se a mesma tarefa disparar o evento duas vezes (um
  // retry do HTTP, dois ouvintes registrados sem querer), o segundo add é
  // ignorado pelo Redis e o usuário recebe UM email.
  const enfileirarEmail = (template, dados, chave) =>
    filaEmail.add(template, dados, { jobId: `email:${template}:${chave}` });

  const enfileirarNotificacao = (dados) =>
    filaNotificacao.add("push", dados, { jobId: `push:${dados.tarefaId}` });

  function ligarAoBarramento(barramento) {
    barramento.on(TAREFA_CRIADA, (tarefa) =>
      enfileirarEmail("tarefaCriada", { para: tarefa.emailDono, tarefa }, tarefa.id)
    );
    barramento.on(USUARIO_CADASTRADO, (usuario) =>
      enfileirarEmail("boasVindas", { para: usuario.email, usuario }, usuario.id)
    );
  }

  return { enfileirarEmail, enfileirarNotificacao, ligarAoBarramento };
}

module.exports = { criarProdutores, TAREFA_CRIADA, USUARIO_CADASTRADO };

// ---- src/workers/emailWorker.js
// 3 — O WORKER DE EMAIL. Processo separado do servidor web em produção:
// `node src/workers/emailWorker.js`, escalado à parte.
const { Worker } = require("bullmq");
const { conexaoRedis } = require("../infrastructure/queue/redisConfig");

const modelos = {
  boasVindas: ({ usuario }) => ({
    assunto: `Bem-vindo(a), ${usuario.nome}!`,
    texto: `Sua conta em ${usuario.email} está pronta.`,
    html: `<h1>Olá, ${usuario.nome}</h1><p>Sua conta está pronta.</p>`,
  }),
  tarefaCriada: ({ tarefa }) => ({
    assunto: `Nova tarefa: ${tarefa.titulo}`,
    texto: `A tarefa "${tarefa.titulo}" foi criada.`,
    html: `<p>A tarefa <strong>${tarefa.titulo}</strong> foi criada.</p>`,
  }),
};

async function processarEmail(job, transporte) {
  const modelo = modelos[job.name];
  // Erro de programação não merece 3 retentativas: se o template não existe,
  // tentar de novo em 1 s, 2 s e 4 s só atrasa a descoberta.
  if (!modelo) throw new Error(`template desconhecido: ${job.name}`);
  const { assunto, texto, html } = modelo(job.data);
  const info = await transporte.sendMail({
    from: process.env.EMAIL_REMETENTE || "nao-responda@exemplo.com",
    to: job.data.para,
    subject: assunto,
    text: texto,
    html,
  });
  return { messageId: info.messageId, aceitos: info.accepted };
}

function criarEmailWorker({ transporte, conexao = conexaoRedis() }) {
  return new Worker("email", (job) => processarEmail(job, transporte), {
    connection: conexao,
    concurrency: 5, // SMTP aguenta paralelismo; banco de dados nem sempre
  });
}

module.exports = { criarEmailWorker, processarEmail, modelos };

// ---- src/workers/importacaoWorker.js
// 5 — IMPORTAÇÃO DE CSV, com progresso.
const { Worker } = require("bullmq");
const { conexaoRedis } = require("../infrastructure/queue/redisConfig");

// Parser mínimo, suficiente para o exercício: cabeçalho na 1a linha, vírgula
// como separador, sem campo com vírgula dentro. Em produção, use csv-parse.
function lerCsv(texto) {
  const linhas = texto.trim().split(/\r?\n/);
  const cabecalho = linhas[0].split(",").map((c) => c.trim());
  return linhas.slice(1).map((linha) => {
    const celulas = linha.split(",");
    return Object.fromEntries(cabecalho.map((c, i) => [c, (celulas[i] ?? "").trim()]));
  });
}

async function processarImportacao(job, { salvarProduto, aCada = 10 }) {
  const linhas = lerCsv(job.data.csv);
  const erros = [];
  let importados = 0;

  for (const [indice, linha] of linhas.entries()) {
    try {
      const preco = Number(linha.preco);
      if (!linha.nome) throw new Error("nome vazio");
      if (!Number.isFinite(preco) || preco <= 0) throw new Error(`preço inválido: "${linha.preco}"`);
      await salvarProduto({ nome: linha.nome, preco, sku: linha.sku });
      importados++;
    } catch (erro) {
      // uma linha ruim não invalida o arquivo: o relatório sai no fim
      erros.push({ linha: indice + 2, motivo: erro.message });
    }
    // updateParticipa do Redis: chamar a cada linha de um arquivo de 100 mil
    // seria uma escrita por linha. Reportar em blocos é o padrão.
    if ((indice + 1) % aCada === 0 || indice === linhas.length - 1) {
      await job.updateProgress(Math.round(((indice + 1) / linhas.length) * 100));
    }
  }
  return { total: linhas.length, importados, erros };
}

function criarImportacaoWorker({ salvarProduto, conexao = conexaoRedis() }) {
  return new Worker("importacao", (job) => processarImportacao(job, { salvarProduto }), {
    connection: conexao,
    concurrency: 1, // importação mexe no banco: uma de cada vez
  });
}

module.exports = { criarImportacaoWorker, processarImportacao, lerCsv };

// ---- src/rotas/importacao.js
const express = require("express");

// A rota NÃO importa nada: enfileira e devolve 202 com o id do job. O cliente
// acompanha depois. Fazer a importação dentro da requisição é o que produz o
// timeout de 30 s do proxy no meio do arquivo.
function criarRotaImportacao({ filaImportacao }) {
  const rotas = express.Router();

  rotas.post("/importacao/produtos", express.text({ type: "*/*", limit: "5mb" }), async (req, res) => {
    const csv = req.body;
    if (!csv || !csv.includes("nome")) {
      return res.status(422).json({ erro: "CSV precisa de cabeçalho com a coluna nome" });
    }
    const job = await filaImportacao.add("csv", { csv, usuarioId: req.usuario?.id ?? "anonimo" });
    res.status(202).json({ jobId: job.id, acompanhar: `/importacao/produtos/${job.id}` });
  });

  rotas.get("/importacao/produtos/:jobId", async (req, res) => {
    const job = await filaImportacao.getJob(req.params.jobId);
    if (!job) return res.status(404).json({ erro: "job não encontrado" });
    res.json({
      jobId: job.id,
      estado: await job.getState(),
      progresso: job.progress,
      resultado: job.returnvalue ?? null,
      erro: job.failedReason ?? null,
    });
  });

  return rotas;
}

module.exports = { criarRotaImportacao };

// ---- src/admin/bullBoard.js
// 6 — O PAINEL. Ele mostra dados de todo mundo e deixa reprocessar job:
// publicar sem proteção é entregar a fila inteira.
const { createBullBoard } = require("@bull-board/api");
const { BullMQAdapter } = require("@bull-board/api/bullMQAdapter");
const { ExpressAdapter } = require("@bull-board/express");

function protegerComChave(chaveEsperada) {
  return (req, res, proximo) => {
    const enviada = req.get("x-admin-key") || req.query.chave;
    // Sem chave configurada, o painel simplesmente não sobe. "Se a variável
    // estiver vazia, deixa passar" é como painéis vazam.
    if (!chaveEsperada) return res.status(503).json({ erro: "ADMIN_KEY não configurada" });
    if (enviada !== chaveEsperada) return res.status(401).json({ erro: "chave inválida" });
    proximo();
  };
}

function montarPainel(app, filas, { caminho = "/admin/filas", chave = process.env.ADMIN_KEY } = {}) {
  const adaptador = new ExpressAdapter();
  adaptador.setBasePath(caminho);
  createBullBoard({
    queues: Object.values(filas).map((fila) => new BullMQAdapter(fila)),
    serverAdapter: adaptador,
  });
  app.use(caminho, protegerComChave(chave), adaptador.getRouter());
  return adaptador;
}

module.exports = { montarPainel, protegerComChave };

// ---- src/jobs/recorrentes.js
// 7 — O JOB RECORRENTE.
// No BullMQ 6 o caminho é `upsertJobScheduler`. O antigo
// `add(nome, dados, { repeat })` ainda é aceito, mas `getRepeatableJobs()`
// deixou de existir — quem inspecionava por ali fica sem nada para ler.
const CHAVE = "verificar-tarefas-atrasadas";

async function agendarVerificacao(fila, { pattern = "* * * * *" } = {}) {
  return fila.upsertJobScheduler(
    CHAVE,
    { pattern }, // '* * * * *' = a cada virada de minuto
    { name: "verificar-atrasadas", data: { origem: "agendador" } }
  );
}

// upsert: rodar isto no boot de três réplicas não cria três agendamentos.
// Com o `repeat` antigo, cada boot criava um novo — e a caixa de entrada
// do usuário contava quantas réplicas havia no ar.
async function listarAgendamentos(fila) {
  return fila.getJobSchedulers();
}

async function cancelarVerificacao(fila) {
  return fila.removeJobScheduler(CHAVE);
}

module.exports = { agendarVerificacao, listarAgendamentos, cancelarVerificacao, CHAVE };

// ---- src/componentes/ImportacaoProdutos.jsx
// 8 — O ACOMPANHAMENTO no front-end.
// A API entra por parâmetro (mesma razão vista no artigo de testes: `import.meta.env` não passa
// pelo Jest, e um componente que chama fetch direto não se testa sem rede).
import { useEffect, useRef, useState } from "react";

const INTERVALO = 1000;

export function ImportacaoProdutos({ api, intervalo = INTERVALO }) {
  const [jobId, setJobId] = useState(null);
  const [estado, setEstado] = useState(null);
  const [progresso, setProgresso] = useState(0);
  const [resultado, setResultado] = useState(null);
  const [erro, setErro] = useState(null);
  const relogio = useRef(null);

  useEffect(() => () => clearInterval(relogio.current), []);

  async function enviar(evento) {
    evento.preventDefault();
    setErro(null);
    setResultado(null);
    const arquivo = evento.target.elements.arquivo.files[0];
    if (!arquivo) return setErro("escolha um arquivo .csv");
    try {
      const { jobId: id } = await api.enviarCsv(await arquivo.text());
      setJobId(id);
      setEstado("waiting");
      setProgresso(0);
    } catch (e) {
      setErro(e.message);
    }
  }

  useEffect(() => {
    if (!jobId) return;
    // Polling: simples e suficiente. Com WebSocket no ar, o mesmo componente
    // ouviria 'importacao:progresso' e não perguntaria nada.
    relogio.current = setInterval(async () => {
      const situacao = await api.consultar(jobId);
      setEstado(situacao.estado);
      setProgresso(situacao.progresso ?? 0);
      if (situacao.estado === "completed" || situacao.estado === "failed") {
        clearInterval(relogio.current);
        setResultado(situacao.resultado);
        if (situacao.erro) setErro(situacao.erro);
      }
    }, intervalo);
    return () => clearInterval(relogio.current);
  }, [jobId, api, intervalo]);

  const emAndamento = jobId && estado !== "completed" && estado !== "failed";

  return (
    <section>
      <h2>Importar produtos</h2>
      <form onSubmit={enviar}>
        <input type="file" name="arquivo" accept=".csv" aria-label="Arquivo CSV" />
        <button type="submit" disabled={Boolean(emAndamento)}>
          {emAndamento ? "Importando…" : "Importar"}
        </button>
      </form>

      {jobId && (
        <div>
          <progress value={progresso} max="100" aria-label="Progresso da importação" />
          <span>{progresso}%</span>
          <p>
            Job <code>{jobId}</code> — {estado}
          </p>
        </div>
      )}

      {resultado && (
        <div role="status">
          <p>
            {resultado.importados} de {resultado.total} produtos importados.
          </p>
          {resultado.erros?.length > 0 && (
            <ul>
              {resultado.erros.map((e) => (
                <li key={e.linha}>
                  Linha {e.linha}: {e.motivo}
                </li>
              ))}
            </ul>
          )}
        </div>
      )}

      {erro && <p role="alert">{erro}</p>}
    </section>
  );
}

// ---- tests/filas.test.js
const { EventEmitter } = require("node:events");
const express = require("express");
const request = require("supertest");
const nodemailer = require("nodemailer");
const { QueueEvents } = require("bullmq");
const { RedisMemoryServer } = require("redis-memory-server");

const { criarFilas } = require("../src/infrastructure/queue/filas");
const { criarProdutores, TAREFA_CRIADA, USUARIO_CADASTRADO } = require("../src/infrastructure/queue/produtores");
const { criarEmailWorker, processarEmail } = require("../src/workers/emailWorker");
const { criarImportacaoWorker, lerCsv } = require("../src/workers/importacaoWorker");
const { criarRotaImportacao } = require("../src/rotas/importacao");
const { montarPainel } = require("../src/admin/bullBoard");
const { agendarVerificacao, listarAgendamentos, cancelarVerificacao, CHAVE } = require("../src/jobs/recorrentes");

jest.setTimeout(120000);

let redis, conexao, filas, produtores, workers, ouvintes, filasAvulsas;

const esperarVazia = async (fila, ms = 5000) => {
  const limite = Date.now() + ms;
  while (Date.now() < limite) {
    const [ativos, esperando, atrasados] = await Promise.all([
      fila.getActiveCount(),
      fila.getWaitingCount(),
      fila.getDelayedCount(),
    ]);
    if (ativos + esperando + atrasados === 0) return;
    await new Promise((ok) => setTimeout(ok, 50));
  }
  throw new Error("a fila não esvaziou");
};

beforeAll(async () => {
  redis = new RedisMemoryServer();
  conexao = { host: await redis.getHost(), port: await redis.getPort(), maxRetriesPerRequest: null };
  filas = criarFilas(conexao);
  produtores = criarProdutores(filas);
  workers = [];
  ouvintes = [];
  filasAvulsas = [];
});

// Toda QueueEvents abre a PRÓPRIA conexão com o Redis. Esquecer de fechar uma
// delas enche a saída de ECONNREFUSED quando o servidor cai no fim da suíte —
// e o Jest fica pendurado esperando o socket.
const escutar = (nome) => {
  const qe = new QueueEvents(nome, { connection: conexao });
  ouvintes.push(qe);
  return qe;
};

afterAll(async () => {
  await Promise.all(workers.map((w) => w.close()));
  await Promise.all(ouvintes.map((o) => o.close()));
  await Promise.all(filasAvulsas.map((f) => f.close()));
  await Promise.all(Object.values(filas).map((f) => f.close()));
  await redis.stop();
});

describe("3 · o worker de email", () => {
  test("os dois templates montam assunto, texto e HTML a partir dos dados", async () => {
    const enviados = [];
    const transporte = {
      sendMail: async (m) => (enviados.push(m), { messageId: "id-1", accepted: [m.to] }),
    };
    await processarEmail(
      { name: "boasVindas", data: { para: "novo@exemplo.com", usuario: { nome: "Ana", email: "novo@exemplo.com" } } },
      transporte
    );
    await processarEmail(
      { name: "tarefaCriada", data: { para: "dono@exemplo.com", tarefa: { titulo: "Revisar PR" } } },
      transporte
    );
    expect(enviados.map((e) => e.subject)).toEqual(["Bem-vindo(a), Ana!", "Nova tarefa: Revisar PR"]);
    expect(enviados[0].html).toContain("<h1>Olá, Ana</h1>");
    expect(enviados[1].text).toBe('A tarefa "Revisar PR" foi criada.');
  });

  test("template inexistente falha de imediato", async () => {
    await expect(processarEmail({ name: "naoExiste", data: {} }, { sendMail: jest.fn() })).rejects.toThrow(
      "template desconhecido: naoExiste"
    );
  });

  test("o job atravessa o Redis de verdade e volta com o messageId", async () => {
    const transporte = { sendMail: async (m) => ({ messageId: `<${m.subject}>`, accepted: [m.to] }) };
    const worker = criarEmailWorker({ transporte, conexao });
    workers.push(worker);
    const job = await filas.filaEmail.add("boasVindas", {
      para: "ana@exemplo.com",
      usuario: { nome: "Ana", email: "ana@exemplo.com" },
    });
    const concluido = await job.waitUntilFinished(escutar("email"));
    expect(concluido).toEqual({ messageId: "<Bem-vindo(a), Ana!>", aceitos: ["ana@exemplo.com"] });
  });
});

describe("4 · jobId único contra email duplicado", () => {
  test("o mesmo evento disparado duas vezes enfileira UM job", async () => {
    const barramento = new EventEmitter();
    produtores.ligarAoBarramento(barramento);
    const tarefa = { id: "t-42", titulo: "Duplicar?", emailDono: "dono@exemplo.com" };

    barramento.emit(TAREFA_CRIADA, tarefa);
    barramento.emit(TAREFA_CRIADA, tarefa);
    await new Promise((ok) => setTimeout(ok, 200));

    const job = await filas.filaEmail.getJob("email:tarefaCriada:t-42");
    expect(job).toBeTruthy();
    barramento.removeAllListeners();
  });

  test("o retorno do segundo add MENTE: mostra o dado novo, o Redis guarda o velho", async () => {
    const primeiro = await filas.filaNotificacao.add("push", { tarefaId: "x1", v: "PRIMEIRO" }, { jobId: "dedup" });
    const segundo = await filas.filaNotificacao.add("push", { tarefaId: "x1", v: "SEGUNDO" }, { jobId: "dedup" });
    expect(segundo.data.v).toBe("SEGUNDO"); // o objeto devolvido
    expect((await filas.filaNotificacao.getJob("dedup")).data.v).toBe("PRIMEIRO"); // o que existe
    expect(segundo.id).toBe(primeiro.id);
    await (await filas.filaNotificacao.getJob("dedup")).remove();
  });

  test("removeOnComplete DESLIGA a deduplicação — a armadilha que ninguém vê", async () => {
    let vezes = 0;
    const { Worker } = require("bullmq");
    const w = new Worker("dedup-teste", async () => void vezes++, { connection: conexao });
    workers.push(w);
    const { Queue } = require("bullmq");
    const q = new Queue("dedup-teste", { connection: conexao });
    filasAvulsas.push(q);

    await q.add("j", {}, { jobId: "mesmo", removeOnComplete: true });
    await esperarVazia(q);
    await q.add("j", {}, { jobId: "mesmo", removeOnComplete: true });
    await esperarVazia(q);
    expect(vezes).toBe(2); // o registro sumiu, e com ele a memória do jobId

    vezes = 0;
    await q.add("k", {}, { jobId: "guardado" });
    await esperarVazia(q);
    await q.add("k", {}, { jobId: "guardado" });
    await esperarVazia(q);
    expect(vezes).toBe(1); // aqui o job concluído continua no Redis
  });
});

describe("retentativa e backoff", () => {
  test("falha 3 vezes com espera crescente e termina em failed", async () => {
    const { Worker, Queue } = require("bullmq");
    const tentativas = [];
    const w = new Worker(
      "falha",
      async () => {
        tentativas.push(Date.now());
        throw new Error("SMTP recusou");
      },
      { connection: conexao }
    );
    workers.push(w);
    const q = new Queue("falha", { connection: conexao });
    filasAvulsas.push(q);
    const job = await q.add("t", {}, { attempts: 3, backoff: { type: "exponential", delay: 100 } });
    await new Promise((ok) => setTimeout(ok, 2500));

    expect(tentativas).toHaveLength(3);
    const intervalos = tentativas.slice(1).map((t, i) => t - tentativas[i]);
    expect(intervalos[0]).toBeGreaterThanOrEqual(100);
    expect(intervalos[1]).toBeGreaterThan(intervalos[0]); // exponencial: ~100, ~200
    expect(await job.getState()).toBe("failed");
    expect((await q.getFailed())[0].failedReason).toBe("SMTP recusou");
  });
});

describe("5 · importação de CSV com progresso", () => {
  const csv = (linhas) =>
    ["nome,preco,sku", ...linhas].join("\n");

  test("o parser lê cabeçalho e linhas", () => {
    expect(lerCsv(csv(["Teclado,349.90,TEC-1"]))).toEqual([
      { nome: "Teclado", preco: "349.90", sku: "TEC-1" },
    ]);
  });

  test("100 produtos entram, o progresso sobe até 100 e as linhas ruins viram relatório", async () => {
    const salvos = [];
    const worker = criarImportacaoWorker({ salvarProduto: async (p) => void salvos.push(p), conexao });
    workers.push(worker);

    const linhas = Array.from({ length: 100 }, (_, i) => `Produto ${i + 1},${(i + 1) * 1.5},SKU-${i + 1}`);
    linhas[9] = ",10,SEM-NOME";
    linhas[19] = "Preço ruim,de graça,SKU-20";

    const eventos = escutar("importacao");
    await eventos.waitUntilReady();
    const progressos = [];
    eventos.on("progress", ({ data }) => progressos.push(data));

    const app = express();
    app.use(criarRotaImportacao({ filaImportacao: filas.filaImportacao }));

    const envio = await request(app)
      .post("/importacao/produtos")
      .set("content-type", "text/csv")
      .send(csv(linhas))
      .expect(202);
    expect(envio.body.jobId).toBeTruthy();

    const job = await filas.filaImportacao.getJob(envio.body.jobId);
    const resultado = await job.waitUntilFinished(eventos);

    expect(resultado).toEqual({
      total: 100,
      importados: 98,
      erros: [
        { linha: 11, motivo: "nome vazio" },
        { linha: 21, motivo: 'preço inválido: "de graça"' },
      ],
    });
    expect(salvos).toHaveLength(98);
    expect(progressos.at(-1)).toBe(100);
    expect(progressos.length).toBeGreaterThan(1); // subiu em blocos, não de uma vez

    const consulta = await request(app).get(`/importacao/produtos/${envio.body.jobId}`).expect(200);
    expect(consulta.body).toMatchObject({ estado: "completed", progresso: 100 });
  });

  test("CSV sem cabeçalho válido é recusado na porta, sem enfileirar nada", async () => {
    const app = express();
    app.use(criarRotaImportacao({ filaImportacao: filas.filaImportacao }));
    const antes = await filas.filaImportacao.getJobCounts();
    await request(app).post("/importacao/produtos").set("content-type", "text/csv").send("a,b,c\n1,2,3").expect(422);
    expect(await filas.filaImportacao.getJobCounts()).toEqual(antes);
  });

  test("job inexistente devolve 404", async () => {
    const app = express();
    app.use(criarRotaImportacao({ filaImportacao: filas.filaImportacao }));
    await request(app).get("/importacao/produtos/9999").expect(404);
  });
});

describe("6 · Bull Board protegido por ADMIN_KEY", () => {
  const montar = (chave) => {
    const app = express();
    montarPainel(app, { email: filas.filaEmail }, { chave });
    return app;
  };

  test("sem chave na requisição: 401", async () => {
    await request(montar("segredo")).get("/admin/filas").expect(401);
  });

  test("chave errada: 401; chave certa: passa", async () => {
    const app = montar("segredo");
    await request(app).get("/admin/filas").set("x-admin-key", "chute").expect(401);
    const ok = await request(app).get("/admin/filas").set("x-admin-key", "segredo");
    expect(ok.status).toBeLessThan(400);
  });

  test("ADMIN_KEY vazia não libera o painel — devolve 503", async () => {
    await request(montar("")).get("/admin/filas").set("x-admin-key", "").expect(503);
    await request(montar(undefined)).get("/admin/filas").expect(503);
  });

  test("a API do painel enxerga a fila registrada", async () => {
    const app = montar("segredo");
    const r = await request(app).get("/admin/filas/api/queues").set("x-admin-key", "segredo");
    expect(r.status).toBe(200);
    expect(JSON.stringify(r.body)).toContain("email");
  });
});

describe("7 · job recorrente", () => {
  afterEach(() => cancelarVerificacao(filas.filaNotificacao));

  test("agenda com o pattern de minuto e aparece na lista de schedulers", async () => {
    await agendarVerificacao(filas.filaNotificacao, { pattern: "* * * * *" });
    const agendados = await listarAgendamentos(filas.filaNotificacao);
    expect(agendados).toEqual([
      expect.objectContaining({ key: CHAVE, pattern: "* * * * *", name: "verificar-atrasadas" }),
    ]);
    expect(agendados[0].next).toBeGreaterThan(Date.now());
  });

  test("upsert: chamar duas vezes (duas réplicas subindo) não duplica o agendamento", async () => {
    await agendarVerificacao(filas.filaNotificacao);
    await agendarVerificacao(filas.filaNotificacao);
    await agendarVerificacao(filas.filaNotificacao);
    expect(await listarAgendamentos(filas.filaNotificacao)).toHaveLength(1);
  });

  test("getRepeatableJobs sumiu no BullMQ 6 — quem inspecionava por ali quebrou", () => {
    expect(filas.filaNotificacao.getRepeatableJobs).toBeUndefined();
    expect(typeof filas.filaNotificacao.getJobSchedulers).toBe("function");
  });
});

describe("3 · envio real por SMTP (Ethereal)", () => {
  test("o email sai de verdade e o servidor devolve messageId e URL de leitura", async () => {
    let conta;
    try {
      conta = await nodemailer.createTestAccount();
    } catch {
      console.warn("Ethereal indisponível — teste de SMTP real pulado");
      return;
    }
    const transporte = nodemailer.createTransport({
      host: conta.smtp.host,
      port: conta.smtp.port,
      secure: conta.smtp.secure,
      auth: { user: conta.user, pass: conta.pass },
    });
    const info = await processarEmail(
      { name: "boasVindas", data: { para: conta.user, usuario: { nome: "Ana", email: conta.user } } },
      transporte
    );
    expect(info.messageId).toMatch(/@/);
    expect(info.aceitos).toContain(conta.user);
    console.log(`SMTP real: ${info.messageId} entregue na caixa ${conta.user}`);
  });
});

// ---- tests/ImportacaoProdutos.test.jsx
/**
 * @jest-environment jsdom
 */
import { act, render, screen, waitFor } from "@testing-library/react";
import userEvent from "@testing-library/user-event";
import { ImportacaoProdutos } from "../src/componentes/ImportacaoProdutos.jsx";

// O jsdom implementa File e FileReader, mas NÃO o Blob.text() — que existe em
// todo navegador desde 2019. Sem este polyfill de três linhas, o componente
// morre com "arquivo.text is not a function" e o erro parece ser do código.
if (typeof File !== "undefined" && !File.prototype.text) {
  File.prototype.text = function () {
    return new Promise((ok, falha) => {
      const leitor = new FileReader();
      leitor.onload = () => ok(leitor.result);
      leitor.onerror = () => falha(leitor.error);
      leitor.readAsText(this);
    });
  };
}

const arquivoCsv = (conteudo) =>
  new File([conteudo], "produtos.csv", { type: "text/csv" });

function apiFalsa(situacoes) {
  const fila = [...situacoes];
  return {
    enviarCsv: jest.fn().mockResolvedValue({ jobId: "77" }),
    consultar: jest.fn().mockImplementation(async () => fila.shift() ?? situacoes.at(-1)),
  };
}

test("8 · envia o CSV, acompanha o progresso e mostra o relatório final", async () => {
  jest.useFakeTimers();
  const usuario = userEvent.setup({ advanceTimers: jest.advanceTimersByTime });
  const api = apiFalsa([
    { estado: "active", progresso: 30 },
    { estado: "active", progresso: 80 },
    {
      estado: "completed",
      progresso: 100,
      resultado: { total: 100, importados: 98, erros: [{ linha: 11, motivo: "nome vazio" }] },
    },
  ]);

  render(<ImportacaoProdutos api={api} intervalo={1000} />);
  await usuario.upload(screen.getByLabelText("Arquivo CSV"), arquivoCsv("nome,preco\nA,1"));
  await usuario.click(screen.getByRole("button", { name: "Importar" }));

  await waitFor(() => expect(api.enviarCsv).toHaveBeenCalledWith("nome,preco\nA,1"));
  expect(screen.getByRole("button", { name: "Importando…" })).toBeDisabled();

  await act(async () => {
    jest.advanceTimersByTime(1000);
  });
  expect(screen.getByLabelText("Progresso da importação")).toHaveValue(30);

  await act(async () => {
    jest.advanceTimersByTime(2000);
  });
  expect(screen.getByRole("status")).toHaveTextContent("98 de 100 produtos importados");
  expect(screen.getByText("Linha 11: nome vazio")).toBeInTheDocument();
  expect(screen.getByRole("button", { name: "Importar" })).toBeEnabled();

  // o polling PARA quando o job acaba — senão o navegador pergunta para sempre
  const chamadasAteAqui = api.consultar.mock.calls.length;
  await act(async () => {
    jest.advanceTimersByTime(5000);
  });
  expect(api.consultar).toHaveBeenCalledTimes(chamadasAteAqui);
  jest.useRealTimers();
});

test("8 · sem arquivo, avisa e não chama a API", async () => {
  const usuario = userEvent.setup();
  const api = apiFalsa([]);
  render(<ImportacaoProdutos api={api} />);
  await usuario.click(screen.getByRole("button", { name: "Importar" }));
  expect(screen.getByRole("alert")).toHaveTextContent("escolha um arquivo .csv");
  expect(api.enviarCsv).not.toHaveBeenCalled();
});

test("8 · erro do servidor aparece para o usuário", async () => {
  const usuario = userEvent.setup();
  const api = apiFalsa([]);
  api.enviarCsv.mockRejectedValue(new Error("CSV precisa de cabeçalho com a coluna nome"));
  render(<ImportacaoProdutos api={api} />);
  await usuario.upload(screen.getByLabelText("Arquivo CSV"), arquivoCsv("a,b\n1,2"));
  await usuario.click(screen.getByRole("button", { name: "Importar" }));
  expect(await screen.findByRole("alert")).toHaveTextContent("CSV precisa de cabeçalho");
});

test("8 · desmontar durante a importação não deixa o intervalo rodando", async () => {
  jest.useFakeTimers();
  const usuario = userEvent.setup({ advanceTimers: jest.advanceTimersByTime });
  const api = apiFalsa([{ estado: "active", progresso: 10 }]);
  const { unmount } = render(<ImportacaoProdutos api={api} intervalo={1000} />);
  await usuario.upload(screen.getByLabelText("Arquivo CSV"), arquivoCsv("nome,preco\nA,1"));
  await usuario.click(screen.getByRole("button", { name: "Importar" }));
  await waitFor(() => expect(api.enviarCsv).toHaveBeenCalled());

  unmount();
  const antes = api.consultar.mock.calls.length;
  await act(async () => {
    jest.advanceTimersByTime(10000);
  });
  expect(api.consultar).toHaveBeenCalledTimes(antes);
  jest.useRealTimers();
});

O item 4 pede "jobId único para evitar emails duplicados", e é aí que mora a armadilha mais cara do artigo: removeOnComplete desliga a deduplicação. O BullMQ lembra do jobId enquanto o registro do job existe no Redis; ao remover o job concluído — que é o que todo mundo configura, e com razão, para o Redis não crescer sem fim — a memória vai junto, e o mesmo jobId roda de novo. O teste mede as duas situações lado a lado: com removeOnComplete, dois processamentos; sem, um. Se o e-mail duplicado importa de verdade, a trava tem de ser sua (uma chave no banco), não do broker. Duas outras: o objeto devolvido pelo segundo add() mente — ele mostra os dados novos, mas o que está no Redis é o job antigo; e no BullMQ 6 o getRepeatableJobs() não existe mais — o repeat antigo ainda é aceito no add, mas quem inspeciona agendamentos usa upsertJobScheduler e getJobSchedulers(). O upsert, aliás, é o que evita três réplicas criarem três agendamentos do mesmo job de minuto.

A fila separa duas coisas que costumam vir grudadas: aceitar o pedido e executá-lo. A API confirma em milissegundos, o trabalho pesado corre em outro processo, e ninguém fica olhando para uma requisição pendurada. O preço é que tudo passa a ser mesmo assíncrono — o cliente precisa de alguma forma de acompanhar o progresso, o job precisa ser seguro para repetir, porque ele será repetido, e a fila vira mais um sistema a monitorar.

Fontes e Referências

Exercícios

Exercício 1

O job cobra o cartão do cliente. A fila está configurada com attempts: 3. Em que situação o cliente é cobrado duas vezes?

new Worker('pagamentos', async (job) => {
  const { pedidoId, valor } = job.data;
  await gateway.cobrar(pedidoId, valor);
  await Pedido.findByIdAndUpdate(pedidoId, { status: 'pago' });
}, { connection: redisConfig });
Ver resposta

✓ Resposta: Sempre que a cobrança funcionar e algo depois dela falhar. Se o gateway processar o pagamento mas a resposta se perder por timeout, ou se o processo morrer entre a cobrança e o findByIdAndUpdate, o job é considerado falho e o BullMQ o executa de novo — cobrando outra vez. A raiz disso é uma propriedade que vale internalizar: filas garantem entrega pelo menos uma vez, não exatamente uma vez. Reprocessamento não é exceção, é parte do funcionamento normal, e todo handler precisa ser escrito para tolerá-lo — é o que se chama de idempotência. Há duas formas de obtê-la aqui. A primeira é verificar o estado antes de agir: se o pedido já está pago, encerrar sem fazer nada. A segunda, mais confiável, é usar uma chave de idempotência — um identificador estável derivado do pedido, enviado ao gateway, que então reconhece a repetição e devolve o resultado da primeira cobrança em vez de criar uma nova. Todo gateway sério suporta isso justamente por causa desse cenário. E vale a inversão de ordem quando possível: registrar a intenção antes da chamada externa dá ao reprocessamento algo em que se apoiar. A regra curta é que operações com efeito no mundo real — cobrar, enviar email, despachar mercadoria — nunca devem depender de o job rodar uma vez só.

Exercício 2

O job recebe o objeto inteiro. Duas coisas dão errado com o tempo. Quais?

await filaEmail.add('boas-vindas', {
  usuario: usuarioCompleto,   // documento inteiro do Mongo
  template: templateHtml,     // ~40 KB de HTML
});
Ver resposta

✓ Resposta: A primeira é o dado obsoleto: o payload é uma fotografia do momento em que o job foi enfileirado, e ele pode ficar minutos ou horas na fila. Se o usuário trocar de email nesse intervalo, a mensagem vai para o endereço antigo; se a conta for desativada, o email sai assim mesmo. A segunda é o consumo de memória do Redis: todo o payload é serializado e guardado ali, e o Redis mantém os dados em RAM. Quarenta quilobytes por job, com dez mil jobs por dia, são centenas de megabytes de conteúdo que já existe em outro lugar. Em escala, isso leva o Redis ao limite de memória e ele começa a evictar chaves — inclusive as da própria fila. A correção é passar apenas identificadores: { usuarioId, tipo: 'boas-vindas' }, e o worker busca o que precisa no momento em que for processar, obtendo de quebra os dados atualizados. O template não deveria trafegar de forma alguma: ele é código, e o worker já o tem. A regra prática é que payload de job é um bilhete, não um pacote — quando algo grande precisa mesmo ser processado, como um arquivo enviado, guarde-o em disco ou em armazenamento de objetos e passe o caminho.

Exercício 3

Tudo funciona por três semanas. Então o Redis começa a recusar escritas e a fila inteira para. O que faltou na configuração?

const filaEmail = new Queue('email', {
  connection: redisConfig,
  defaultJobOptions: {
    attempts: 3,
    backoff: { type: 'exponential', delay: 2000 },
  },
});
Ver resposta

✓ Resposta: Faltaram removeOnComplete e removeOnFail. Por padrão o BullMQ mantém os jobs concluídos, com todo o payload e o histórico de tentativas, para que você possa inspecioná-los. É útil em desenvolvimento e é uma bomba-relógio em produção: nada os apaga, o Redis guarda tudo em memória, e três semanas depois a instância bate no limite configurado. O que acontece então depende da política de eviction — ou o Redis passa a recusar escritas, e nenhum job novo entra, ou ele começa a descartar chaves e a fila perde trabalho. O sintoma é cruel porque o sistema funcionou perfeitamente por semanas e quebra sem que nada tenha mudado no código. A configuração saudável limita por idade e por quantidade: removeOnComplete: { age: 86400, count: 1000 } guarda um dia ou os mil últimos, o que vier primeiro. Para as falhas costuma-se guardar mais tempo, porque é delas que se precisa para investigar. Duas dicas que acompanham: monitore o uso de memória do Redis desde o primeiro dia, porque esse é o recurso que acaba antes; e defina maxmemory-policy explicitamente — com noeviction, o Redis recusa escritas em vez de descartar em silêncio, o que é bem melhor do que perder jobs sem ninguém saber.

Exercício 4

O job de gerar relatório leva cerca de 2 minutos. Ele é processado, termina — e é executado de novo. E de novo. Por quê?

new Worker('relatorios', async (job) => {
  const dados = await consultarTudo();     // ~40s
  const pdf = await gerarPDF(dados);       // ~80s
  await salvar(pdf);
}, { connection: redisConfig });
Ver resposta

✓ Resposta: O job está sendo considerado travado (stalled). Para saber se um worker morreu no meio do processamento, o BullMQ trabalha com uma trava renovada periodicamente; o padrão do lockDuration é de 30 segundos, e se a trava não for renovada nesse intervalo a fila conclui que o worker caiu e entrega o job a outro. Um processamento de dois minutos que não renova a trava estoura esse limite, é reenfileirado enquanto ainda está rodando, e aí passam a existir duas execuções simultâneas do mesmo relatório — que acabam, são reenfileiradas de novo, e o ciclo se repete. Há duas correções, e a boa prática é usar as duas. A primeira é aumentar o lockDuration para algo confortavelmente acima do tempo real do job. A segunda, melhor, é renovar a trava durante o trabalho, chamando job.updateProgress() de tempo em tempo — o que renova a trava e, de quebra, alimenta a barra de progresso do cliente. A dica geral que fica: job longo precisa dar sinal de vida; do ponto de vista da fila, um worker silencioso por muito tempo é indistinguível de um worker morto. E, quando o processamento passa de alguns minutos, costuma valer quebrá-lo em jobs menores e encadeados, que são mais fáceis de retomar quando algo falha.

Exercício 5

Um deploy trocou o formato do payload. Durante alguns minutos, jobs falham com Cannot read properties of undefined. O que aconteceu, e como evitar?

// versão antiga, ainda na fila
{ usuarioId: '123', tipo: 'boas-vindas' }

// versão nova, que o worker novo espera
{ destinatario: { id: '123' }, template: 'boas-vindas' }
Ver resposta

✓ Resposta: A fila tinha jobs do formato antigo esperando quando o worker novo entrou no ar. Diferente de uma requisição HTTP, que nasce e morre em milissegundos, um job pode ficar minutos ou horas parado — a fila é, na prática, um buffer que atravessa deploys. O worker novo lê job.data.destinatario.id num payload que só tem usuarioId, e falha. Pior: com attempts configurado, ele tenta várias vezes antes de desistir, enchendo o log de erros idênticos. A prevenção é a mesma de mudança de API, aplicada aqui: mudança de formato de payload é mudança incompatível e precisa de transição. O caminho seguro tem três etapas — o worker passa a aceitar os dois formatos, o produtor passa a emitir o novo, e só quando não houver mais nada do formato antigo na fila o suporte ao velho é removido. Uma alternativa mais explícita é versionar o payload, com um campo versao: 2, o que torna o tratamento legível em vez de adivinhado. E há o caso em que nada disso é possível: então drene a fila antes de publicar — pare de produzir, espere esvaziar, então faça o deploy. Para filas grandes isso costuma ser inviável, e é por isso que compatibilidade retroativa acaba sendo a regra.

Comentários

Mais em Javascript

ESLint e Prettier: código limpo e padronizado
ESLint e Prettier: código limpo e padronizado

Entrar num projeto em que cada arquivo tem um estilo custa caro: revisão que…

Performance em aplicações web
Performance em aplicações web

É fácil otimizar a coisa errada, e por isso a ordem importa: medir antes. O…

Dominando o JavaScript
Dominando o JavaScript

Cinquenta e dois artigos em seis módulos, do primeiro console.log ao dashboard…