Esta es la lección en la que el diseño de los módulos 2 y 3 se convierte en un servicio que funciona de punta a punta. Construimos el corazón de servicio-pedidos (puerto 3002, equipo de Luis) sobre la misma plantilla que servicio-catalogo, pero con todo lo que Catálogo no tenía: PostgreSQL con transacciones, llamadas HTTP salientes a Catálogo y Clientes con su timeout y su ACL, el caso de uso crearPedido que guarda el pedido y su evento pedido.creado en la misma transacción (outbox), el relay que publica en RabbitMQ con confirmaciones, y los consumidores que reciben los eventos de la saga y mueven el pedido por su máquina de estados hasta CONFIRMADO o CANCELADO. Terminamos con una prueba manual completa: crear el pedido de Ana con curl, verlo PENDIENTE, inyectar un stock.reservado y verlo STOCK_RESERVADO. Mostramos completo lo esencial y resumimos lo repetitivo; los ficheros omitidos son variantes directas de lo que ya viste en 04-02.

Contenido

  1. Estructura del proyecto y dependencias
  2. PostgreSQL: pool, transacciones y migraciones
  3. Clientes HTTP salientes: Catálogo (con ACL) y Clientes
  4. Dominio: el agregado Pedido y la máquina de estados
  5. Repositorio con outbox: guardarConEventos() y procesarUnaVez()
  6. El caso de uso crearPedido y la ruta POST /v1/pedidos
  7. El relay del outbox
  8. Consumidores: la saga y la réplica de clientes
  9. GET /v1/pedidos/{id} con ETag y el arranque completo
  10. Prueba de punta a punta

  1. Estructura del proyecto y dependencias

npm install express pg amqplib pino pino-http zod @techcorp/comun-http && npm install -D nodemon dotenv
servicio-pedidos/
├── src/
│   ├── servidor.js, app.js, config.js (04-03), salud.js
│   ├── rutas/pedidos.js                    # POST /v1/pedidos, GET /v1/pedidos/:id
│   ├── casos-uso/crearPedido.js
│   ├── dominio/pedido.js                   # agregado: crear, calcular total, aplicar evento
│   ├── dominio/maquinaEstadosPedido.js     # TRANSICIONES, MOTIVOS, transicionar (02-05, tal cual)
│   ├── repositorios/pedidoRepositorio.js   # guardarConEventos, obtener, claves de idempotencia, procesarUnaVez
│   ├── clientes/catalogoCliente.js, clientes/clientesCliente.js
│   ├── traductores/traductorProducto.js    # ACL (02-03, 03-06)
│   ├── infra/postgres.js
│   └── mensajeria/relayOutbox.js, consumidorSaga.js, consumidorClientes.js
├── migraciones/001-esquema-inicial.sql … 004-claves-idempotencia.sql
├── scripts/migrar.js, scripts/publicarEvento.js
└── contratos/openapi.yaml, contratos/asyncapi.yaml

Respecto a Catálogo cambian el driver (pg en lugar de mongodb), aparece amqplib, y hay tres carpetas nuevas: casos-uso/ (la lógica de aplicación es más rica que en Catálogo), dominio/ y mensajeria/. mensajeria/topologia.js y publicador.js de 03-02 vienen de @techcorp/comun-http.

  1. PostgreSQL: pool, transacciones y migraciones

// src/infra/postgres.js
const { Pool } = require('pg');

function crearPoolPostgres({ url, logger }) {
  const pool = new Pool({ connectionString: url, max: 10, connectionTimeoutMillis: 5000, idleTimeoutMillis: 30000 });
  pool.on('error', (err) => logger.error({ err }, 'error en conexión inactiva del pool'));   // sin esto, un error tumba el proceso
  return {
    consultar: (sql, params) => pool.query(sql, params),                 // fuera de transacción
    // Ejecuta fn(tx) dentro de BEGIN/COMMIT; ROLLBACK si fn lanza. tx.consultar usa SIEMPRE la misma conexión.
    async transaccion(fn) {
      const cliente = await pool.connect();
      try {
        await cliente.query('BEGIN');
        const resultado = await fn({ consultar: (sql, params) => cliente.query(sql, params) });
        await cliente.query('COMMIT');
        return resultado;
      } catch (err) { await cliente.query('ROLLBACK'); throw err; }
      finally { cliente.release(); }
    },
    ping: () => pool.query('SELECT 1'),
    cerrar: () => pool.end()
  };
}
module.exports = { crearPoolPostgres };

transaccion(fn) es la pieza que hace posible el outbox: todo lo que se ejecute con el tx recibido va en la misma transacción, o se confirma entero o no se confirma nada.

