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
| Forma | Significa | Handlers |
|---|---|---|
@Task | comando: "faça isto" | exatamente um |
@OnEvent | fato: "isto aconteceu" | zero ou mais; cada assinante recebe sua própria entrega |
@Schedule | relógio: dispara uma task no horário | nenhum (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ção | Significa |
|---|---|
name | Nome da unidade: segmentos separados por ponto de [a-z][a-z0-9_]*, até 64 caracteres |
queue | critical, default, bulk ou maintenance |
concurrency | Execuções simultâneas por processo |
timeoutSeconds | Prazo da execução |
context | tenant (padrão), platform ou system |
overlap | Só 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
| Unidade | Tipo | Fila | O que faz |
|---|---|---|---|
maintenance.purge | @Task + @Schedule("rate(1 hour)") | maintenance | Remove 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 v1 | critical | E-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ável | Significa |
|---|---|
QUEUES | Filas a consumir, separadas por vírgula (obrigatória) |
MONALISA_RELAY | true em exatamente um serviço por ambiente: ele move o outbox para a fila |
MONALISA_LOCAL_SCHEDULER | true para disparar agendamentos no processo; só desenvolvimento |
MONALISA_SQS_REGION, MONALISA_SQS_QUEUE_URL_PREFIX, MONALISA_SQS_ENDPOINT | Onde estão as filas (o nome da fila é anexado ao prefixo) |
MONALISA_SQS_LOCAL_CREDENTIALS | true para credenciais locais fixas em vez da cadeia da AWS |
MONALISA_SQS_WAIT_SECONDS | Long polling por recebimento, 1 a 20 (padrão 20) |
MONALISA_STOP_TIMEOUT_MS | Orçamento de parada; a drenagem recebe esse valor menos 5 s e precisa superar o long polling |
MONALISA_HEALTH_PORT | Porta do /healthz (padrão 8081) |
DATABASE_POOL_SIZE | Pelo 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.