Monalisa
Background work

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, sem DelaySeconds.
  • 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álogo tasks.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

FilaPara quêVisibilidade
criticalo que alguém espera agora60 s
defaulttrabalho comum120 s
bulkimportações, relatórios15 min
maintenancelimpeza e rotinas5 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.

ColunaTipoNuloObservação
idbigintnãoidentity (GENERATED ALWAYS)
tenant_idbigintnão
run_iduuidnão
queuetextnão
envelopejsonbnão
available_attimestamptznão
sent_attimestamptzsim
attemptsintegernãopadrão 0
created_attimestamptznãopadrão clock_timestamp()

Chaves estrangeiras

  • tenant_id → tenants.id

Garantido pelo banco

TipoNomeDefinição
CHECKoutbox_messages_attempts_checkCHECK ((attempts >= 0))
CHECKoutbox_messages_envelope_checkCHECK ((jsonb_typeof(envelope) = 'object'))
CHECKoutbox_messages_queue_checkCHECK ((queue IN ('critical', 'default', 'bulk', 'maintenance')))
PKoutbox_messages_pkeyPRIMARY KEY (id)
UNIQUEoutbox_messages_run_id_keyUNIQUE (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.

ColunaTipoNuloObservação
idbigintnãoidentity (GENERATED ALWAYS)
run_iduuidnão
queuetextnão
envelopejsonbnão
available_attimestamptznão
sent_attimestamptzsim
attemptsintegernãopadrão 0
created_attimestamptznãopadrão clock_timestamp()

Chaves estrangeiras

  • nenhuma

Garantido pelo banco

TipoNomeDefinição
CHECKplatform_outbox_messages_attempts_checkCHECK ((attempts >= 0))
CHECKplatform_outbox_messages_envelope_checkCHECK ((jsonb_typeof(envelope) = 'object'))
CHECKplatform_outbox_messages_queue_checkCHECK ((queue IN ('critical', 'default', 'bulk', 'maintenance')))
PKplatform_outbox_messages_pkeyPRIMARY KEY (id)
UNIQUEplatform_outbox_messages_run_id_keyUNIQUE (run_id)

Fora do banco

  • Linhas enviadas há mais de 7 dias são removidas pela task maintenance.purge; pendentes nunca.

Nesta página