Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,8 @@ Hexagonal, com dependências sempre apontando para dentro:
- `config` — `AsyncJobsProperties` em pacote neutro (adapters não dependem de `autoconfigure`).
- `autoconfigure` — composition root; 4 auto-configurations em
`src/main/resources/META-INF/spring/*.imports`.
- `spi` — o que o projeto consumidor implementa (`JobHandler`, `AsyncJobHandler`,
`JobReporter`).
- `spi` — a fronteira com o projeto consumidor: ele **implementa** `JobHandler` /
`AsyncJobHandler`, e **injeta** `JobReporter` e `JobFreshness`.

Regras que já custaram bugs e devem ser preservadas:

Expand All @@ -91,6 +91,12 @@ Regras que já custaram bugs e devem ser preservadas:
8. **A lib não configura `DataSource` nem pool**, e não aplica DDL. O esquema é
do consumidor; `src/main/resources/async-jobs-schema.sql` é a referência, e é
o mesmo arquivo que os testes aplicam.
9. **Cancelar não interrompe a thread da rotina.** É cooperativo por decisão
(ADR 0004): abortar no meio deixaria a base do consumidor parcialmente
atualizada. Não introduza `Thread.interrupt()` nem `Future.cancel(true)`.
10. **Frescor depende de `coalesce-in-flight`**, porque a janela procura a última
carga concluída pela `coalescing_key` — que só é gravada quando o coalescing
está ligado. Ligar só o frescor era um no-op silencioso; hoje falha o startup.

## Testes

Expand Down
81 changes: 72 additions & 9 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -197,8 +197,11 @@ Possiveis respostas:
- `404 Not Found`: job inexistente.

