job.completed cuando el resultado llega al almacenamiento, y job.failed cuando la ejecución se detiene. Suscríbete a esos eventos en lugar de consultar GET /v1/fetcher/{id} de forma periódica.
Esta página cubre el lado del consumidor: qué llega, qué significa y qué hacer con ello.
Dónde llegan los eventos
El Worker publica en un exchange topic durable.
RABBITMQ_JOB_EVENTS_EXCHANGE lo nombra, y el valor distribuido es fetcher.job.events.
Vincula tu propia cola a las routing keys que te interesan. La definición de infraestructura local vincula dos colas como ejemplo:
fetcher.job.completed.queue y fetcher.job.failed.queue.
Los eventos de job son un contrato de producto, no una funcionalidad opcional. El Worker se niega a arrancar sin streaming activado y sin un nombre de exchange. Consulta Despliegue para el lado del operador.
job.completed
Este evento significa una sola cosa: la extracción llegó hasta el final, y el resultado cifrado ya está en el almacenamiento de objetos. El Worker lo publica después de dos pasos previos. Primero escribe el objeto del resultado, luego registra el estado terminal en el job. Solo entonces emite el evento.
metadata lleva la metadata que enviaste en el job, así que source y cualquier identificador de correlación vuelven a ti sin cambios.
sizeBytes y rowCount describen el resultado en texto plano, antes del cifrado. path es la clave de objeto del resultado almacenado.
job.failed
Este evento significa que la ejecución se detuvo y Fetcher no almacenó ningún resultado. La extracción se detiene ante el primer error, así que el primer datasource que falla termina el job entero. Un resultado parcial nunca llega al almacenamiento. El evento se dispara ante cualquier fallo del camino: un nombre de conexión que no se resuelve, un error de datasource, un esquema que no corresponde o una escritura en almacenamiento que no se completó.
result ni completedAt.
metadata.error.message pasa primero por un paso de redacción. Fetcher reemplaza cuatro formas de fuga por [redacted]: las URIs de conexión, el operando de dirección de un error de red de Go, el operando Addr: de un error del driver de MongoDB y una dirección IPv4 literal. El resto del texto sobrevive, así que el mensaje sigue siendo accionable. Dirígelo a los operadores.
El sobre CloudEvents
Cada mensaje viaja en modo binario de CloudEvents, versión 1.0. Los atributos de contexto viajan como headers AMQP.
Un despliegue de tenant único también lleva un valor de tenant. Emite el literal
single-tenant, así que un mismo consumidor maneja ambas formas de despliegue con el mismo código.
Configuración del source de CloudEvents
STREAMING_CLOUDEVENTS_SOURCE no tiene valor por defecto. El Worker lo exige siempre que el streaming esté activo, y se detiene al arrancar cuando el valor está vacío. El ejemplo distribuido usa //lerian.fetcher/worker.
Fetcher copia el valor en ce-source tal cual. Dale a cada despliegue del Worker su propio valor de source cuando varios productores comparten un broker, y enruta por ese header.
Contrato de entrega
La entrega es at-least-once. Deduplica sobre
ce-id.
- El Worker escribe el evento en un outbox durable antes de publicarlo. Una caída del broker retrasa el evento, no lo pierde.
- Un reparador busca eventos terminales que nunca se publicaron y los reemite cada 30 segundos.
ce-idse mantiene idéntico en cada reemisión del mismo job y el mismo estado. Nada más es lo bastante estable para usarlo como clave.- El orden no está garantizado. Dos jobs pueden completarse en un orden y llegar en otro.
- El Worker registra el estado terminal del job antes de emitir.
GET /v1/fetcher/{id}sigue siendo la autoridad sobre el estado del job.
Verificar lo que recibes
Dos firmas independientes protegen el camino entre el Worker y tú. Ambas usan HMAC-SHA256 con la clave HMAC externa. HKDF-SHA256 deriva esa clave a partir de la clave maestra
APP_ENC_KEY. El repositorio de Fetcher trae una pequeña herramienta que la imprime, así que un consumidor verifica firmas sin tener nunca la clave maestra.
El mensaje. El Worker firma cada mensaje publicado y le estampa tres headers: x-message-signature, t para el timestamp de la firma y signature-version. El payload firmado une el timestamp, la versión de firma, el tenant, el identificador del job, el exchange y la routing key, seguidos del cuerpo del mensaje. Un replay del mismo cuerpo bajo otro tenant o por otra ruta falla la verificación.
El resultado. result.hmac es el HMAC-SHA256 con clave sobre el JSON del resultado en texto plano, calculado antes del cifrado. El bloque integrity declara el mismo valor junto con su algoritmo. Verifícalo después de descifrar y antes de confiar en las filas.
El bloque protection describe los bytes almacenados: encrypted es verdadero, el adaptador de almacenamiento aplicó el cifrado y el modo es adapter-managed. Describe solo el resultado, nunca las credenciales del datasource.
Próximos pasos
Jobs de extracción
Qué solicita un job y los estados por los que pasa.
API REST de Fetcher
Crea un job, lee un job y gestiona conexiones.
Configuración
Las variables de streaming, exchange y cifrado detrás de estos eventos.
Arquitectura
El Manager, el Worker y el Engine que ambos ejecutan.

