Outbox, relay e filas
Como a API agenda trabalho sem falar com a fila, e como ele chega aos workers.
O trabalho em segundo plano roda em apps/workers, um executável separado sobre framework/jobs. A API nunca roda
tasks nem fala com a fila: ela grava o pedido no PostgreSQL, na mesma transação do caso de uso, e um relay o leva até
o SQS depois do commit.
Outbox: a verdade
work.dispatch(...) e work.publish(...) gravam uma linha em outbox_messages (ou platform_outbox_messages, para
trabalho sem tenant) e a linha da run em task_runs com status queued, na transação do caso de uso. Se o commit
falha, não existe mensagem; se o commit passa, a mensagem existe mesmo que o processo morra no instante seguinte.
run_id(uuid) é a identidade de tudo: a linha do outbox, o envelope na fila e a run.- Atraso de qualquer duração vira
available_at: o relay só envia o que já venceu, semDelaySeconds. - Um evento com vários assinantes vira uma linha de outbox por assinante, cada uma com seu
run_id; a falha de um assinante não afeta os outros. A API descobre os assinantes pelo catálogotasks.json, sem importar o código deles.
Valkey: só um acelerador
Depois do commit, a API faz LPUSH na chave monalisa:outbox:wake. O relay espera nela com BLPOP e acorda em
milissegundos. Um sinal perdido não perde nada: o relay também varre o outbox a cada 5 s.
Relay
O relay roda em exatamente um serviço de workers por ambiente (MONALISA_RELAY=true). Ele lê lotes com
WHERE sent_at IS NULL AND available_at <= clock_timestamp() … FOR UPDATE SKIP LOCKED, envia com SendMessageBatch
(até 10 mensagens e 256 KiB por chamada) e marca sent_at na mesma transação do lote. Falha ao enviar soma
attempts e deixa a linha pendente. Um lote cheio dispara outra passada na hora, então um acúmulo drena na velocidade
da fila. Como o consumidor é idempotente, reenviar depois de uma falha parcial é seguro.
Filas
| Fila | Para quê | Visibilidade |
|---|---|---|
critical | o que alguém espera agora | 60 s |
default | trabalho comum | 120 s |
bulk | importações, relatórios | 15 min |
maintenance | limpeza e rotinas | 5 min |
Cada fila tem uma dead-letter queue. Um processo consome as filas listadas em QUEUES. Enquanto uma mensagem está
em processamento, um heartbeat estende a visibilidade.
Falhas
- Permanente (validação, proibido, não encontrado…): o runtime envia a mensagem para a dead-letter queue com o
contexto do erro, registra o evento na trilha e a run vira
failed. - Transitória: a mensagem volta para a fila com backoff exponencial (pela visibilidade, sem segurar o slot) e a run
vira
retrying. No 5º recebimento, segue o caminho da falha permanente. - Processo morre no meio: a mensagem reaparece pela visibilidade. A redrive nativa das filas manda para a dead-letter queue depois de 7 recebimentos (5 do runtime mais 2 de margem).
- Desligamento (SIGTERM): o worker para de buscar, sai de pronto, espera o que está em andamento até o orçamento
de parada, e o que não terminou volta para a fila com a run em
queued.
Tabelas do outbox
outbox_messages
Outbox de tenant: cada dispatch ou evento publicado vira uma linha na transação do use case; sent_at nulo significa pendente.
| Coluna | Tipo | Nulo | Observação |
|---|---|---|---|
id | bigint | não | identity (GENERATED ALWAYS) |
tenant_id | bigint | não | |
run_id | uuid | não | |
queue | text | não | |
envelope | jsonb | não | |
available_at | timestamptz | não | |
sent_at | timestamptz | sim | |
attempts | integer | não | padrão 0 |
created_at | timestamptz | não | padrão clock_timestamp() |
Chaves estrangeiras
tenant_id→tenants.id
Garantido pelo banco
| Tipo | Nome | Definição |
|---|---|---|
| CHECK | outbox_messages_attempts_check | CHECK ((attempts >= 0)) |
| CHECK | outbox_messages_envelope_check | CHECK ((jsonb_typeof(envelope) = 'object')) |
| CHECK | outbox_messages_queue_check | CHECK ((queue IN ('critical', 'default', 'bulk', 'maintenance'))) |
| PK | outbox_messages_pkey | PRIMARY KEY (id) |
| UNIQUE | outbox_messages_run_id_key | UNIQUE (run_id) |
Fora do banco
- Linhas enviadas há mais de 7 dias são removidas pela task
maintenance.purge; pendentes nunca.
platform_outbox_messages
Outbox do trabalho sem tenant (plataforma, sistema, agendamentos), na mesma forma.
| Coluna | Tipo | Nulo | Observação |
|---|---|---|---|
id | bigint | não | identity (GENERATED ALWAYS) |
run_id | uuid | não | |
queue | text | não | |
envelope | jsonb | não | |
available_at | timestamptz | não | |
sent_at | timestamptz | sim | |
attempts | integer | não | padrão 0 |
created_at | timestamptz | não | padrão clock_timestamp() |
Chaves estrangeiras
- nenhuma
Garantido pelo banco
| Tipo | Nome | Definição |
|---|---|---|
| CHECK | platform_outbox_messages_attempts_check | CHECK ((attempts >= 0)) |
| CHECK | platform_outbox_messages_envelope_check | CHECK ((jsonb_typeof(envelope) = 'object')) |
| CHECK | platform_outbox_messages_queue_check | CHECK ((queue IN ('critical', 'default', 'bulk', 'maintenance'))) |
| PK | platform_outbox_messages_pkey | PRIMARY KEY (id) |
| UNIQUE | platform_outbox_messages_run_id_key | UNIQUE (run_id) |
Fora do banco
- Linhas enviadas há mais de 7 dias são removidas pela task
maintenance.purge; pendentes nunca.