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
- Estructura del proyecto y dependencias
- PostgreSQL: pool, transacciones y migraciones
- Clientes HTTP salientes: Catálogo (con ACL) y Clientes
- Dominio: el agregado
Pedidoy la máquina de estados - Repositorio con outbox:
guardarConEventos()yprocesarUnaVez() - El caso de uso
crearPedidoy la rutaPOST /v1/pedidos - El relay del outbox
- Consumidores: la saga y la réplica de clientes
GET /v1/pedidos/{id}con ETag y el arranque completo- Prueba de punta a punta
- Estructura del proyecto y dependencias
npm install express pg amqplib pino pino-http zod @techcorp/comun-http && npm install -D nodemon dotenvservicio-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.
- 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).
- 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; 404 → ErrorNegocio('CLIENTE_NO_EXISTE', ..., 404); 5xx/red → DEPENDENCIA_NO_DISPONIBLE; devuelve { clienteId, nombre, email, direcciones }. Sin reintentos ni circuit breaker (06-03): solo timeout.
- Dominio: el agregado
Pedido y la máquina de estados
Pedido y la máquina de estadosdominio/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 };
- Repositorio con outbox:
guardarConEventos() y procesarUnaVez()
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.
- El caso de uso
crearPedido y la ruta POST /v1/pedidos
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.
- 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.
- 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).
GET /v1/pedidos/{id} con ETag y el arranque completo
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
- 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). Solooutbox; el relay publica. - Un
eventoIdnuevo en cada intento del relay. Los consumidores no reconocerían el duplicado. El sobre reutilizaoutbox.evento_id. ackantes de procesar o efecto fuera de la transacción deprocesarUnaVez. Se pierden eventos o se aplican dos veces. Efecto y registro con el mismotx;ackal final.- Consultar Clientes y Catálogo en serie.
Promise.all: son independientes. (Y si uno falla,Promise.allrechaza en cuanto falla el primero: correcto aquí, porque sin ambos no hay pedido). - Calcular el total con decimales flotantes.
59.90 + 19.80no siempre es79.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 UPDATEsinSKIP LOCKED. Con dos réplicas, la segunda espera a la primera en cada ciclo; conSKIP LOCKEDtrabajan 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
- Conceptos Básicos de Microservicios
- Ventajas y Desventajas de los Microservicios
- Comparación con la Arquitectura Monolítica
- Cuándo Adoptar Microservicios: Criterios de Decisión
- El Caso Práctico del Curso: la Tienda Online de TechCorp
Módulo 2: Diseño de Microservicios
- Principios de Diseño de Microservicios
- Descomposición de Aplicaciones Monolíticas
- Definición de Bounded Contexts
- Gestión de Datos: una Base de Datos por Servicio
- Consistencia Distribuida: Sagas, CQRS y Event Sourcing
Módulo 3: Comunicación entre Microservicios
- APIs RESTful
- Mensajería Asíncrona
- Protocolos de Comunicación: gRPC, GraphQL
- API Gateway y Backend for Frontend
- Descubrimiento de Servicios y Balanceo de Carga
- Contratos y Versionado de APIs
Módulo 4: Implementación de Microservicios
- Elección de Tecnologías y Herramientas
- Desarrollo de un Microservicio Simple
- Gestión de Configuración
- Integración Práctica: Consumir APIs y Publicar Eventos
- Pruebas en Microservicios: Unitarias, de Integración y de Contrato
Módulo 5: Despliegue y Orquestación
- Contenedores y Docker
- Orquestación con Kubernetes
- CI/CD para Microservicios
- Estrategias de Despliegue: Rolling, Blue-Green y Canary
- Service Mesh: Istio y Linkerd
Módulo 6: Monitoreo y Mantenimiento
- Monitoreo y Logging
- Trazabilidad Distribuida con OpenTelemetry
- Gestión de Errores y Recuperación
- Escalabilidad y Rendimiento
- SLOs, Alertas y Gestión de Incidentes
Módulo 7: Seguridad en Microservicios
- Autenticación y Autorización
- Seguridad en la Comunicación
- Prácticas de Seguridad
- Seguridad en Contenedores y Kubernetes