Las migraciones son los esquemas ya diseñados, sin cambios de fondo: 001-esquema-inicial.sql (pedidos, lineas_pedido, clientes_ref de 02-04 §7.1), 002-outbox.sql (outbox con su índice parcial de pendientes, 02-05 §7), 003-eventos-procesados.sql (eventos_procesados con clave primaria compuesta (evento_id, consumidor)), 004-claves-idempotencia.sql (claves_idempotencia de 02-05 §8, a la que añadimos una columna huella TEXT NOT NULL con el hash del cuerpo para detectar la reutilización de clave con otro cuerpo, 03-01). scripts/migrar.js (~30 líneas) las aplica en orden y anota cada una en migraciones_aplicadas para no repetirla; se ejecuta con npm run migrar antes de arrancar (en Kubernetes será un Job o un init container, 05-02).

  1. Clientes HTTP salientes: Catálogo (con ACL) y Clientes

El cliente de Catálogo evoluciona el de 03-01: recibe la URL y el timeout por parámetro (04-03), llama a /v1/productos?ids= y aplica el ACL traductorProducto (03-06) antes de devolver, de modo que el resto de Pedidos nunca ve el JSON de Catálogo.

// src/traductores/traductorProducto.js — el lector tolerante de 03-06, sin cambios
function aProductoDePedidos(dto) {
  return { productoId: dto.id, nombre: dto.nombre, precioUnitario: Number(dto.precio), disponible: dto.disponible !== false };
}
module.exports = { aProductoDePedidos };
// src/clientes/catalogoCliente.js
const { ErrorNegocio } = require('@techcorp/comun-http');
const { aProductoDePedidos } = require('../traductores/traductorProducto');

function crearCatalogoCliente({ urlBase, timeoutMs }) {
  return {
    // Devuelve un Map productoId → { productoId, nombre, precioUnitario, disponible }
    async obtenerProductos(ids, { requestId }) {
      const url = `${urlBase}/v1/productos?ids=${encodeURIComponent(ids.join(','))}`;
      let respuesta;
      try {
        respuesta = await fetch(url, { headers: { Accept: 'application/json', 'X-Request-Id': requestId }, signal: AbortSignal.timeout(timeoutMs) });
      } catch (err) {                                     // timeout o red: 503 para nuestro cliente (03-01)
        throw new ErrorNegocio('DEPENDENCIA_NO_DISPONIBLE', `Catálogo no disponible: ${err.name}`, 503);
      }
      if (!respuesta.ok) {
        const problema = await respuesta.json().catch(() => ({}));
        if (respuesta.status >= 500) throw new ErrorNegocio('DEPENDENCIA_NO_DISPONIBLE', `Catálogo respondió ${respuesta.status}`, 503);
        throw new ErrorNegocio('PETICION_INVALIDA', `Catálogo rechazó la petición (${problema.codigo ?? respuesta.status})`, 400);
      }
      const { datos, noEncontrados } = await respuesta.json();
      const productos = datos.map(aProductoDePedidos);
      const noVendibles = [...noEncontrados, ...productos.filter((p) => !p.disponible).map((p) => p.productoId)];
      if (noVendibles.length > 0) throw new ErrorNegocio('PRODUCTO_NO_DISPONIBLE', `Productos no disponibles: ${noVendibles.join(', ')}`, 422);
      return new Map(productos.map((p) => [p.productoId, p]));
    }
  };
}
module.exports = { crearCatalogoCliente };

El cliente de Clientes es simétrico y más corto: GET {CLIENTES_URL}/v1/clientes/{id} con el mismo timeout; 404ErrorNegocio('CLIENTE_NO_EXISTE', ..., 404); 5xx/red → DEPENDENCIA_NO_DISPONIBLE; devuelve { clienteId, nombre, email, direcciones }. Sin reintentos ni circuit breaker (06-03): solo timeout.

  1. Dominio: el agregado Pedido y la máquina de estados

dominio/maquinaEstadosPedido.js es literalmente el de 02-05 (TRANSICIONES, MOTIVOS, transicionar). El agregado añade la construcción y el cálculo del total con céntimos enteros para evitar 59.90 + 9.90 * 2 = 79.69999...:

// src/dominio/pedido.js
const { randomUUID } = require('node:crypto');
const { transicionar, MOTIVOS } = require('./maquinaEstadosPedido');

const aCentimos = (n) => Math.round(n * 100);

// Construye un Pedido PENDIENTE a partir de la petición validada y de lo que dijeron Clientes y Catálogo
function crearPedido({ clienteId, lineas, direccionEnvio }, { cliente, productos }) {
  const lineasCongeladas = lineas.map((l, i) => {
    const p = productos.get(l.productoId);            // catalogoCliente ya garantizó que existe y es vendible
    return { linea: i + 1, productoId: p.productoId, nombreProducto: p.nombre, precioUnitario: p.precioUnitario, cantidad: l.cantidad };
  });
  const totalCentimos = lineasCongeladas.reduce((acc, l) => acc + aCentimos(l.precioUnitario) * l.cantidad, 0);
  return {
    pedidoId: `ped-${randomUUID().slice(0, 8)}`,       // id opaco generado por el dueño (02-04)
    clienteId, cliente: { nombre: cliente.nombre, email: cliente.email },
    estado: 'PENDIENTE', lineas: lineasCongeladas, direccionEnvio,
    total: totalCentimos / 100, motivoCancelacion: null, creadoEn: new Date().toISOString()
  };
}