> Nota: o cancelamento marca o job como `CANCELLED` e impede que uma conclusao
> posterior sobrescreva o estado, mas **nao interrompe** uma execucao ja em
> andamento — a rotina continua rodando ate o fim.
> posterior sobrescreva o estado, mas **nao interrompe** a thread da rotina.
> Interromper trabalho a meio caminho deixaria a base do consumidor em estado
> parcial, sem ninguem para consertar — entao quem decide parar e a propria
> rotina, consultando `ctx.isCancelled()` entre lotes
> ([cancelamento cooperativo](#cancelamento-cooperativo)).

## Estados do job

Expand Down Expand Up @@ -226,16 +229,21 @@ public class ManterContasQuentesHandler implements JobHandler {
}

@Override
public void handle() {
public void handle(JobContext ctx) {
// regra de negocio do projeto consumidor: varre as contas,
// avalia a aptidao de cada uma e grava o resultado na base
contas.reavaliarTodas();
for (var lote : contas.emLotes(500)) {
if (ctx.isCancelled()) {
return; // para num ponto consistente
}
contas.reavaliar(lote);
}
}
}
```

O `handle()` nao retorna nada: a rotina deixa o resultado na base do consumidor,
e a lib apenas registra que a carga terminou.
O `handle` nao retorna nada: a rotina deixa o resultado na base do consumidor, e
a lib apenas registra que a carga terminou.

Regras importantes:

Expand All @@ -256,6 +264,41 @@ Se o worker morrer sem reportar, a recuperacao automatica marca o job como falho
apos `async-jobs.recovery.processing-timeout`. Rotinas legitimamente longas devem
chamar `progress` periodicamente para renovar esse prazo.

### Cancelamento cooperativo

`ctx.isCancelled()` diz se o job foi cancelado enquanto a rotina roda. A lib nao
interrompe a thread de proposito: abortar no meio deixaria a base do consumidor
parcialmente atualizada, e so a rotina sabe onde e seguro parar.

- **Cada chamada le o storage** — pergunte entre lotes, nao a cada item.
- Rotinas que ja chamam `progress` periodicamente **nao precisam disto**: o
`false` devolvido por `progress` carrega a mesma informacao, sem leitura extra.
- Parar nao muda o estado: o job permanece `CANCELLED`. Uma rotina que retorna
normalmente depois de parar nao "descancela" o job — a transicao para
`COMPLETED` e recusada em estado terminal.

### Saber de quando sao os dados

Como a lib nao serve dados, nada impede o cliente de ler a base antes de a carga
terminar. O `JobFreshness` existe para a resposta de dominio poder dizer isso:

```java
@GetMapping("/contas/aptas")
ResponseEntity<ContasResponse> aptas() {
var contas = repository.buscarAptas(); // SQL de dominio, otimizado
return ResponseEntity.ok(new ContasResponse(
contas,
freshness.lastRefreshedAt("contas").orElse(null), // "dados de"
freshness.isFresh("contas"))); // dentro da janela?
}
```

- `lastRefreshedAt(type)`: quando a ultima carga concluiu — a idade real dos
dados. Vazio se nenhuma carga concluiu; uma carga que falhou nao conta.
- `isFresh(type)`: se essa conclusao esta dentro da janela configurada. Sempre
`false` sem janela configurada, pela mesma razao que toda submissao dispara
carga nesse caso.

## Politicas implementadas

- **Polling hint**: respostas usam `Retry-After` para orientar quando o client deve consultar novamente.
Expand All @@ -267,6 +310,8 @@ chamar `progress` periodicamente para renovar esse prazo.
- **Eventos em tempo real**: `GET /jobs/{id}/events` (SSE) entrega snapshot + transicoes. Faz parte do contrato, sem flag para desligar.
- **Execucao isolada**: os jobs rodam em um executor proprio da lib com **threads virtuais**, nunca no executor default da aplicacao; o limite de concorrencia (`async-jobs.processing.concurrency-limit`) e o backpressure.
- **Recuperacao de jobs orfaos**: a varredura reenfileira jobs que ficaram `PENDING` (instancia caiu antes de processar) e falha `PROCESSING` sem atualizacao ha muito tempo (worker morreu sem reportar). Ela so age sobre `type`s registrados na instancia, para nao interferir em jobs de outra aplicacao no mesmo banco.
- **Cancelamento cooperativo**: `ctx.isCancelled()` deixa a rotina parar num ponto consistente. A lib nao interrompe a thread — ver a secao da SPI.
- **Frescura consultavel**: `JobFreshness` permite ao endpoint de dominio dizer de quando sao os dados que esta devolvendo.

Como o submit nao recebe parametros, o dedupe/single-flight e por `type`.
Operacoes logicamente distintas devem usar `type`s distintos.
Expand All @@ -280,6 +325,7 @@ segue quente e a lib devolve aquele job em vez de disparar uma nova execucao.
```yaml
async-jobs:
retention: PT6H
coalesce-in-flight: true
freshness:
enabled: true
default-window: PT1H
Expand All @@ -288,6 +334,10 @@ async-jobs:
```

- E **opt-in**: sem `freshness.enabled=true`, todo submit dispara carga nova.
- Exige `coalesce-in-flight=true`, e o startup falha se so o frescor for ligado:
a janela localiza a ultima carga concluida pelo escopo de coalescing, que so e
gravado quando o coalescing esta ligado. Sem a validacao, ligar so o frescor
seria um no-op silencioso.
- Sem janela configurada para o `type` (nem `default-window`), o frescor nao se aplica.
- So conta job `COMPLETED`: uma carga que falhou nao suprime a proxima tentativa.
- `Cache-Control: no-cache` no submit ignora a janela e forca uma carga nova.
Expand Down Expand Up @@ -343,7 +393,7 @@ Parametros proprios:
| `async-jobs.retention` | `PT1H` | Por quanto tempo o registro de controle do job segue relevante. Alimenta o header `Expires` (a partir da ultima atualizacao) e o timeout do stream SSE. Aceita formato `Duration` do Spring, como `PT10M`, `PT1H` ou `P1D`. |
| `async-jobs.retry-after-seconds` | `5` | Hint enviado no header `Retry-After` em submissao e consulta de status enquanto o job esta ativo. |
| `async-jobs.coalesce-in-flight` | `false` | Quando `true`, chamadas equivalentes enquanto um job ainda esta ativo reutilizam o mesmo job em andamento em vez de criar outro. |
| `async-jobs.freshness.enabled` | `false` | Liga a janela de frescor: uma carga concluida dentro da janela dispensa carga nova. |
| `async-jobs.freshness.enabled` | `false` | Liga a janela de frescor: uma carga concluida dentro da janela dispensa carga nova. Exige `coalesce-in-flight=true`. |
| `async-jobs.freshness.default-window` | — | Janela aplicada aos `type`s sem configuracao propria. Sem valor, o frescor nao se aplica a eles. |
| `async-jobs.freshness.per-type.<type>` | — | Janela especifica de um `type`, sobrepondo a default. |
| `async-jobs.processing.concurrency-limit` | `256` | Jobs processados simultaneamente no executor proprio da lib (threads virtuais). Ao saturar, a submissao aguarda vaga — backpressure em vez de acumulo ilimitado. |
Expand Down Expand Up @@ -404,10 +454,12 @@ A suite de testes registra rotinas de exemplo e cobre:
- `404 Not Found` para job inexistente, inclusive quando o id nem tem forma de UUID (nao `500`);
- single-flight (coalescing) por indice unico, inclusive com submits concorrentes;
- janela de frescor: carga suprimida com dado quente, refeita com `Cache-Control: no-cache`;
- frescura consultavel: `lastRefreshedAt` ignora carga ativa e carga que falhou, e `isFresh` e falso sem janela configurada;
- cancelamento cooperativo: a rotina em execucao ve `isCancelled()` virar `true`, para no meio, e o `complete` seguinte e recusado;
- job com falha (`422` + Problem Detail no status), inclusive com titulo em branco;
- report recusado em job terminal (`JobReporter` devolvendo `false`);
- recuperacao de jobs orfaos: reenfileiramento de `PENDING`, falha de `PROCESSING` zumbi e nao-interferencia em `type` de outra aplicacao;
- validacao de configuracao: frescor maior que a retencao falha o startup;
- validacao de configuracao: frescor maior que a retencao, ou frescor sem coalescing, falham o startup;
- timestamps e `Expires` determinísticos com um `Clock` fixo injetado;
- fluxo fire-and-forget via `JobReporter`.

Expand All @@ -428,6 +480,17 @@ Ultima verificacao local:
./mvnw clean test
```

Resultado: `Tests run: 123, Failures: 0, Errors: 0, Skipped: 0`.
Resultado: `Tests run: 143, Failures: 0, Errors: 0, Skipped: 0`.

Para rodar só os rápidos: `./mvnw test -Dgroups='!integration'`.

## Evolucao registrada

O dispatch e in-process: o job roda na instancia que recebeu o `POST`, e a
varredura de recuperacao e a rede de segurancia. A alternativa — a propria tabela
como fila, via `FOR UPDATE SKIP LOCKED` — esta desenhada no
`docs/adr/0005-dispatch-por-fila-na-propria-tabela.md`, junto com o critério de
quando vale implementar. Resumo: ganha distribuicao real entre instancias, custa
latencia de um ciclo de poll e uma peca viva a mais; para uma rotina de
reaquecimento com single-flight ligado (um job ativo por escopo), nao ha fila a
balancear.
14 changes: 12 additions & 2 deletions docs/adr/0004-storage-em-aurora-postgresql-e-escopo-de-controle.md
Original file line number Diff line number Diff line change
Expand Up @@ -189,14 +189,24 @@ Desvios que permanecem, conscientes:
Aurora existe uma evolução natural — a própria tabela `async_jobs` como fila via
`FOR UPDATE SKIP LOCKED`, permitindo que qualquer instância puxe `PENDING` e
tornando a varredura o mecanismo principal de distribuição em vez de rede de
segurança. Fora do escopo desta ADR; registrado como opção.
segurança. → **detalhado no ADR 0005**, com o critério de quando implementar;
permanece desvio consciente por ora.
- **O gate deixa de ser garantia**: com a leitura indo para um endpoint de domínio
compartilhado, "só leia quando terminar" passa a ser conselho — nada impede o
cliente de ler dado morno antes. É inerente à semântica de cache quente; por
isso a frescura deve ser exposta na resposta de domínio.
isso a frescura deve ser exposta na resposta de domínio. → **endereçado**: a SPI
`JobFreshness` (`lastRefreshedAt`/`isFresh`) permite ao endpoint de domínio
declarar de quando são os dados. O desvio deixa de ser silencioso: continua
possível ler dado morno, mas o cliente agora consegue saber que é.
- **Cancelar não interrompe a rotina em andamento**, e o doc pede para avaliar
rollback parcial ou transação compensatória. Para reaquecimento não há o que
compensar; se surgir rotina com efeito colateral externo, reabrir.
→ **endereçado de forma cooperativa**: `JobContext.isCancelled()` deixa a
rotina parar num ponto consistente que ela escolhe. Não interrompemos a thread
de propósito — abortar no meio deixaria a base do consumidor parcialmente
atualizada, sem ninguém para consertar; a rotina é quem sabe onde é seguro
parar. Rotinas que já reportam progresso têm o mesmo sinal de graça, pelo
`false` do `progress`.
- **SSE como complemento, não alternativa**: o doc posiciona SSE na seção de
"quando este padrão não é adequado". Mantemos o polling como baseline e SSE como
transporte adicional — superconjunto deliberado, não contradição.
Expand Down
116 changes: 116 additions & 0 deletions docs/adr/0005-dispatch-por-fila-na-propria-tabela.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
# ADR 0005 — Dispatch por fila na própria tabela (`FOR UPDATE SKIP LOCKED`)

- **Status:** Proposto (não implementado)
- **Data:** 2026-07
- **Contexto:** desvio consciente registrado no ADR 0004 — o offload do padrão

## Contexto

O padrão Asynchronous Request-Reply pede que a requisição seja repassada "para
outro componente, como uma fila". Hoje o `SubmitJobUseCase` chama
`processor.process(job)`, que é `@Async` no executor da própria lib. Consequência:
**o job fica amarrado à instância que recebeu o `POST`**.

Isso produz três problemas, em ordem de gravidade:

1. **Distribuição desigual.** Se o load balancer manda dois `POST` para a mesma
instância, os dois jobs rodam lá, mesmo com outra instância ociosa. Não há
como uma instância "puxar" trabalho.
2. **Recuperação como caminho principal, não como rede de segurança.** Quando a
instância cai, o job só volta a andar quando o `StaleJobRecoveryService` o
encontra — depois de `redispatch-after` (default `PT1M`). O que deveria ser
exceção virou o mecanismo normal de redistribuição.
3. **Backpressure que rejeita em vez de enfileirar.** Ao saturar o
`concurrency-limit`, o submit espera vaga; se o executor recusar, o job fica
`PENDING` até a varredura. Já é tratado (a submissão devolve `202` de todo
jeito), mas é acidente, não desenho.

Com o storage relacional, a própria tabela `async_jobs` pode ser a fila — não é
preciso introduzir SQS, RabbitMQ ou Kafka para resolver isso.

## Decisão proposta

Inverter o dispatch de **push** para **pull**. Cada instância roda um laço que
reclama jobs `PENDING` com `FOR UPDATE SKIP LOCKED`:

```sql
update async_jobs
set status = 'PROCESSING', last_updated_at = now()
where id in (
select id from async_jobs
where status = 'PENDING'
order by created_at
limit :batch
for update skip locked)
returning id, type;
```

`SKIP LOCKED` é o ponto central: duas instâncias rodando este statement ao mesmo
tempo **não bloqueiam uma à outra e não pegam as mesmas linhas** — a segunda
simplesmente pula as linhas travadas pela primeira. É o mesmo mecanismo que
bibliotecas de fila em Postgres usam, e evita tanto o lock global quanto o
processamento em duplicidade.

O `SubmitJobUseCase` passa a apenas persistir `PENDING` e devolver `202`. O
`processor.process(job)` direto some do caminho de submissão.

## Consequências

**A favor:**

- Qualquer instância puxa qualquer job: distribuição real, sem afinidade com quem
atendeu o `POST`.
- A instância que cai deixa de ser um problema de recuperação e passa a ser um
não-evento: a transação não commitou, o lock caiu, e o job volta a ser visível
para os outros no próximo ciclo.
- O `StaleJobRecoveryService` volta a ser o que o nome diz — rede de segurança
para o zumbi em `PROCESSING`, não distribuidor de trabalho.
- O `202` deixa de depender de o executor local aceitar a tarefa.

**Contra:**

- **Latência mínima de um ciclo de poll.** Hoje o dispatch é imediato; com fila
passa a esperar até o intervalo do laço. Para rotinas de minutos (a premissa do
padrão) é irrelevante, mas piora a experiência dos testes e de jobs triviais.
- **Carga constante no banco.** Um `SELECT ... FOR UPDATE SKIP LOCKED` por
instância por ciclo, mesmo sem trabalho. Mitigável com backoff quando a fila
vem vazia, e é uma tabela pequena com índice parcial.
- **Mais uma peça viva.** Um laço por instância, com shutdown ordenado, que não
existe hoje.
- **Atenção a autovacuum.** A tabela passa a ter `UPDATE` frequente; se o volume
crescer, dá bloat.

## Por que não agora

O ganho é de **distribuição e resiliência sob múltiplas instâncias**, e o custo é
latência mais uma peça viva no ciclo de vida. Para o caso de uso que originou a
lib — uma rotina de reaquecimento por `type`, com single-flight ligado, ou seja
**um job ativo por escopo de cada vez** — a distribuição não é o gargalo: não há
fila de trabalho para balancear.

Vale implementar quando aparecer pelo menos um destes:

- mais de um job ativo por instância de forma rotineira (single-flight desligado,
ou muitos `type`s distintos);
- jobs perdendo tempo em `PENDING` porque a instância que os aceitou está
saturada enquanto outras estão ociosas;
- necessidade de que a queda de uma instância seja transparente, sem esperar
`redispatch-after`.

Até então, o dispatch in-process com varredura de recuperação entrega a mesma
garantia funcional — nenhum job fica órfão — com menos partes móveis.

## Alternativas consideradas

- **Broker dedicado (SQS/Rabbit/Kafka).** Resolve o mesmo problema e adiciona
infraestrutura, um segundo storage a manter consistente com o banco, e o
problema de outbox. O ADR 0004 acabou de remover um storage; reintroduzir outro
para dispatch anda na direção contrária.
- **`LISTEN/NOTIFY` para acordar o laço** e evitar o poll constante. Atrativo,
mas não passa por RDS Proxy, perde notificação em failover e exige conexão
dedicada por instância. Serviria como otimização de latência **em cima** da
fila, nunca como o mecanismo de reivindicação — a garantia tem que vir do
`SKIP LOCKED`.
- **Advisory locks** (`pg_try_advisory_lock`) por job. Funciona, mas o lock vive
fora da linha e não aparece em `SELECT`; diagnosticar "quem está com este job"
fica pior do que com estado na própria tabela.
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,11 @@ where status in ('PENDING', 'PROCESSING') and last_updated_at < :olderThan
limit :limit
""";

private static final String SELECT_LAST_COMPLETED_AT = """
select max(last_updated_at) from async_jobs
where coalescing_key = :coalescingKey and status = 'COMPLETED'
""";

private static final String SELECT_FRESH_COMPLETED = """
select * from async_jobs
where coalescing_key = :coalescingKey
Expand Down Expand Up @@ -247,6 +252,19 @@ public Optional<Job> findFreshCompleted(String coalescingKey, Instant completedA
.optional();
}

@Override
public Optional<Instant> findLastCompletedAt(String coalescingKey) {
if (coalescingKey == null) {
return Optional.empty();
}
// max() sempre devolve uma linha; sem carga concluida, o valor e null
return jdbc.sql(SELECT_LAST_COMPLETED_AT)
.param("coalescingKey", coalescingKey)
.query(OffsetDateTime.class)
.optional()
.map(OffsetDateTime::toInstant);
}

/**
* Check-and-set em um único statement: o {@code WHERE} carrega os estados de
* origem permitidos, e o {@code RETURNING} devolve o instante gravado quando
Expand Down
Loading