Monalisa
Background work

Tasks, eventos e agendamentos

Os decorators @Task, @OnEvent e @Schedule, o catálogo tasks.json e a configuração dos workers.

Três formas de trabalho, um runtime

FormaSignificaHandlers
@Taskcomando: "faça isto"exatamente um
@OnEventfato: "isto aconteceu"zero ou mais; cada assinante recebe sua própria entrega
@Schedulerelógio: dispara uma task no horárionenhum (só despacha a task)
@Schedule("rate(1 hour)")
@Task({
  name: "maintenance.purge",
  queue: "maintenance",
  context: "system",
  overlap: "skip",
  timeoutSeconds: 900,
})
export class MaintenancePurge implements TaskHandler<EmptyPayload> {
  parsePayload(): EmptyPayload { return {}; }
  async run(_payload: EmptyPayload, run: TaskRun): Promise<PurgeOutcome> { /* … */ }
}

@OnEvent({
  type: "identity.session_started",
  version: 1,
  name: "identity.login_notification_email",
  queue: "critical",
})
export class LoginNotificationEmail implements EventHandler<SessionStartedPayload> { /* … */ }

// Dentro de um caso de uso, na transação dele (GenerateReport e InvitationCreated são ilustrativos):
await work.dispatch(GenerateReport, { reportId }, { delaySeconds: 30 });
await work.publish(new InvitationCreated({ invitationId }));

Opções de @Task:

OpçãoSignifica
nameNome da unidade: segmentos separados por ponto de [a-z][a-z0-9_]*, até 64 caracteres
queuecritical, default, bulk ou maintenance
concurrencyExecuções simultâneas por processo
timeoutSecondsPrazo da execução
contexttenant (padrão), platform ou system
overlapSó para agendadas: skip (padrão) pula uma execução sobreposta, allow deixa rodar

Despachar uma task cujo context não bate com quem chama é recusado (RUN_CONTEXT_MISMATCH): uma task tenant precisa de tenant, e uma platform ou system não pode ter.

O handler recebe um TaskRun, nunca o envelope cru: runId, name, attempt, scheduledFor, now(), logger, progress(percent) e transaction. transaction é a transação do efeito, já com o tenant da run e app.system_process = <nome da task>: o que o handler grava por ela faz commit junto com os eventos started e succeeded da run.

Unidades atuais

UnidadeTipoFilaO que faz
maintenance.purge@Task + @Schedule("rate(1 hour)")maintenanceRemove em lotes chaves de idempotência vencidas e linhas de outbox enviadas há mais de 7 dias. Nunca toca nas runs nem na trilha.
identity.login_notification_email@OnEvent de identity.session_started v1criticalE-mail de "novo login". Continua registrado, mas nada publica o evento desde que o login antigo saiu.

O catálogo tasks.json

tasks.json lista cada task (nome, fila, context, agendamento) e cada assinatura. É gerado da descoberta das unidades (*.task.ts e *.handler.ts), nunca commitado: a imagem da API o recebe no build e o carrega de MONALISA_TASKS_CATALOG; sem ele a API não sobe. bun run tasks:check valida que a geração funciona (descoberta, nomes únicos, agendamentos válidos).

Agendamentos

Em desenvolvimento, MONALISA_LOCAL_SCHEDULER=true dispara os agendamentos dentro do processo; isso é recusado com NODE_ENV=production. Em staging a imagem roda em produção e não há agendador externo, então maintenance.purge não roda lá. Uma run agendada tem origem única por (agendamento, horário previsto), o que impede disparo duplicado.

Configuração dos workers

VariávelSignifica
QUEUESFilas a consumir, separadas por vírgula (obrigatória)
MONALISA_RELAYtrue em exatamente um serviço por ambiente: ele move o outbox para a fila
MONALISA_LOCAL_SCHEDULERtrue para disparar agendamentos no processo; só desenvolvimento
MONALISA_SQS_REGION, MONALISA_SQS_QUEUE_URL_PREFIX, MONALISA_SQS_ENDPOINTOnde estão as filas (o nome da fila é anexado ao prefixo)
MONALISA_SQS_LOCAL_CREDENTIALStrue para credenciais locais fixas em vez da cadeia da AWS
MONALISA_SQS_WAIT_SECONDSLong polling por recebimento, 1 a 20 (padrão 20)
MONALISA_STOP_TIMEOUT_MSOrçamento de parada; a drenagem recebe esse valor menos 5 s e precisa superar o long polling
MONALISA_HEALTH_PORTPorta do /healthz (padrão 8081)
DATABASE_POOL_SIZEPelo menos 2 × slots + 2, sendo slots a soma da concorrência das unidades das filas consumidas

A regra do pool é verificada no boot: cada slot pode segurar ao mesmo tempo a transação do efeito e uma de registro, mais uma conexão para o relay e uma para o agendador. Um worker com pool pequeno demais não sobe.

Ordem de deploy

Quando uma versão adiciona uma task ou um assinante, publique os workers antes da API. Uma mensagem para uma unidade que os workers em execução ainda não conhecem volta para a fila com backoff (UNIT_NOT_REGISTERED) e só vai para a dead-letter queue depois de cinco recebimentos; os workers novos normalmente a pegam antes, mas a ordem elimina a corrida.

Nesta página