// Aplica un evento de la saga; devuelve el nuevo estado o null si no hay transición (evento tardío/duplicado)
function aplicarEventoSaga(pedido, tipoEvento) {
  const nuevoEstado = transicionar(pedido.estado, tipoEvento);
  if (!nuevoEstado) return null;
  pedido.estado = nuevoEstado;
  if (nuevoEstado === 'CANCELADO') pedido.motivoCancelacion = MOTIVOS[tipoEvento];
  return nuevoEstado;
}

// Carga que viaja en pedido.creado / pedido.confirmado / pedido.cancelado (contrato de 02-05 y AsyncAPI de 03-06)
function datosParaConsumidores(p) {
  return { pedidoId: p.pedidoId, clienteId: p.clienteId, cliente: p.cliente, direccionEnvio: p.direccionEnvio, total: p.total,
           lineas: p.lineas.map((l) => ({ productoId: l.productoId, nombre: l.nombreProducto, cantidad: l.cantidad, precioUnitario: l.precioUnitario })) };
}
module.exports = { crearPedido, aplicarEventoSaga, datosParaConsumidores };

  1. Repositorio con outbox: guardarConEventos() y procesarUnaVez()

// src/repositorios/pedidoRepositorio.js
const { randomUUID } = require('node:crypto');

function crearPedidoRepositorio(bd) {
  // Guarda pedido + líneas + réplica del cliente + eventos en el outbox, en UNA transacción (02-05 §7).
  // `tx` opcional: si el llamante ya está en una transacción (procesarUnaVez), reutilizamos la suya.
  async function guardarConEventos(pedido, eventos, tx) {
    const trabajo = async (t) => {
      await t.consultar(`INSERT INTO pedidos (pedido_id, cliente_id, estado, total, direccion_envio, motivo_cancelacion, creado_en, actualizado_en)
                         VALUES ($1,$2,$3,$4,$5,$6,$7,NOW())
                         ON CONFLICT (pedido_id) DO UPDATE SET estado = EXCLUDED.estado, motivo_cancelacion = EXCLUDED.motivo_cancelacion, actualizado_en = NOW()`,
        [pedido.pedidoId, pedido.clienteId, pedido.estado, pedido.total, pedido.direccionEnvio, pedido.motivoCancelacion, pedido.creadoEn]);
      for (const l of pedido.lineas) {                      // las líneas nunca cambian tras la creación: insert idempotente
        await t.consultar(`INSERT INTO lineas_pedido (pedido_id, linea, producto_id, nombre_producto, precio_unitario, cantidad)
                           VALUES ($1,$2,$3,$4,$5,$6) ON CONFLICT DO NOTHING`, [pedido.pedidoId, l.linea, l.productoId, l.nombreProducto, l.precioUnitario, l.cantidad]);
      }
      if (pedido.cliente) {                                 // réplica clientes_ref (02-04): la calentamos con lo que ya sabemos
        await t.consultar(`INSERT INTO clientes_ref (cliente_id, nombre, email, actualizado_en) VALUES ($1,$2,$3,NOW())
                           ON CONFLICT (cliente_id) DO NOTHING`, [pedido.clienteId, pedido.cliente.nombre, pedido.cliente.email]);
      }
      for (const ev of eventos) {                           // el outbox: mismos INSERT, misma transacción
        await t.consultar(`INSERT INTO outbox (evento_id, agregado_tipo, agregado_id, tipo, version, carga) VALUES ($1,'Pedido',$2,$3,$4,$5)`,
          [`evt-${randomUUID()}`, pedido.pedidoId, ev.tipo, ev.version ?? 1, ev.carga]);
      }
    };
    return tx ? trabajo(tx) : bd.transaccion(trabajo);
  }

  async function obtener(pedidoId, tx = bd) {
    const { rows } = await tx.consultar(`SELECT p.*, c.nombre AS cliente_nombre, c.email AS cliente_email
                                          FROM pedidos p LEFT JOIN clientes_ref c ON c.cliente_id = p.cliente_id WHERE p.pedido_id = $1`, [pedidoId]);
    if (rows.length === 0) return null;
    const lineas = (await tx.consultar('SELECT * FROM lineas_pedido WHERE pedido_id = $1 ORDER BY linea', [pedidoId])).rows;
    return aPedido(rows[0], lineas);                       // fila → agregado (nombres camelCase; ~10 líneas, omitidas)
  }

  // Idempotencia de la API (02-05 §8c, 03-01): clave → respuesta guardada
  const buscarClave = async (clave) => (await bd.consultar('SELECT huella, respuesta FROM claves_idempotencia WHERE clave = $1', [clave])).rows[0] ?? null;
  const guardarClave = (t, clave, pedidoId, huella, respuesta) =>
    t.consultar('INSERT INTO claves_idempotencia (clave, pedido_id, huella, respuesta) VALUES ($1,$2,$3,$4)', [clave, pedidoId, huella, respuesta]);

  // Idempotencia de consumidores (02-05 §8b, 03-02): efecto + registro en la MISMA transacción
  async function procesarUnaVez(eventoId, consumidor, fn) {
    return bd.transaccion(async (tx) => {
      const { rowCount } = await tx.consultar('INSERT INTO eventos_procesados (evento_id, consumidor) VALUES ($1,$2) ON CONFLICT DO NOTHING', [eventoId, consumidor]);
      if (rowCount === 0) return 'DUPLICADO';              // ya procesado: la transacción no hace nada más
      await fn(tx);                                        // el efecto real, con el mismo tx
      return 'PROCESADO';
    });
  }

  return { guardarConEventos, obtener, buscarClave, guardarClave, procesarUnaVez, transaccion: bd.transaccion };
}
module.exports = { crearPedidoRepositorio };

