job.completed quando o resultado chega ao armazenamento, e job.failed quando a execução para. Assine esses eventos em vez de consultar GET /v1/fetcher/{id} periodicamente.
Esta página cobre o lado do consumidor: o que chega, o que significa e o que fazer com isso.
Onde os eventos chegam
O Worker publica em um exchange topic durável.
RABBITMQ_JOB_EVENTS_EXCHANGE o nomeia, e o valor distribuído é fetcher.job.events.
Vincule a sua própria fila às routing keys que lhe interessam. A definição de infraestrutura local vincula duas filas como exemplo:
fetcher.job.completed.queue e fetcher.job.failed.queue.
Os eventos de job são um contrato do produto, não um recurso opcional. O Worker recusa iniciar sem o streaming habilitado e sem um nome de exchange. Veja Implantação para o lado do operador.
job.completed
Este evento significa uma coisa: a extração rodou até o fim, e o resultado criptografado já está no armazenamento de objetos. O Worker o publica depois de dois passos anteriores. Primeiro ele grava o objeto de resultado, depois registra o status terminal no job. Só então emite o evento.
metadata carrega os metadados que você enviou no job, então source e qualquer identificador de correlação voltam para você sem alteração.
sizeBytes e rowCount descrevem o resultado em texto claro, antes da criptografia. path é a chave do objeto de resultado armazenado.
job.failed
Este evento significa que a execução parou e o Fetcher não armazenou resultado algum. A extração falha rápido, então a primeira fonte de dados que falha encerra o job inteiro. Um resultado parcial nunca chega ao armazenamento. O evento dispara para qualquer falha ao longo do caminho: um nome de conexão que não resolve, um erro da fonte de dados, um esquema incompatível ou uma gravação no armazenamento que não se completou.
result nem completedAt.
metadata.error.message passa antes por uma etapa de ocultação. O Fetcher troca quatro formatos de vazamento por [redacted]: URIs de conexão, o operando de endereço de um erro de rede do Go, o operando Addr: de um erro do driver do MongoDB e um literal IPv4. O restante do texto permanece, então a mensagem continua acionável. Encaminhe-a para os operadores.
O envelope CloudEvents
Toda mensagem viaja em modo binário do CloudEvents, versão 1.0. Os atributos de contexto vão como cabeçalhos AMQP.
Uma implantação single-tenant também carrega um valor de tenant. Ela emite o literal
single-tenant, de modo que um mesmo consumidor atende às duas formas de implantação com o mesmo código.
Configuração do source de CloudEvents
STREAMING_CLOUDEVENTS_SOURCE não tem valor padrão. O Worker o exige sempre que o streaming está ligado, e para na inicialização quando o valor está vazio. O exemplo distribuído usa //lerian.fetcher/worker.
O Fetcher copia o valor para ce-source sem alteração. Dê a cada implantação do Worker o seu próprio valor de source quando vários produtores compartilharem um broker, e roteie por esse cabeçalho.
Contrato de entrega
A entrega é at-least-once. Deduplique por
ce-id.
- O Worker grava o evento em um outbox durável antes de publicar. Uma queda do broker atrasa o evento, não o perde.
- Um reparador varre eventos terminais que nunca foram publicados e os reemite a cada 30 segundos.
- O
ce-idpermanece idêntico em toda reemissão do mesmo job e status. Nada mais é estável o bastante para servir de chave. - A ordem não é garantida. Dois jobs podem terminar em uma ordem e chegar em outra.
- O Worker registra o status terminal do job antes de emitir.
GET /v1/fetcher/{id}continua sendo a autoridade sobre o estado do job.
Verificar o que você recebe
Duas assinaturas independentes protegem o caminho entre o Worker e você. As duas usam HMAC-SHA256 com a chave HMAC externa. HKDF-SHA256 deriva essa chave da chave-mestra
APP_ENC_KEY. O repositório do Fetcher traz uma pequena ferramenta que imprime essa chave, de modo que um consumidor verifica as assinaturas sem nunca guardar a chave-mestra.
A mensagem. O Worker assina cada mensagem publicada e carimba três cabeçalhos nela: x-message-signature, t para o timestamp da assinatura, e signature-version. O payload assinado amarra o timestamp, a versão da assinatura, o tenant, o identificador do job, o exchange e a routing key, seguidos do corpo da mensagem. Um replay do mesmo corpo sob outro tenant ou outra rota falha na verificação.
O resultado. result.hmac é o HMAC-SHA256 com chave sobre o JSON do resultado em texto claro, calculado antes da criptografia. O bloco integrity declara o mesmo valor junto com o algoritmo. Verifique-o depois de descriptografar e antes de confiar nas linhas.
O bloco protection descreve os bytes armazenados: encrypted é verdadeiro, o adaptador de armazenamento aplicou a criptografia, e o modo é adapter-managed. Ele descreve apenas o resultado, nunca as credenciais da fonte de dados.
Próximos passos
Jobs de extração
O que um job pede, e os estados pelos quais ele passa.
API REST do Fetcher
Criar um job, ler um job e gerenciar conexões.
Configuração
As variáveis de streaming, exchange e criptografia por trás desses eventos.
Arquitetura
O Manager, o Worker e o Engine que os dois executam.

