diff --git a/AGENTS.md b/AGENTS.md
index 2a9670b..b5a83ec 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -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:
@@ -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
diff --git a/README.md b/README.md
index 8ff9acd..c110489 100644
--- a/README.md
+++ b/README.md
@@ -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
@@ -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:
@@ -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 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.
@@ -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.
@@ -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
@@ -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.
@@ -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.` | — | 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. |
@@ -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`.
@@ -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.
diff --git a/docs/adr/0004-storage-em-aurora-postgresql-e-escopo-de-controle.md b/docs/adr/0004-storage-em-aurora-postgresql-e-escopo-de-controle.md
index fa258d6..492060f 100644
--- a/docs/adr/0004-storage-em-aurora-postgresql-e-escopo-de-controle.md
+++ b/docs/adr/0004-storage-em-aurora-postgresql-e-escopo-de-controle.md
@@ -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.
diff --git a/docs/adr/0005-dispatch-por-fila-na-propria-tabela.md b/docs/adr/0005-dispatch-por-fila-na-propria-tabela.md
new file mode 100644
index 0000000..b654c82
--- /dev/null
+++ b/docs/adr/0005-dispatch-por-fila-na-propria-tabela.md
@@ -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.
diff --git a/src/main/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryAdapterOut.java b/src/main/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryAdapterOut.java
index d543d29..af573d6 100644
--- a/src/main/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryAdapterOut.java
+++ b/src/main/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryAdapterOut.java
@@ -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
@@ -247,6 +252,19 @@ public Optional findFreshCompleted(String coalescingKey, Instant completedA
.optional();
}
+ @Override
+ public Optional 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
diff --git a/src/main/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOut.java b/src/main/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOut.java
index 078ec5d..2a0f1d4 100644
--- a/src/main/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOut.java
+++ b/src/main/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOut.java
@@ -1,6 +1,7 @@
package com.async.request.reply.adapter.out.processing;
import com.async.request.reply.core.domain.Job;
+import com.async.request.reply.core.enums.JobStatus;
import com.async.request.reply.core.port.out.JobProcessorPortOut;
import com.async.request.reply.core.port.out.JobRepositoryPortOut;
import com.async.request.reply.spi.AsyncJobHandler;
@@ -63,8 +64,11 @@ public void process(Job job) {
try {
switch (routine) {
case JobHandler sync -> {
- sync.handle(); // o efeito é a escrita na base do consumidor
- repository.complete(job.getId()); // o retorno do handler é o sinal de conclusão
+ sync.handle(context(job)); // o efeito é a escrita na base do consumidor
+ // o retorno do handler é o sinal de conclusão; se ele parou por
+ // cancelamento, o UPDATE condicional recusa e o estado terminal
+ // permanece — não há "descancelar"
+ repository.complete(job.getId());
}
case AsyncJobHandler async -> runAsync(async, job);
default -> repository.fail(job.getId(), "Unsupported routine",
@@ -76,8 +80,28 @@ public void process(Job job) {
}
private void runAsync(AsyncJobHandler handler, Job job) {
- JobContext ctx = job::getId; // expõe apenas o jobId
- handler.start(ctx);
+ handler.start(context(job));
// NÃO completa: aguarda JobReporter.complete(jobId) do worker
}
+
+ /**
+ * Contexto da rotina. O {@code isCancelled} lê o estado a cada chamada — é o
+ * preço de não manter cache que possa mentir sobre um cancelamento recente.
+ */
+ private JobContext context(Job job) {
+ return new JobContext() {
+
+ @Override
+ public String jobId() {
+ return job.getId();
+ }
+
+ @Override
+ public boolean isCancelled() {
+ return repository.findById(job.getId())
+ .map(current -> current.getStatus() == JobStatus.CANCELLED)
+ .orElse(false);
+ }
+ };
+ }
}
diff --git a/src/main/java/com/async/request/reply/autoconfigure/AsyncJobsAutoConfiguration.java b/src/main/java/com/async/request/reply/autoconfigure/AsyncJobsAutoConfiguration.java
index 8850c4a..9e1554c 100644
--- a/src/main/java/com/async/request/reply/autoconfigure/AsyncJobsAutoConfiguration.java
+++ b/src/main/java/com/async/request/reply/autoconfigure/AsyncJobsAutoConfiguration.java
@@ -22,8 +22,10 @@
import com.async.request.reply.core.service.StaleJobRecoveryService;
import com.async.request.reply.core.usecase.CancelJobUseCase;
import com.async.request.reply.core.usecase.GetJobStatusUseCase;
+import com.async.request.reply.core.usecase.JobFreshnessUseCase;
import com.async.request.reply.core.usecase.JobReporterUseCase;
import com.async.request.reply.core.usecase.SubmitJobUseCase;
+import com.async.request.reply.spi.JobFreshness;
import com.async.request.reply.spi.JobReporter;
import com.async.request.reply.spi.Routine;
import org.springframework.boot.autoconfigure.AutoConfiguration;
@@ -163,6 +165,19 @@ JobReporter jobReporter(JobRepositoryPortOut repository) {
return new JobReporterUseCase(repository);
}
+ /**
+ * SPI de leitura: permite ao endpoint de domínio do consumidor dizer de
+ * quando são os dados que está devolvendo (ADR 0004).
+ */
+ @Bean
+ @ConditionalOnMissingBean
+ JobFreshness jobFreshness(JobRepositoryPortOut repository,
+ JobFreshnessPolicyPortOut freshnessPolicy,
+ CoalescingKeyPortOut coalescingKey,
+ Clock clock) {
+ return new JobFreshnessUseCase(repository, freshnessPolicy, coalescingKey, clock);
+ }
+
// --- recuperação de jobs órfãos ----------------------------------------
@Bean
diff --git a/src/main/java/com/async/request/reply/config/AsyncJobsProperties.java b/src/main/java/com/async/request/reply/config/AsyncJobsProperties.java
index 2d44ba6..577faa3 100644
--- a/src/main/java/com/async/request/reply/config/AsyncJobsProperties.java
+++ b/src/main/java/com/async/request/reply/config/AsyncJobsProperties.java
@@ -45,6 +45,26 @@ public record AsyncJobsProperties(
recovery = (recovery == null) ? new Recovery(null, null, null, null) : recovery;
requireFreshnessWithinRetention(retention, freshness);
+ requireCoalescingForFreshness(coalesceInFlight, freshness);
+ }
+
+ /**
+ * O frescor localiza a última carga concluída pela {@code coalescing_key}, e
+ * essa coluna só é gravada quando o coalescing está ligado. Sem esta
+ * validação, ligar apenas {@code freshness.enabled} seria um no-op
+ * silencioso: a configuração pede para poupar carga e nada acontece.
+ *
+ * Não é acidental que os dois andem juntos — ambos falam do mesmo escopo:
+ * coalescing cobre carga em andamento, frescor cobre carga já
+ * concluída.
+ */
+ private static void requireCoalescingForFreshness(boolean coalesceInFlight, Freshness freshness) {
+ if (freshness.enabled() && !coalesceInFlight) {
+ throw new IllegalArgumentException(
+ "async-jobs.freshness.enabled=true exige async-jobs.coalesce-in-flight=true: "
+ + "a janela de frescor localiza a ultima carga concluida pelo escopo de "
+ + "coalescing, que so e gravado quando o coalescing esta ligado.");
+ }
}
/**
diff --git a/src/main/java/com/async/request/reply/core/port/out/JobRepositoryPortOut.java b/src/main/java/com/async/request/reply/core/port/out/JobRepositoryPortOut.java
index a86bcac..cce686e 100644
--- a/src/main/java/com/async/request/reply/core/port/out/JobRepositoryPortOut.java
+++ b/src/main/java/com/async/request/reply/core/port/out/JobRepositoryPortOut.java
@@ -64,4 +64,11 @@ public interface JobRepositoryPortOut {
* próxima tentativa.
*/
Optional findFreshCompleted(String coalescingKey, Instant completedAfter);
+
+ /**
+ * Quando a última carga daquele escopo concluiu, sem filtro de janela.
+ * É a idade real dos dados, exposta ao consumidor via
+ * {@code JobFreshness} para a resposta de domínio poder dizer "dados de".
+ */
+ Optional findLastCompletedAt(String coalescingKey);
}
diff --git a/src/main/java/com/async/request/reply/core/usecase/JobFreshnessUseCase.java b/src/main/java/com/async/request/reply/core/usecase/JobFreshnessUseCase.java
new file mode 100644
index 0000000..dbf1a22
--- /dev/null
+++ b/src/main/java/com/async/request/reply/core/usecase/JobFreshnessUseCase.java
@@ -0,0 +1,52 @@
+package com.async.request.reply.core.usecase;
+
+import com.async.request.reply.core.port.out.CoalescingKeyPortOut;
+import com.async.request.reply.core.port.out.JobFreshnessPolicyPortOut;
+import com.async.request.reply.core.port.out.JobRepositoryPortOut;
+import com.async.request.reply.spi.JobFreshness;
+
+import java.time.Clock;
+import java.time.Instant;
+import java.util.Optional;
+
+/**
+ * Implementação do {@link JobFreshness}. Responde sobre o mesmo escopo que o
+ * {@code SubmitJobUseCase} usa para decidir se poupa uma carga — a chave de
+ * coalescing do {@code type} — para que a resposta de domínio e a decisão da lib
+ * nunca discordem sobre o que está quente.
+ */
+public class JobFreshnessUseCase implements JobFreshness {
+
+ private final JobRepositoryPortOut repository;
+ private final JobFreshnessPolicyPortOut policy;
+ private final CoalescingKeyPortOut coalescingKey;
+ private final Clock clock;
+
+ public JobFreshnessUseCase(JobRepositoryPortOut repository,
+ JobFreshnessPolicyPortOut policy,
+ CoalescingKeyPortOut coalescingKey,
+ Clock clock) {
+ this.repository = repository;
+ this.policy = policy;
+ this.coalescingKey = coalescingKey;
+ this.clock = clock;
+ }
+
+ @Override
+ public Optional lastRefreshedAt(String type) {
+ return repository.findLastCompletedAt(coalescingKey.keyFor(type));
+ }
+
+ @Override
+ public boolean isFresh(String type) {
+ Optional lastRefreshedAt = lastRefreshedAt(type);
+ if (lastRefreshedAt.isEmpty()) {
+ return false;
+ }
+ // sem janela nada e quente — a mesma razao pela qual toda submissao
+ // dispara carga; prometer frescura aqui seria mentir para o cliente
+ return policy.freshnessFor(type)
+ .map(window -> lastRefreshedAt.get().isAfter(clock.instant().minus(window)))
+ .orElse(false);
+ }
+}
diff --git a/src/main/java/com/async/request/reply/spi/JobContext.java b/src/main/java/com/async/request/reply/spi/JobContext.java
index 565a77f..00b66a2 100644
--- a/src/main/java/com/async/request/reply/spi/JobContext.java
+++ b/src/main/java/com/async/request/reply/spi/JobContext.java
@@ -1,11 +1,43 @@
package com.async.request.reply.spi;
/**
- * Contexto entregue a uma {@link AsyncJobHandler}. Expõe o {@code jobId}
- * gerado pela lib para que o projeto consumidor possa amarrá-lo ao seu
- * worker e, mais tarde, reportar o resultado via {@link JobReporter}.
+ * Contexto entregue à rotina em execução. Expõe o {@code jobId} gerado pela lib
+ * — para o projeto consumidor amarrá-lo ao seu worker e reportar depois via
+ * {@link JobReporter} — e o sinal de cancelamento.
*/
public interface JobContext {
String jobId();
+
+ /**
+ * Se o job foi cancelado enquanto a rotina roda. É a base do
+ * cancelamento cooperativo: a lib não interrompe a thread da rotina
+ * (interromper trabalho a meio caminho deixaria a base do consumidor em
+ * estado parcial, sem ninguém para consertar), então quem decide parar é a
+ * própria rotina — tipicamente entre lotes:
+ *
+ * {@code
+ * public void handle(JobContext ctx) {
+ * for (var lote : contas.emLotes(500)) {
+ * if (ctx.isCancelled()) {
+ * return; // para num ponto consistente
+ * }
+ * contas.reavaliar(lote);
+ * }
+ * }
+ * }
+ *
+ * Cada chamada lê o storage, então pergunte entre lotes, não a cada
+ * item. Rotinas que já chamam {@link JobReporter#progress} periodicamente não
+ * precisam disto: o {@code false} devolvido por {@code progress} carrega a
+ * mesma informação, sem leitura extra.
+ *
+ * Parar não muda o estado do job: ele permanece {@code CANCELLED}. Uma
+ * rotina que retorna normalmente após parar não corre risco de "descancelar"
+ * o job — a transição para {@code COMPLETED} é recusada em estado terminal.
+ *
+ * @return {@code false} também quando o job não é encontrado — não há o que
+ * parar por conta de um job que não existe mais
+ */
+ boolean isCancelled();
}
diff --git a/src/main/java/com/async/request/reply/spi/JobFreshness.java b/src/main/java/com/async/request/reply/spi/JobFreshness.java
new file mode 100644
index 0000000..ba324ea
--- /dev/null
+++ b/src/main/java/com/async/request/reply/spi/JobFreshness.java
@@ -0,0 +1,44 @@
+package com.async.request.reply.spi;
+
+import java.time.Instant;
+import java.util.Optional;
+
+/**
+ * Consulta de frescura, injetável pelo projeto consumidor no endpoint de domínio
+ * dele.
+ *
+ * Existe porque a lib não serve dados (ADR 0004): o cliente lê o resultado na
+ * base do consumidor, e nada o impede de ler antes de a carga terminar. O
+ * "só leia quando estiver quente" deixa de ser garantia e passa a ser
+ * informação — então a resposta de domínio precisa poder dizer de quando
+ * são os dados que está devolvendo.
+ *
+ * {@code
+ * @GetMapping("/contas/aptas")
+ * ResponseEntity 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?
+ * }
+ * }
+ */
+public interface JobFreshness {
+
+ /**
+ * Quando a última carga daquele {@code type} concluiu — a idade real dos
+ * dados. Vazio significa que nenhuma carga concluiu (nunca rodou, ainda está
+ * rodando, ou a última falhou); uma carga que falhou não deixa dado quente.
+ */
+ Optional lastRefreshedAt(String type);
+
+ /**
+ * Se a última carga concluiu dentro da janela configurada para o
+ * {@code type} — a mesma janela que faz a lib poupar uma carga nova.
+ *
+ * Sempre {@code false} quando não há janela configurada: sem janela, nada
+ * é considerado quente, e é por isso que toda submissão dispara carga.
+ */
+ boolean isFresh(String type);
+}
diff --git a/src/main/java/com/async/request/reply/spi/JobHandler.java b/src/main/java/com/async/request/reply/spi/JobHandler.java
index 1800f8c..dd0c698 100644
--- a/src/main/java/com/async/request/reply/spi/JobHandler.java
+++ b/src/main/java/com/async/request/reply/spi/JobHandler.java
@@ -1,19 +1,23 @@
package com.async.request.reply.spi;
/**
- * SPI síncrono: a rotina roda inteira dentro de {@link #handle()} e o retorno do
- * método é o sinal de conclusão — a lib marca o job como COMPLETED, ou FAILED se
- * a rotina lançar exceção.
+ * SPI síncrono: a rotina roda inteira dentro de {@link #handle(JobContext)} e o
+ * retorno do método é o sinal de conclusão — a lib marca o job como COMPLETED, ou
+ * FAILED se a rotina lançar exceção.
*
* Não devolve valor: a biblioteca controla execução, não serve dados
* (ADR 0004). O efeito da rotina é a escrita na base do próprio projeto
* consumidor, e o cliente lê o resultado pelo endpoint de domínio dele — com o
* SQL e os filtros dele — depois de saber que o job terminou.
*
+ * O {@link JobContext} traz o {@code jobId} e o sinal de cancelamento
+ * ({@link JobContext#isCancelled()}): rotinas longas devem consultá-lo entre
+ * lotes para poderem parar num ponto consistente.
+ *
* Use {@link AsyncJobHandler} quando a rotina apenas dispara o trabalho e a
* conclusão chega depois (evento/webhook), reportada via {@link JobReporter}.
*/
public interface JobHandler extends Routine {
- void handle();
+ void handle(JobContext ctx);
}
diff --git a/src/test/java/com/async/request/reply/AsynchronousRequestReplyPatternApplicationTests.java b/src/test/java/com/async/request/reply/AsynchronousRequestReplyPatternApplicationTests.java
index d45bf26..5d51a4c 100644
--- a/src/test/java/com/async/request/reply/AsynchronousRequestReplyPatternApplicationTests.java
+++ b/src/test/java/com/async/request/reply/AsynchronousRequestReplyPatternApplicationTests.java
@@ -21,6 +21,7 @@
import java.time.Duration;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import static org.hamcrest.Matchers.containsString;
@@ -48,16 +49,25 @@ class AsynchronousRequestReplyPatternApplicationTests extends PostgresContainerT
/** Gate que mantém o handler "lento" ativo até o teste liberar (sem sleep fixo). */
static final AtomicReference SLOW_GATE = new AtomicReference<>(new CountDownLatch(0));
+ /** Coordenação do teste de cancelamento cooperativo, sem sleep fixo. */
+ static final AtomicReference COOP_STARTED = new AtomicReference<>(new CountDownLatch(0));
+ static final AtomicReference COOP_GATE = new AtomicReference<>(new CountDownLatch(0));
+ static final AtomicReference COOP_SAW_CANCELLATION = new AtomicReference<>();
+
@BeforeEach
void setup() {
truncateJobs();
mvc = MockMvcBuilders.webAppContextSetup(wac).build();
SLOW_GATE.set(new CountDownLatch(1));
+ COOP_STARTED.set(new CountDownLatch(1));
+ COOP_GATE.set(new CountDownLatch(1));
+ COOP_SAW_CANCELLATION.set(null);
}
@AfterEach
void releaseSlowGate() {
SLOW_GATE.get().countDown();
+ COOP_GATE.get().countDown();
}
// -------------------------------------------------------------------------
@@ -70,7 +80,10 @@ static class TestHandlers {
JobHandler fastHandler() {
return new JobHandler() {
public String type() { return "test"; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: estes testes verificam o ciclo de vida
+ // do job (submit, status, headers), nao o efeito da rotina
+ }
};
}
@@ -78,7 +91,9 @@ public void handle() { }
JobHandler idempotentHandler() {
return new JobHandler() {
public String type() { return "idempotent-test"; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: só a Idempotency-Key importa aqui
+ }
};
}
@@ -86,7 +101,7 @@ public void handle() { }
JobHandler slowHandler() {
return new JobHandler() {
public String type() { return "cancel-test"; }
- public void handle() {
+ public void handle(JobContext ctx) {
// bloqueia até o teste liberar — mantém o job ativo de forma determinística
TestGate.await(SLOW_GATE.get());
}
@@ -104,19 +119,13 @@ public void start(JobContext ctx) {
};
}
- @Bean
- JobHandler nullHandler() {
- return new JobHandler() {
- public String type() { return "null-test"; }
- public void handle() { }
- };
- }
-
@Bean
JobHandler numericNameHandler() {
return new JobHandler() {
public String type() { return "123"; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: só o formato do type importa aqui
+ }
};
}
@@ -124,7 +133,9 @@ public void handle() { }
JobHandler dottedTypeHandler() {
return new JobHandler() {
public String type() { return "report.v1_all-items"; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: só o formato do type importa aqui
+ }
};
}
@@ -132,7 +143,20 @@ public void handle() { }
JobHandler failingHandler() {
return new JobHandler() {
public String type() { return "fail-test"; }
- public void handle() { throw new IllegalStateException("boom simulado"); }
+ public void handle(JobContext ctx) { throw new IllegalStateException("boom simulado"); }
+ };
+ }
+
+ /** Rotina que decide parar sozinha ao ver o cancelamento (ADR 0004). */
+ @Bean
+ JobHandler cooperativeHandler() {
+ return new JobHandler() {
+ public String type() { return "coop-cancel-test"; }
+ public void handle(JobContext ctx) {
+ COOP_STARTED.get().countDown();
+ TestGate.await(COOP_GATE.get()); // espera o teste cancelar
+ COOP_SAW_CANCELLATION.set(ctx.isCancelled());
+ }
};
}
}
@@ -341,6 +365,31 @@ void asyncHandlerCompletesViaReporter() throws Exception {
.andExpect(jsonPath("$.percentComplete", is(100)));
}
+ // 15. Cancelamento cooperativo: a rotina ve o cancelamento e para sozinha
+ @Test
+ void cancelledJobIsVisibleToTheRunningRoutine() throws Exception {
+ MvcResult post = mvc.perform(post("/jobs/coop-cancel-test"))
+ .andExpect(status().isAccepted()).andReturn();
+ String jobId = post.getResponse().getContentAsString().replaceAll(".*\"jobId\":\"([^\"]+)\".*", "$1");
+
+ // garante que a rotina esta rodando antes de cancelar
+ assertTrue(COOP_STARTED.get().await(10, TimeUnit.SECONDS), "a rotina nao comecou");
+
+ mvc.perform(delete("/jobs/{id}", jobId)).andExpect(status().isAccepted());
+ COOP_GATE.get().countDown(); // libera a rotina para consultar o contexto
+
+ Awaitility.await().atMost(Duration.ofSeconds(10)).pollInterval(Duration.ofMillis(50))
+ .until(() -> COOP_SAW_CANCELLATION.get() != null);
+ assertTrue(COOP_SAW_CANCELLATION.get(),
+ "a rotina precisa ver o cancelamento para poder parar num ponto consistente");
+
+ // a rotina retornou normalmente depois de parar; o complete que a lib
+ // tenta em seguida e recusado, e o estado terminal permanece
+ mvc.perform(get("/jobs/{id}/status", jobId))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.status", is("CANCELLED")));
+ }
+
// --- helpers de espera (polling com timeout em vez de Thread.sleep) -----
/** Aguarda o job concluir (status passa a reportar COMPLETED). */
diff --git a/src/test/java/com/async/request/reply/FixedClockJobTest.java b/src/test/java/com/async/request/reply/FixedClockJobTest.java
index 50ffc4d..8c07de9 100644
--- a/src/test/java/com/async/request/reply/FixedClockJobTest.java
+++ b/src/test/java/com/async/request/reply/FixedClockJobTest.java
@@ -1,5 +1,6 @@
package com.async.request.reply;
+import com.async.request.reply.spi.JobContext;
import com.async.request.reply.spi.JobHandler;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
@@ -73,7 +74,7 @@ Clock fixedClock() {
JobHandler clockHandler() {
return new JobHandler() {
public String type() { return "clock-test"; }
- public void handle() {
+ public void handle(JobContext ctx) {
TestGate.await(GATE.get());
}
};
diff --git a/src/test/java/com/async/request/reply/SingleFlightTest.java b/src/test/java/com/async/request/reply/SingleFlightTest.java
index 362a634..1a22210 100644
--- a/src/test/java/com/async/request/reply/SingleFlightTest.java
+++ b/src/test/java/com/async/request/reply/SingleFlightTest.java
@@ -1,5 +1,6 @@
package com.async.request.reply;
+import com.async.request.reply.spi.JobContext;
import com.async.request.reply.spi.JobHandler;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
@@ -57,7 +58,7 @@ static class Handlers {
JobHandler slowSfHandler() {
return new JobHandler() {
public String type() { return "sf-test"; }
- public void handle() {
+ public void handle(JobContext ctx) {
TestGate.await(GATE.get());
}
};
diff --git a/src/test/java/com/async/request/reply/StaleJobRecoveryTest.java b/src/test/java/com/async/request/reply/StaleJobRecoveryTest.java
index bc47fbc..fc83a18 100644
--- a/src/test/java/com/async/request/reply/StaleJobRecoveryTest.java
+++ b/src/test/java/com/async/request/reply/StaleJobRecoveryTest.java
@@ -4,6 +4,7 @@
import com.async.request.reply.core.enums.JobStatus;
import com.async.request.reply.core.port.out.JobRepositoryPortOut;
import com.async.request.reply.core.service.StaleJobRecoveryService;
+import com.async.request.reply.spi.JobContext;
import com.async.request.reply.spi.JobHandler;
import org.awaitility.Awaitility;
import org.junit.jupiter.api.BeforeEach;
@@ -55,7 +56,11 @@ static class Handlers {
JobHandler reapHandler() {
return new JobHandler() {
public String type() { return TYPE; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: o teste escreve o job orfao direto na
+ // tabela e chama recovery.recover() manualmente — a rotina
+ // so precisa existir para o type ser reconhecido
+ }
};
}
}
diff --git a/src/test/java/com/async/request/reply/TomcatEndToEndIntegrationTest.java b/src/test/java/com/async/request/reply/TomcatEndToEndIntegrationTest.java
index 2169f9c..a5df298 100644
--- a/src/test/java/com/async/request/reply/TomcatEndToEndIntegrationTest.java
+++ b/src/test/java/com/async/request/reply/TomcatEndToEndIntegrationTest.java
@@ -1,5 +1,6 @@
package com.async.request.reply;
+import com.async.request.reply.spi.JobContext;
import com.async.request.reply.spi.JobHandler;
import com.jayway.jsonpath.JsonPath;
import org.awaitility.Awaitility;
@@ -78,7 +79,10 @@ static class Handlers {
JobHandler e2eReportHandler() {
return new JobHandler() {
public String type() { return "e2e-report"; }
- public void handle() { }
+ public void handle(JobContext ctx) {
+ // vazio de proposito: este teste cobre o ciclo de vida pela
+ // borda HTTP real, nao o efeito da rotina
+ }
};
}
@@ -86,7 +90,7 @@ public void handle() { }
JobHandler e2eSlowHandler() {
return new JobHandler() {
public String type() { return "e2e-slow"; }
- public void handle() {
+ public void handle(JobContext ctx) {
TestGate.await(GATE.get());
}
};
diff --git a/src/test/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryTransitionsTest.java b/src/test/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryTransitionsTest.java
index 32df13b..ef9bca1 100644
--- a/src/test/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryTransitionsTest.java
+++ b/src/test/java/com/async/request/reply/adapter/out/persistence/jdbc/JdbcJobRepositoryTransitionsTest.java
@@ -163,6 +163,32 @@ void findFreshCompleted_should_ignore_active_jobs() {
assertThat(repository.findFreshCompleted(TYPE, NOW.minus(Duration.ofHours(1)))).isEmpty();
}
+ @Test
+ @DisplayName("findLastCompletedAt devolve a conclusão da última carga, sem janela")
+ void findLastCompletedAt_should_return_the_last_completed_instant() {
+ assertThat(repository.findLastCompletedAt(TYPE))
+ .as("nenhuma carga concluida ainda")
+ .isEmpty();
+
+ String id = newJob();
+ repository.complete(id);
+
+ // sem janela: e a pergunta "quando esses dados ficaram quentes?", que o
+ // consumidor responde no endpoint de dominio dele
+ assertThat(repository.findLastCompletedAt(TYPE)).contains(NOW);
+ }
+
+ @Test
+ @DisplayName("findLastCompletedAt ignora job ativo e job que falhou")
+ void findLastCompletedAt_should_ignore_active_and_failed_jobs() {
+ String failed = newJob();
+ repository.fail(failed, "Processing error", "estourou");
+
+ assertThat(repository.findLastCompletedAt(TYPE))
+ .as("carga que falhou nao deixou dado quente")
+ .isEmpty();
+ }
+
@Test
@DisplayName("id malformado é apenas um job inexistente, não um erro")
void findById_should_return_empty_when_id_is_not_a_uuid() {
diff --git a/src/test/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOutTest.java b/src/test/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOutTest.java
new file mode 100644
index 0000000..cf3f975
--- /dev/null
+++ b/src/test/java/com/async/request/reply/adapter/out/processing/AsyncJobProcessorAdapterOutTest.java
@@ -0,0 +1,194 @@
+package com.async.request.reply.adapter.out.processing;
+
+import com.async.request.reply.core.domain.Job;
+import com.async.request.reply.core.enums.JobStatus;
+import com.async.request.reply.core.port.out.JobRepositoryPortOut;
+import com.async.request.reply.spi.AsyncJobHandler;
+import com.async.request.reply.spi.JobContext;
+import com.async.request.reply.spi.JobHandler;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import java.time.Instant;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.contains;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@DisplayName("AsyncJobProcessorAdapterOut")
+class AsyncJobProcessorAdapterOutTest {
+
+ private static final String JOB_ID = "job-1";
+ private static final String TYPE = "contas";
+ private static final Instant NOW = Instant.parse("2026-07-01T12:00:00Z");
+
+ private final JobRepositoryPortOut repository = mock(JobRepositoryPortOut.class);
+
+ private final Job job = Job.pending(JOB_ID, TYPE, NOW);
+
+ private AsyncJobProcessorAdapterOut processor(Object... routines) {
+ when(repository.start(JOB_ID)).thenReturn(Optional.of(NOW));
+ JobHandlerRegistry registry = new JobHandlerRegistry(
+ List.of(routines).stream().map(r -> (com.async.request.reply.spi.Routine) r).toList());
+ return new AsyncJobProcessorAdapterOut(registry, repository);
+ }
+
+ @Test
+ @DisplayName("roda a rotina sincrona e completa o job no retorno")
+ void process_should_run_the_sync_routine_and_complete() {
+ AtomicInteger runs = new AtomicInteger();
+ processor(handler(_ -> runs.incrementAndGet())).process(job);
+
+ assertThat(runs).hasValue(1);
+ verify(repository).complete(JOB_ID);
+ }
+
+ @Test
+ @DisplayName("nao roda a rotina quando o start e recusado — job ja cancelado")
+ void process_should_not_run_the_routine_when_start_is_refused() {
+ AtomicInteger runs = new AtomicInteger();
+ JobHandlerRegistry registry = new JobHandlerRegistry(List.of(handler(_ -> runs.incrementAndGet())));
+ when(repository.start(JOB_ID)).thenReturn(Optional.empty());
+
+ new AsyncJobProcessorAdapterOut(registry, repository).process(job);
+
+ assertThat(runs).hasValue(0);
+ verify(repository, never()).complete(anyString());
+ }
+
+ @Test
+ @DisplayName("marca como falho quando a rotina lanca excecao")
+ void process_should_fail_the_job_when_the_routine_throws() {
+ processor(handler(_ -> {
+ throw new IllegalStateException("estourou");
+ })).process(job);
+
+ verify(repository).fail(JOB_ID, "Processing error", "estourou");
+ verify(repository, never()).complete(anyString());
+ }
+
+ @Test
+ @DisplayName("marca como falho quando nao existe rotina para o type")
+ void process_should_fail_the_job_when_no_routine_matches_the_type() {
+ JobHandlerRegistry registry = new JobHandlerRegistry(List.of(handler("outro-type", _ -> { })));
+ when(repository.start(JOB_ID)).thenReturn(Optional.of(NOW));
+
+ new AsyncJobProcessorAdapterOut(registry, repository).process(job);
+
+ verify(repository).fail(eq(JOB_ID), eq("Unknown job type"), contains(TYPE));
+ }
+
+ @Test
+ @DisplayName("rotina fire-and-forget nao e completada pela lib")
+ void process_should_not_complete_a_fire_and_forget_routine() {
+ AtomicReference seen = new AtomicReference<>();
+ AsyncJobHandler async = new AsyncJobHandler() {
+ public String type() { return TYPE; }
+ public void start(JobContext ctx) { seen.set(ctx.jobId()); }
+ };
+ when(repository.start(JOB_ID)).thenReturn(Optional.of(NOW));
+
+ new AsyncJobProcessorAdapterOut(new JobHandlerRegistry(List.of(async)), repository).process(job);
+
+ assertThat(seen).hasValue(JOB_ID);
+ verify(repository, never()).complete(anyString());
+ }
+
+ // --- cancelamento cooperativo -------------------------------------------
+
+ /**
+ * Cancelar não interrompe a thread da rotina: o desvio consciente do ADR 0004
+ * é resolvido de forma cooperativa — a rotina pergunta e decide parar.
+ */
+ @Test
+ @DisplayName("a rotina ve isCancelled() virar true quando o job e cancelado no meio")
+ void context_should_report_cancellation_while_the_routine_runs() {
+ when(repository.findById(JOB_ID))
+ .thenReturn(Optional.of(Job.pending(JOB_ID, TYPE, NOW)))
+ .thenReturn(Optional.of(cancelled()));
+
+ AtomicInteger chunks = new AtomicInteger();
+ processor(handler(ctx -> {
+ for (int i = 0; i < 10; i++) {
+ if (ctx.isCancelled()) {
+ return; // para no meio, sem processar os chunks restantes
+ }
+ chunks.incrementAndGet();
+ }
+ })).process(job);
+
+ assertThat(chunks)
+ .as("a rotina deve parar na segunda pergunta, nao rodar os 10 chunks")
+ .hasValue(1);
+ }
+
+ /**
+ * Uma rotina que para sozinha ainda retorna normalmente, e a lib chamaria
+ * {@code complete}. Quem protege o estado é o {@code UPDATE} condicional: ele
+ * recusa a transição em job terminal, então o cancelamento não é sobrescrito.
+ */
+ @Test
+ @DisplayName("complete de rotina que parou por cancelamento e recusado pelo repositorio")
+ void complete_should_be_refused_after_cooperative_stop() {
+ when(repository.findById(JOB_ID)).thenReturn(Optional.of(cancelled()));
+ when(repository.complete(JOB_ID)).thenReturn(Optional.empty());
+
+ processor(handler(ctx -> {
+ if (ctx.isCancelled()) {
+ return;
+ }
+ throw new AssertionError("deveria ter visto o cancelamento");
+ })).process(job);
+
+ assertThat(repository.complete(JOB_ID)).isEmpty();
+ }
+
+ @Test
+ @DisplayName("isCancelled() e false para job que segue ativo")
+ void context_should_report_not_cancelled_while_the_job_is_active() {
+ when(repository.findById(JOB_ID)).thenReturn(Optional.of(Job.pending(JOB_ID, TYPE, NOW)));
+
+ AtomicReference cancelled = new AtomicReference<>();
+ processor(handler(ctx -> cancelled.set(ctx.isCancelled()))).process(job);
+
+ assertThat(cancelled).hasValue(false);
+ }
+
+ /** Job que desapareceu do storage não é "cancelado" — não há o que parar. */
+ @Test
+ @DisplayName("isCancelled() e false quando o job nao e encontrado")
+ void context_should_report_not_cancelled_when_the_job_is_gone() {
+ when(repository.findById(JOB_ID)).thenReturn(Optional.empty());
+
+ AtomicReference cancelled = new AtomicReference<>();
+ processor(handler(ctx -> cancelled.set(ctx.isCancelled()))).process(job);
+
+ assertThat(cancelled).hasValue(false);
+ }
+
+ // --- helpers -------------------------------------------------------------
+
+ private static Job cancelled() {
+ return Job.restore(JOB_ID, TYPE, JobStatus.CANCELLED, NOW, NOW, null, null);
+ }
+
+ private static JobHandler handler(java.util.function.Consumer body) {
+ return handler(TYPE, body);
+ }
+
+ private static JobHandler handler(String type, java.util.function.Consumer body) {
+ return new JobHandler() {
+ public String type() { return type; }
+ public void handle(JobContext ctx) { body.accept(ctx); }
+ };
+ }
+}
diff --git a/src/test/java/com/async/request/reply/adapter/out/processing/JobHandlerRegistryTest.java b/src/test/java/com/async/request/reply/adapter/out/processing/JobHandlerRegistryTest.java
index bb7770d..6efb60e 100644
--- a/src/test/java/com/async/request/reply/adapter/out/processing/JobHandlerRegistryTest.java
+++ b/src/test/java/com/async/request/reply/adapter/out/processing/JobHandlerRegistryTest.java
@@ -1,5 +1,6 @@
package com.async.request.reply.adapter.out.processing;
+import com.async.request.reply.spi.JobContext;
import com.async.request.reply.spi.JobHandler;
import com.async.request.reply.spi.Routine;
import org.junit.jupiter.api.DisplayName;
@@ -63,7 +64,9 @@ public String type() {
}
@Override
- public void handle() {
+ public void handle(JobContext ctx) {
+ // vazio de proposito: este teste cobre só a indexacao por type,
+ // a rotina nunca chega a rodar
}
};
}
diff --git a/src/test/java/com/async/request/reply/config/AsyncJobsPropertiesTest.java b/src/test/java/com/async/request/reply/config/AsyncJobsPropertiesTest.java
index 4b27b4c..6fcb317 100644
--- a/src/test/java/com/async/request/reply/config/AsyncJobsPropertiesTest.java
+++ b/src/test/java/com/async/request/reply/config/AsyncJobsPropertiesTest.java
@@ -77,6 +77,36 @@ void should_accept_freshness_window_within_retention() {
.isEqualTo(Duration.ofMinutes(30));
}
+ /**
+ * O frescor procura a última carga concluída pela `coalescing_key`, que só é
+ * gravada quando o coalescing está ligado. Sem essa validação, ligar apenas
+ * `freshness.enabled` seria um no-op silencioso: a configuração pede para
+ * poupar carga e nada acontece.
+ */
+ @Test
+ @DisplayName("recusa frescor ligado sem coalescing: seria um no-op silencioso")
+ void should_reject_freshness_without_coalescing() {
+ AsyncJobsProperties.Freshness freshness =
+ new AsyncJobsProperties.Freshness(true, Duration.ofMinutes(30), Map.of());
+
+ assertThatThrownBy(() -> new AsyncJobsProperties(
+ null, null, false, null, null, freshness, null))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("async-jobs.coalesce-in-flight");
+ }
+
+ @Test
+ @DisplayName("aceita frescor ligado junto com coalescing")
+ void should_accept_freshness_with_coalescing() {
+ AsyncJobsProperties.Freshness freshness =
+ new AsyncJobsProperties.Freshness(true, Duration.ofMinutes(30), Map.of());
+
+ AsyncJobsProperties properties = new AsyncJobsProperties(
+ null, null, true, null, null, freshness, null);
+
+ assertThat(properties.freshness().enabled()).isTrue();
+ }
+
/** Frescor desligado não impõe relação com a retenção: a janela é inerte. */
@Test
@DisplayName("ignora a relacao quando o frescor esta desligado")
@@ -87,8 +117,9 @@ void should_ignore_the_relation_when_freshness_is_disabled() {
assertThat(properties(Duration.ofHours(1), freshness).freshness().enabled()).isFalse();
}
+ /** Coalescing ligado porque o frescor depende dele (ver teste acima). */
private static AsyncJobsProperties properties(Duration retention,
AsyncJobsProperties.Freshness freshness) {
- return new AsyncJobsProperties(retention, null, false, null, null, freshness, null);
+ return new AsyncJobsProperties(retention, null, true, null, null, freshness, null);
}
}
diff --git a/src/test/java/com/async/request/reply/core/usecase/JobFreshnessUseCaseTest.java b/src/test/java/com/async/request/reply/core/usecase/JobFreshnessUseCaseTest.java
new file mode 100644
index 0000000..9d1fd37
--- /dev/null
+++ b/src/test/java/com/async/request/reply/core/usecase/JobFreshnessUseCaseTest.java
@@ -0,0 +1,90 @@
+package com.async.request.reply.core.usecase;
+
+import com.async.request.reply.core.port.out.CoalescingKeyPortOut;
+import com.async.request.reply.core.port.out.JobFreshnessPolicyPortOut;
+import com.async.request.reply.core.port.out.JobRepositoryPortOut;
+import com.async.request.reply.spi.JobFreshness;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import java.time.Clock;
+import java.time.Duration;
+import java.time.Instant;
+import java.time.ZoneOffset;
+import java.util.Optional;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+@DisplayName("JobFreshnessUseCase")
+class JobFreshnessUseCaseTest {
+
+ private static final String TYPE = "contas";
+ private static final Instant NOW = Instant.parse("2026-07-01T12:00:00Z");
+
+ private final JobRepositoryPortOut repository = mock(JobRepositoryPortOut.class);
+ private final CoalescingKeyPortOut coalescingKey = type -> type;
+ private final Clock clock = Clock.fixed(NOW, ZoneOffset.UTC);
+
+ private JobFreshness freshness(Duration window) {
+ JobFreshnessPolicyPortOut policy = _ -> Optional.ofNullable(window);
+ return new JobFreshnessUseCase(repository, policy, coalescingKey, clock);
+ }
+
+ @Test
+ @DisplayName("devolve o instante da última carga concluída")
+ void lastRefreshedAt_should_return_the_last_completed_instant() {
+ Instant completedAt = NOW.minus(Duration.ofMinutes(10));
+ when(repository.findLastCompletedAt(TYPE)).thenReturn(Optional.of(completedAt));
+
+ assertThat(freshness(Duration.ofHours(1)).lastRefreshedAt(TYPE)).contains(completedAt);
+ }
+
+ @Test
+ @DisplayName("vazio quando nenhuma carga concluiu")
+ void lastRefreshedAt_should_be_empty_when_no_load_has_completed() {
+ when(repository.findLastCompletedAt(TYPE)).thenReturn(Optional.empty());
+
+ assertThat(freshness(Duration.ofHours(1)).lastRefreshedAt(TYPE)).isEmpty();
+ }
+
+ @Test
+ @DisplayName("quente quando a última carga está dentro da janela")
+ void isFresh_should_be_true_within_the_window() {
+ when(repository.findLastCompletedAt(TYPE))
+ .thenReturn(Optional.of(NOW.minus(Duration.ofMinutes(30))));
+
+ assertThat(freshness(Duration.ofHours(1)).isFresh(TYPE)).isTrue();
+ }
+
+ @Test
+ @DisplayName("frio quando a última carga é mais velha que a janela")
+ void isFresh_should_be_false_outside_the_window() {
+ when(repository.findLastCompletedAt(TYPE))
+ .thenReturn(Optional.of(NOW.minus(Duration.ofHours(2))));
+
+ assertThat(freshness(Duration.ofHours(1)).isFresh(TYPE)).isFalse();
+ }
+
+ /**
+ * Sem janela nada é quente — a mesma razão pela qual toda submissão dispara
+ * carga. Devolver {@code true} aqui faria a resposta de domínio prometer
+ * frescura que a lib não está garantindo.
+ */
+ @Test
+ @DisplayName("frio quando não há janela configurada, mesmo com carga recente")
+ void isFresh_should_be_false_when_no_window_is_configured() {
+ when(repository.findLastCompletedAt(TYPE)).thenReturn(Optional.of(NOW));
+
+ assertThat(freshness(null).isFresh(TYPE)).isFalse();
+ }
+
+ @Test
+ @DisplayName("frio quando nenhuma carga concluiu")
+ void isFresh_should_be_false_when_no_load_has_completed() {
+ when(repository.findLastCompletedAt(TYPE)).thenReturn(Optional.empty());
+
+ assertThat(freshness(Duration.ofHours(1)).isFresh(TYPE)).isFalse();
+ }
+}