Dos matices: procesarUnaVez inserta primero en eventos_procesados (si dos réplicas reciben el mismo evento a la vez, la segunda se bloquea en la fila y al hacer COMMIT la primera ve rowCount = 0), y el efecto se ejecuta con el mismo tx, de modo que guardarConEventos(pedido, eventos, tx) va en esa transacción: el cambio de estado, el evento de salida y la marca de procesado se confirman juntos o no se confirma ninguno.

  1. El caso de uso crearPedido y la ruta POST /v1/pedidos

// src/casos-uso/crearPedido.js
const { createHash } = require('node:crypto');
const { ErrorNegocio } = require('@techcorp/comun-http');
const Pedido = require('../dominio/pedido');

const huellaDe = (cuerpo) => createHash('sha256').update(JSON.stringify(cuerpo)).digest('hex');

function crearCasoUsoCrearPedido({ repositorio, catalogoCliente, clientesCliente, logger }) {
  return async function crearPedido(peticion, { claveIdempotencia, requestId }) {
    // 1. Idempotencia: misma clave + mismo cuerpo → misma respuesta; misma clave + otro cuerpo → 422 (03-01)
    const huella = huellaDe(peticion);
    const previa = await repositorio.buscarClave(claveIdempotencia);
    if (previa) {
      if (previa.huella !== huella) throw new ErrorNegocio('CLAVE_IDEMPOTENCIA_REUTILIZADA', 'La Idempotency-Key ya se usó con otro cuerpo', 422);
      return { pedido: previa.respuesta, repetida: true };
    }
    // 2. Cliente y productos EN PARALELO: son independientes y así la latencia es la del más lento, no la suma
    const [cliente, productos] = await Promise.all([
      clientesCliente.obtenerCliente(peticion.clienteId, { requestId }),                            // CLIENTE_NO_EXISTE → 404
      catalogoCliente.obtenerProductos([...new Set(peticion.lineas.map((l) => l.productoId))], { requestId })   // PRODUCTO_NO_DISPONIBLE → 422
    ]);
    // 3. El agregado, en PENDIENTE, con nombres y precios congelados y el total calculado
    const pedido = Pedido.crearPedido(peticion, { cliente, productos });
    const respuesta = aRepresentacion(pedido);           // el JSON de 03-01 (id, estado, lineas, total, _links); ~8 líneas, omitidas
    // 4. Pedido + pedido.creado + clave de idempotencia: UNA transacción. Ni RabbitMQ ni HTTP aquí dentro.
    await repositorio.transaccion(async (tx) => {
      await repositorio.guardarConEventos(pedido, [{ tipo: 'pedido.creado', version: 1, carga: Pedido.datosParaConsumidores(pedido) }], tx);
      await repositorio.guardarClave(tx, claveIdempotencia, pedido.pedidoId, huella, respuesta);
    });
    logger.info({ pedidoId: pedido.pedidoId, clienteId: pedido.clienteId, total: pedido.total, requestId }, 'pedido creado');
    return { pedido: respuesta, repetida: false };
  };
}
module.exports = { crearCasoUsoCrearPedido };

La ruta es la de 03-01 con /v1/ y sin el try/catch de mapeo de errores (ahora lo hace middlewareErrores de 04-02, porque los clientes lanzan ErrorNegocio con status):

// src/rutas/pedidos.js (fragmento POST)
enrutador.post('/v1/pedidos', async (req, res, next) => {
  try {
    const claveIdempotencia = req.get('Idempotency-Key');
    if (!claveIdempotencia) throw new ErrorNegocio('PETICION_INVALIDA', 'Falta la cabecera Idempotency-Key', 400);
    const peticion = esquemaNuevoPedido.parse(req.body);   // zod: clienteId, lineas[{productoId, cantidad ≥ 1}], direccionEnvio (400/422 vía middleware)
    const { pedido } = await crearPedido(peticion, { claveIdempotencia, requestId: req.id });
    res.status(202).location(`/v1/pedidos/${pedido.id}`).json(pedido);
  } catch (err) { next(err); }
});

Al responder 202, el evento pedido.creado está en la tabla outbox, no en RabbitMQ. Eso es deliberado y es lo que hace robusto el diseño: si RabbitMQ está caído, el pedido se acepta igual y el evento saldrá cuando vuelva.

  1. El relay del outbox

// src/mensajeria/relayOutbox.js
const { construirSobre, publicarEvento } = require('@techcorp/comun-http/mensajeria/publicador');   // 03-02

