Skip to main content
Cada job de extracción termina en un único evento terminal. El Worker publica 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.
Fetcher copia los metadatos enviados al evento, así que source y un identificador de correlación siguen siendo campos de tu aplicación. En job.failed, metadata.error está reservado para los detalles de error saneados por Fetcher y reemplaza un valor enviado por quien llama para esa clave. size_bytes y row_count 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. La extracción se detiene ante el primer error, así que el primer datasource que falla termina el job entero. Un payload job.failed no incluye result ni completed_at, pero no prueba que no exista un objeto: el Worker escribe el objeto cifrado antes de persistir completed, y un fallo al persistir el estado terminal puede producir job.failed después de que el almacenamiento haya tenido éxito. 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ó.
Un evento de fallo no lleva bloque result ni completed_at. 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.
El esquema 2.0.0 es una versión de payload incompatible. Actualiza tu consumidor para los campos de Fetcher en snake_case que aparecen arriba. Las routing keys job.completed y job.failed, y el formato de ce-id, no cambian.

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.
  • Cada 30 segundos, un reparador revisa jobs terminales cuyo marcador de evento pendiente sigue activo. Eso incluye un fallo al limpiar el marcador después de una publicación exitosa, por lo que puede reemitir un evento que ya llegó al broker.
  • ce-id se 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.
Trata un ce-id repetido como un duplicado y confírmalo sin reprocesarlo. Un consumidor que use como clave el identificador de mensaje del broker procesará el mismo job dos veces.

Verificar lo que recibes


Dos artefactos protegidos con HMAC usan 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 los verifica sin tener 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 sobre firma el timestamp, la versión de firma, el tenant, el exchange, la routing key y el cuerpo. Los eventos v2 actuales emiten job_id, mientras que el extractor del sobre reconoce jobId; no dependas de un campo de ID de job extraído por separado. 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.