Ingesta de eventos desde el flujo
Streaming Hub consume el flujo de la plataforma como mensajes CloudEvents 1.0 en modo de contenido binario: los atributos de contexto de CloudEvents viajan en los headers del registro de Kafka (cada uno con el prefijo
ce-) y el cuerpo del evento es el valor del registro. El hub lee el enrutamiento y la identidad desde los headers sin deserializar el payload.
Cada registro que el hub acepta debe llevar este conjunto de contexto:
Para cada registro, el hub devuelve uno de tres veredictos:
El tenant admitido en la ingesta proviene únicamente del atributo
ce-tenantid validado que extrajo el parser. Un tenant desconocido o inactivo se descarta (drop) en la puerta del roster: el único punto donde se aplica el aislamiento de tenants en el camino de ingesta.
Tres strings relacionados se parecen pero son distintos; nunca los confundas:
ce-type— el tipo de CloudEvents:studio.lerian.<resource>.<event>(por ejemplo,studio.lerian.transaction.posted).- Tópico de Kafka — dónde viaja el evento en el bus interno. No existe un único prefijo común a toda la plataforma: cada productor tiene su propio namespace de tópicos. Midaz enruta a tópicos explícitos con un segmento de servicio fijo — los eventos del ledger viajan en
lerian.streaming.ledger_<resource>.<event>, con el servicio plegado en el primer segmento (así quetransaction.postedviaja enlerian.streaming.ledger_transaction.posted) y cualquier guion en el recurso o el evento convertido en un guion bajo en el tópico (balance.config-changedviaja enlerian.streaming.ledger_balance.config_changed). Otros productores — Consignado y los comandos entre productos — derivan su namespace de suce-source. X-Lerian-Event-Type— el tipo estampado en un webhook entregado: la cola escueta<resource>.<event>delce-type(transaction.posted). Refleja la identidad delce-type, no el tópico de Kafka — no lleva ningún segmento de servicio del productor (sin el plegadoledger_).
Deduplicación consume-once
El flujo interno es at-least-once, así que el mismo registro puede llegar más de una vez. En el camino hacia el inbox de eventos, el hub deduplica por
ce-id: un id duplicado no escribe nada y se cuenta como un descarte por deduplicación. Los registros se persisten por partición de Kafka en una sola transacción, y la posición del flujo de una partición se confirma solo después de que la escritura de esa partición se haya confirmado, de modo que un fallo en una partición nunca pierde ni confirma dos veces eventos de otra.
Coincidencia de suscripciones
Una vez que un evento está almacenado, el hub resuelve cuál de tus suscripciones debe recibirlo. La coincidencia de suscripciones es interna y no depende del catálogo: evalúa cada evento contra tus suscripciones por
(tenant, tipo de evento, major de esquema) directamente, sin consultar el catálogo de eventos. Un tipo de evento que ninguna suscripción quiere produce cero entregas: un resultado correcto, no un error.
La coincidencia solo admite una suscripción cuando está a la vez enabled y en verification_state = active. Son dos condiciones independientes —consulta el modelo de suscripción— y la coincidencia requiere ambas.
Envío a tus destinos
Por cada coincidencia, el hub crea un trabajo de entrega y lo pasa a un pool de workers (ocho workers por defecto). Los workers reclaman los trabajos pendientes de forma atómica con un lease corto y se intercalan entre tenants, así que el backlog de un solo tenant nunca deja sin recursos a los demás. Cada worker carga el payload almacenado, descifra en memoria el secreto de firma o la credencial del destino, y realiza exactamente un intento de entrega al sink. La entrega es at-least-once: el hub puede entregar el mismo evento más de una vez (mediante reintentos o reenvíos después de reiniciar un worker). Cada entrega lleva un
X-Lerian-Event-Id estable (el ce-id) sobre el que deduplicar, y un X-Lerian-Delivery-Id por intento que cambia en cada reintento. Trata un X-Lerian-Event-Id repetido como un duplicado y confírmalo sin reprocesarlo. Consulta Consumo de eventos para el lado del consumidor de este contrato.
Reintentos y back-off
Cuando un intento falla, el hub programa el siguiente en una curva de back-off fija con jitter completo:
Dead-lettering
Un evento cuyos reintentos se agotan pasa a dead-letter: el hub deja de intentarlo y registra el resultado terminal. Un trabajo que, en cambio, sigue haciendo caer a un worker a mitad de un intento —sin registrar nunca un resultado— se reclama un número acotado de veces (cinco por defecto) y luego pasa a dead-letter como poison, así que un único trabajo tóxico nunca puede ocupar un worker para siempre.
Circuit breaker
Cada par
(tenant, destino) tiene un circuit breaker. Cuando un destino falla repetidamente, su breaker se abre y los trabajos posteriores para ese destino se descartan —se reprograman sin intento—, así que un único endpoint roto no consume capacidad de los workers ni martillea a un destino en apuros. El breaker se recupera por sí solo en cuanto el destino vuelve a tener éxito.
Auto-desactivación de un destino roto
Un destino que permanece roto durante mucho tiempo se desactiva automáticamente. La auto-desactivación se dispara solo cuando el tramo de fallos abierto es a la vez:
- sostenido — ha estado fallando durante al menos la ventana de fallos (120 horas por defecto); y
- disperso — al menos 12 horas separan el primer y el último fallo del tramo.
enabled de la suscripción a false; nunca toca verification_state. La entrega se detiene (la coincidencia requiere enabled) y el hub registra el motivo.
Recuperas una suscripción auto-desactivada con POST /v1/subscriptions/:id/verify, que vuelve a sondear el destino y, si tiene éxito, la vuelve a habilitar en su sitio. La auto-desactivación está gobernada por un kill switch (STREAMING_HUB_AUTODISABLE_ENABLED), así que un operador puede desplegarla desactivada. Consulta recuperar una suscripción auto-desactivada.
Garantías de orden
El flujo interno está particionado por tenant, así que los eventos de un tenant normalmente viajan por una sola partición y el hub preserva su orden primero en entrar, primero en salir de extremo a extremo. El hub atribuye y deduplica según los headers de CloudEvents, no según la record key de Kafka, así que un productor que reparte (hace salting de) un tenant activo entre varias particiones cambia solo la ubicación física: el hub sigue atribuyendo y deduplicando correctamente. La única contrapartida es el orden: los eventos de un tenant con salting abarcan varias particiones sin garantía de orden entre particiones, así que ese tenant renuncia al FIFO estricto. La secuencia de llegada que el hub asigna refleja el orden en que los eventos llegaron al hub, no el orden en que se produjeron. Un tenant sin salting conserva el FIFO de una sola partición en todo momento.
Próximos pasos
Gestión de suscripciones
El modelo de suscripción, los flujos de incorporación y la rotación de secretos.
Consumo de eventos
Verifica firmas, deduplica y consulta eventos mediante pull.