function crearRelayOutbox({ bd, canalConfirm, intervaloMs, lote = 50, logger }) {
  let temporizador = null, parado = false;

  async function publicarPendientes() {
    // Una transacción por lote. FOR UPDATE SKIP LOCKED: si hay dos réplicas de Pedidos, cada una toma filas distintas.
    const publicados = await bd.transaccion(async (tx) => {
      const { rows } = await tx.consultar(
        `SELECT evento_id, tipo, version, carga FROM outbox WHERE publicado_en IS NULL ORDER BY creado_en LIMIT $1 FOR UPDATE SKIP LOCKED`, [lote]);
      if (rows.length === 0) return 0;
      for (const fila of rows) {
        // El sobre lleva el eventoId del outbox (no uno nuevo): si el relay reintenta, el consumidor lo reconoce como duplicado
        const sobre = { ...construirSobre(fila.tipo, fila.carga, { version: fila.version }), eventoId: fila.evento_id };
        publicarEvento(canalConfirm, sobre);
      }
      await canalConfirm.waitForConfirms();                // el broker confirma que ha recibido (y persistido) el lote
      await tx.consultar(`UPDATE outbox SET publicado_en = NOW() WHERE evento_id = ANY($1)`, [rows.map((r) => r.evento_id)]);
      return rows.length;                                  // COMMIT: solo ahora quedan marcadas
    });
    if (publicados > 0) logger.debug({ publicados }, 'outbox publicado');
    return publicados;
  }

  async function ciclo() {
    if (parado) return;
    try {
      const n = await publicarPendientes();
      temporizador = setTimeout(ciclo, n === lote ? 0 : intervaloMs);   // si el lote venía lleno, seguir sin esperar
    } catch (err) {
      logger.error({ err }, 'relay outbox: fallo, reintento en el siguiente ciclo');   // RabbitMQ caído: las filas siguen pendientes
      temporizador = setTimeout(ciclo, intervaloMs);
    }
  }
  return { iniciar: () => ciclo(), parar: () => { parado = true; clearTimeout(temporizador); } };
}
module.exports = { crearRelayOutbox };

Si el proceso muere entre waitForConfirms y el COMMIT, las filas se vuelven a publicar en el siguiente ciclo: es el at-least-once que aceptamos en 03-02, absorbido por procesarUnaVez en los consumidores.

  1. Consumidores: la saga y la réplica de clientes

El consumidor de la saga sigue el patrón de 03-02 (prefetch, ack tras procesar, nack a DLQ) sobre la cola pedidos.saga; lo que cambia es el manejador, que ahora usa el dominio y el repositorio reales:

// src/mensajeria/consumidorSaga.js
const { declararColaConsumidor } = require('@techcorp/comun-http/mensajeria/topologia');   // 03-02
const Pedido = require('../dominio/pedido');

const COLA = 'pedidos.saga';
const EVENTOS = ['stock.reservado', 'stock.rechazado', 'pago.confirmado', 'pago.rechazado'];

function crearConsumidorSaga({ canal, repositorio, logger }) {
  // Efecto de un evento de la saga. Se ejecuta DENTRO de procesarUnaVez, con su tx.
  async function manejar(sobre, tx) {
    const pedido = await repositorio.obtener(sobre.carga.pedidoId, tx);
    if (!pedido) { logger.warn({ sobre }, 'evento para pedido desconocido'); return; }   // no reintentar: a ack (queda registrado como procesado)

    let nuevoEstado = Pedido.aplicarEventoSaga(pedido, sobre.tipo);
    if (!nuevoEstado) { logger.info({ pedidoId: pedido.pedidoId, de: pedido.estado, evento: sobre.tipo }, 'transicion_ignorada'); return; }
    if (nuevoEstado === 'PAGADO') nuevoEstado = Pedido.aplicarEventoSaga(pedido, 'confirmar');   // T4 de 02-05: hoy es inmediata

    const eventos = [];
    if (nuevoEstado === 'CONFIRMADO') eventos.push({ tipo: 'pedido.confirmado', carga: Pedido.datosParaConsumidores(pedido) });
    if (nuevoEstado === 'CANCELADO')  eventos.push({ tipo: 'pedido.cancelado',  carga: { ...Pedido.datosParaConsumidores(pedido), motivo: pedido.motivoCancelacion } });
    await repositorio.guardarConEventos(pedido, eventos, tx);        // estado nuevo + eventos de salida, misma transacción que eventos_procesados
    logger.info({ pedidoId: pedido.pedidoId, estado: nuevoEstado, evento: sobre.tipo }, 'pedido actualizado');
  }

  async function iniciar() {
    await declararColaConsumidor(canal, COLA, EVENTOS);
    await canal.prefetch(10);
    await canal.consume(COLA, async (msg) => {
      if (!msg) return;
      let sobre;
      try { sobre = JSON.parse(msg.content.toString()); } catch { return canal.nack(msg, false, false); }   // ilegible → DLQ
      try {
        await repositorio.procesarUnaVez(sobre.eventoId, 'pedidos.saga', (tx) => manejar(sobre, tx));
        canal.ack(msg);
      } catch (err) {
        logger.error({ err, eventoId: sobre.eventoId }, 'error procesando evento de saga');
        canal.nack(msg, false, !msg.fields.redelivered);              // 1.ª vez: reencolar; 2.ª: DLQ (03-02)
      }
    });
  }
  return { iniciar };
}
module.exports = { crearConsumidorSaga };

El consumidor de pedidos.clientes (consumidorClientes.js) es la misma estructura con un solo evento, cliente.actualizado, y un manejador de una sentencia: INSERT INTO clientes_ref ... ON CONFLICT (cliente_id) DO UPDATE SET nombre, email, actualizado_en = EXCLUDED.actualizado_en WHERE clientes_ref.actualizado_en < EXCLUDED.actualizado_en (la condición descarta eventos que lleguen desordenados). Ambos consumidores usan un canal propio sobre la misma conexión AMQP; el relay usa un tercer canal de confirmación (createConfirmChannel).

  1. GET /v1/pedidos/{id} con ETag y el arranque completo

// src/rutas/pedidos.js (fragmento GET)
enrutador.get('/v1/pedidos/:id', async (req, res, next) => {
  try {
    const pedido = await repositorio.obtener(req.params.id);
    if (!pedido) throw new ErrorNegocio('PEDIDO_NO_EXISTE', `No existe ${req.params.id}`, 404);
    const etag = `"${pedido.pedidoId}:${new Date(pedido.actualizadoEn).getTime()}"`;   // cambia con cada transición
    if (req.get('If-None-Match') === etag) return res.status(304).end();               // polling barato (03-01)
    res.set('ETag', etag).set('Cache-Control', 'no-cache').json(aRepresentacion(pedido));
  } catch (err) { next(err); }
});

servidor.js amplía el de 04-02: carga config (04-03), crea el pool y el repositorio, conecta a RabbitMQ con conectar(config.RABBITMQ_URL) (03-02) y abre canalConfirm = await conexion.createConfirmChannel(), construye los clientes HTTP, compone crearApp({ repositorio, catalogoCliente, clientesCliente, logger, comprobacionesSalud: { postgres: bd.ping, rabbitmq: () => canal.closed ? Promise.reject(new Error('canal cerrado')) : Promise.resolve() } }), y arranca relay y consumidores después de listen (crearApp construye internamente el caso de uso con crearCasoUsoCrearPedido({ repositorio, catalogoCliente, clientesCliente, logger }), de modo que las pruebas de 04-05 puedan sustituir cualquiera de las tres piezas). En el apagado, el orden es el inverso: relay.parar(), cerrar canales y conexión AMQP (los mensajes sin ack se reentregan), servidor.close(), bd.cerrar(). Como decidió el ejercicio 1 de 03-05, /health/ready comprueba PostgreSQL y RabbitMQ, no Catálogo ni Clientes.

sequenceDiagram
    participant W as Web
    participant R as rutas/pedidos.js
    participant CU as crearPedido
    participant CL as servicio-clientes:3004
    participant CA as servicio-catalogo:3001
    participant PG as PostgreSQL (pedidos)
    participant RL as relayOutbox
    participant MQ as RabbitMQ techcorp.eventos
    participant CS as consumidorSaga
    W->>R: POST /v1/pedidos (Idempotency-Key)
    R->>CU: crearPedido(peticion)
    par en paralelo
        CU->>CL: GET /v1/clientes/c-1024
        CU->>CA: GET /v1/productos?ids=p-501,p-777
    end
    CU->>PG: BEGIN; pedidos+lineas+clientes_ref+outbox(pedido.creado)+claves_idempotencia; COMMIT
    R-->>W: 202 Location: /v1/pedidos/ped-…
    RL->>PG: SELECT … FOR UPDATE SKIP LOCKED
    RL->>MQ: publish pedido.creado (confirm)
    RL->>PG: UPDATE outbox SET publicado_en
    Note over MQ: Inventario reserva y publica stock.reservado
    MQ->>CS: stock.reservado (cola pedidos.saga)
    CS->>PG: procesarUnaVez: eventos_procesados + estado STOCK_RESERVADO
    CS->>MQ: ack

  1. Prueba de punta a punta

Dependencias locales (04-01) y arranque; Catálogo (04-02) debe estar corriendo en 3001. Como servicio-clientes aún no existe, en desarrollo se apunta CLIENTES_URL a un stub de 20 líneas (scripts/stubClientes.js, Express que responde GET /v1/clientes/c-1024 con Ana Ruiz y 404 para el resto):

docker run -d --name pg-pedidos -p 5432:5432 -e POSTGRES_USER=svc_pedidos -e POSTGRES_PASSWORD=dev-pedidos -e POSTGRES_DB=pedidos postgres:16
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
cp .env.ejemplo .env && npm run migrar && node scripts/stubClientes.js &   # stub en 3004
npm run dev
# 1. Crear el pedido de Ana → 202
curl -s -i -X POST http://localhost:3002/v1/pedidos -H 'Content-Type: application/json' -H 'Idempotency-Key: 7f3c9a2e-1b4d-4e8f-9c21-5a6b7c8d9e0f' \
  -d '{"clienteId":"c-1024","lineas":[{"productoId":"p-501","cantidad":1},{"productoId":"p-777","cantidad":2}],
       "direccionEnvio":{"calle":"Gran Vía 12","codigoPostal":"28013","ciudad":"Madrid","pais":"ES"}}'
# HTTP/1.1 202 Accepted   Location: /v1/pedidos/ped-3f9a1c2b   cuerpo: {"id":"ped-3f9a1c2b","estado":"PENDIENTE","total":79.7,...}

# 2. Repetir la misma petición → 202 con el MISMO id (idempotencia); cambiar cantidad con la misma clave → 422 CLAVE_IDEMPOTENCIA_REUTILIZADA

# 3. Consultar → PENDIENTE (y en RabbitMQ Management, cola inventario.pedidos: 1 mensaje pedido.creado si Inventario no está arrancado)
curl -s http://localhost:3002/v1/pedidos/ped-3f9a1c2b | jq .estado      # "PENDIENTE"

# 4. Simular a Inventario: publicar stock.reservado con el script
node scripts/publicarEvento.js stock.reservado '{"pedidoId":"ped-3f9a1c2b","reservaId":"res-4471"}'
curl -s http://localhost:3002/v1/pedidos/ped-3f9a1c2b | jq .estado      # "STOCK_RESERVADO"

# 5. Simular a Pagos → CONFIRMADO, y en la cola notificaciones.pedidos aparece pedido.confirmado
node scripts/publicarEvento.js pago.confirmado '{"pedidoId":"ped-3f9a1c2b","pagoId":"pag-9001","importe":79.70}'
curl -s http://localhost:3002/v1/pedidos/ped-3f9a1c2b | jq .estado      # "CONFIRMADO"

# 6. Repetir el paso 4 → el estado NO cambia (transicion_ignorada en el log): idempotencia + máquina de estados
// scripts/publicarEvento.js — publica un sobre estándar en techcorp.eventos: node scripts/publicarEvento.js <tipo> '<carga JSON>'
const { conectar } = require('@techcorp/comun-http/mensajeria/topologia');
const { construirSobre, publicarEvento } = require('@techcorp/comun-http/mensajeria/publicador');
(async () => {
  const [tipo, cargaJson] = process.argv.slice(2);
  const { conexion, canal } = await conectar(process.env.RABBITMQ_URL ?? 'amqp://localhost:5672');
  const sobre = construirSobre(tipo, JSON.parse(cargaJson));
  publicarEvento(canal, sobre);
  console.log('publicado', sobre.eventoId, tipo);
  await canal.close(); await conexion.close();
})();

Si los seis pasos se comportan así, el flujo "un cliente hace un pedido" funciona de punta a punta en la parte de Pedidos, con las garantías diseñadas: sin dual write, sin dobles pedidos y sin transiciones imposibles.

Errores Comunes y Consejos

  • Publicar en RabbitMQ dentro de crearPedido. Es el dual write de 02-05: pedido guardado y evento perdido (o al revés). Solo outbox; el relay publica.
  • Un eventoId nuevo en cada intento del relay. Los consumidores no reconocerían el duplicado. El sobre reutiliza outbox.evento_id.
  • ack antes de procesar o efecto fuera de la transacción de procesarUnaVez. Se pierden eventos o se aplican dos veces. Efecto y registro con el mismo tx; ack al final.
  • Consultar Clientes y Catálogo en serie. Promise.all: son independientes. (Y si uno falla, Promise.all rechaza en cuanto falla el primero: correcto aquí, porque sin ambos no hay pedido).
  • Calcular el total con decimales flotantes. 59.90 + 19.80 no siempre es 79.70. Céntimos enteros y división al final; en la BD, NUMERIC(10,2).
  • Olvidar pool.on('error'). Un error en una conexión inactiva es un 'error' sin manejador: el proceso muere.
  • FOR UPDATE sin SKIP LOCKED. Con dos réplicas, la segunda espera a la primera en cada ciclo; con SKIP LOCKED trabajan en paralelo sin duplicar.
  • Bloquear la respuesta HTTP hasta que la saga termina. El contrato es 202 + polling con ETag. Esperar dentro de la petición recrea el acoplamiento temporal.

Ejercicios

Ejercicio 1. Dos réplicas de Pedidos reciben a la vez el mismo POST /v1/pedidos con la misma Idempotency-Key (el cliente reintentó por un timeout de red). Recorre el código de crearPedido y explica qué ocurre en cada réplica; identifica el punto en el que la clave primaria de claves_idempotencia decide el resultado y qué debería devolver la réplica "perdedora" (pista: código de error de PostgreSQL 23505).

Ejercicio 2. Escribe el manejador del consumidor pedidos.clientes completo (manejar(sobre, tx)) para cliente.actualizado con carga { clienteId, nombre, email, actualizadoEn }, con la protección contra desorden del apartado 8, y explica por qué no hace falta la máquina de estados aquí.

Ejercicio 3. El relay tiene una réplica y publica 50 eventos por ciclo cada 500 ms. En el pico de campaña (×20 → 60.000 pedidos/día ≈ 0,7 pedidos/s de media, con ráfagas de 10/s) ¿es suficiente? Calcula y propón dos ajustes de configuración (04-03) sin cambiar código.

Soluciones

Solución 1. Ambas réplicas ejecutan buscarClave casi a la vez y ninguna encuentra la clave; ambas llaman a Clientes y Catálogo y construyen un pedido con ids distintos (ped-a…, ped-b…); ambas abren su transacción. La primera en hacer COMMIT inserta su fila en claves_idempotencia; la segunda, al ejecutar guardarClave, choca con la clave primaria (23505 unique_violation) y su transacción entera hace ROLLBACK: su pedido y su pedido.creado desaparecen, que es exactamente lo deseado. Falta tratar el error: en crearPedido, capturar err.code === '23505' alrededor de la transacción, volver a leer con buscarClave y devolver la respuesta guardada por la ganadora (con repetida: true). Sin ese catch, la réplica perdedora respondería 500 a un cliente que, si reintenta, ya recibirá el pedido correcto; con él, responde 202 con el mismo pedido.

Solución 2.

async function manejar(sobre, tx) {
  const { clienteId, nombre, email, actualizadoEn } = sobre.carga;
  await tx.consultar(
    `INSERT INTO clientes_ref (cliente_id, nombre, email, actualizado_en) VALUES ($1,$2,$3,$4)
     ON CONFLICT (cliente_id) DO UPDATE SET nombre = EXCLUDED.nombre, email = EXCLUDED.email, actualizado_en = EXCLUDED.actualizado_en
     WHERE clientes_ref.actualizado_en < EXCLUDED.actualizado_en`,
    [clienteId, nombre, email, actualizadoEn]);
}

No hay máquina de estados porque clientes_ref no es un agregado con invariantes: es una réplica de solo lectura (02-04) cuyo único requisito es converger al último valor conocido. La condición WHERE actualizado_en < EXCLUDED.actualizado_en garantiza que un evento antiguo que llegue tarde no pise a uno más nuevo; procesarUnaVez cubre el duplicado exacto.

Solución 3. Capacidad del relay: 50 eventos cada 500 ms = 100 eventos/s (y más, porque con lote lleno encadena ciclos sin esperar). Cada pedido genera 2 eventos de Pedidos (pedido.creado y pedido.confirmado/cancelado): la ráfaga de 10 pedidos/s son 20 eventos/s, cinco veces por debajo. Es suficiente. Lo que sí importa es la latencia: con 500 ms de intervalo, un evento espera de media 250 ms antes de salir. Ajustes de configuración: bajar OUTBOX_INTERVALO_MS a 250 en producción (ya en la tabla de 04-03) y, si se quiere margen, exponer el tamaño de lote como OUTBOX_LOTE (100). Añadir una segunda réplica de Pedidos también duplica el relay sin cambios gracias a SKIP LOCKED.

Conclusión

servicio-pedidos está construido en su parte esencial y sigue, pieza a pieza, lo diseñado antes: el pool de PostgreSQL con transaccion(fn); las migraciones con los esquemas de 02-04 y 02-05 (más la huella en claves_idempotencia); los clientes HTTP a GET /v1/productos?ids= y GET /v1/clientes/{id} con timeout, mapeo a DEPENDENCIA_NO_DISPONIBLE/CLIENTE_NO_EXISTE/PRODUCTO_NO_DISPONIBLE y el ACL traductorProducto; el agregado Pedido con precios congelados y total en céntimos; crearPedido con Idempotency-Key, Promise.all y una sola transacción para pedido, pedido.creado y clave; el relay del outbox con FOR UPDATE SKIP LOCKED y waitForConfirms; los consumidores de pedidos.saga y pedidos.clientes con procesarUnaVez y la máquina de estados; GET /v1/pedidos/{id} con ETag; y una prueba manual que lleva el pedido de Ana de PENDIENTE a CONFIRMADO.

Todo eso lo hemos comprobado a mano, con curl y un script. No sirve como red de seguridad: mañana alguien tocará traductorProducto o el consumidor de la saga y nadie repetirá los seis pasos. La siguiente lección convierte estas comprobaciones en pruebas automáticas a distintos niveles: unitarias del dominio (Pedido, transicionar), de componente contra crearApp con dobles, de integración con PostgreSQL y RabbitMQ reales mediante Testcontainers, y de contrato con Pact entre Pedidos y Catálogo, para que el GET /v1/productos?ids= que hoy funciona no se rompa en silencio.

Curso de Microservicios

Módulo 1: Introducción a los Microservicios

Módulo 2: Diseño de Microservicios

Módulo 3: Comunicación entre Microservicios

Módulo 4: Implementación de Microservicios

Módulo 5: Despliegue y Orquestación

Módulo 6: Monitoreo y Mantenimiento

Módulo 7: Seguridad en Microservicios

Módulo 8: Casos de Estudio y Ejemplos Prácticos

© Copyright 2026. Todos los derechos reservados